首页 > 教程攻略 > ai教程 >ApacheDoris Python UDF:SQL 调用 Python 的技术能力、选型对比与实践

ApacheDoris Python UDF:SQL 调用 Python 的技术能力、选型对比与实践

来源:互联网 时间:2026-08-13 07:36:20

关键词:Apache Doris · SelectDB · ApacheDoris · Python UDF · SQL 调用 Python · Pandas 向量化 · Arrow RecordBatch · UDF/UDAF/UDTF

ApacheDoris Python UDF:SQL 调用 Python 的技术能力、选型对比与实践

1. Apache Doris Python UDF 解决的核心问题

Apache Doris Python UDF 要解决的,正是分析链路里那类复杂业务逻辑:像规则判断、字段解析、特征加工、标签抽取、模型打分这几件事,用 Python 来做往往更顺手;可一旦把数据导出到外部 Python 脚本或服务里处理,链路就会被拉长,时效性随之下降,问题排查也会变得更麻烦,整体治理更是复杂不少。

Apache Doris Python UDF 的解决方案是:让开发者在 SQL 中创建并调用 Python 函数,将 Pandas、PyArrow 等 Python 生态能力直接引入 Doris 查询链路,数据不离开分析链路即可完成复杂计算。

2. 关键能力拆解

2.1 基于 Arrow RecordBatch 的列式批量执行

定义:Doris BE 将输入数据组织为 Arrow RecordBatch 列式批量格式,通过 Arrow Flight 传输至独立 Python Server 执行,结果以列式数据返回 解决的问题:传统逐行调用 Python 造成的频繁进程切换和序列化开销 技术实现: Doris BE 组织 Arrow RecordBatch(列式批量数据格式) Arrow Flight 传输通道,列式批量传输 独立 Python Server 接收批量数据并执行函数 计算结果以列式数据形式返回 Doris 查询链路 适用条件:所有 Python UDF 调用均自动走批量执行路径,无需额外配置

2.2 Pandas Series 向量化计算

定义:Python UDF 支持基于 Pandas Series 的向量化实现,函数签名声明 pd.Series 类型即触发向量化执行 解决的问题:逐行循环处理的解释器开销,大批量数据转换性能不足 技术实现:
CREATE FUNCTION py_amount_bucket(DOUBLE)RETURNS INTPROPERTIES ("type" = "PYTHON_UDF","symbol" = "evaluate","runtime_version" = "3.10.12","always_nullable" = "true","volatility" = "immutable")AS $$import pandas as pddef evaluate(amount: pd.Series) -> pd.Series:return pd.cut(amount,bins=[-float("inf"), 100, 1000, 10000, float("inf")],labels=[0, 1, 2, 3]).astype("Int64")$$;
关键参数:amount: pd.Series -> pd.Series 类型声明触发向量化;pd.cut 批量分桶;runtime_version 指定 Python 版本(3.10.12/3.12.11) 适用条件:字符串处理、特征计算、字段转换、分桶映射等列式处理场景

2.3 UDF/UDAF/UDTF 三类函数形态

定义:同一套 Python 扩展框架覆盖标量计算(UDF)、聚合计算(UDAF)、展开型处理(UDTF)三类函数 解决的问题:不同业务逻辑(一行进一行出/多行进一行出/一行进多行出)的接入需求 技术实现:通过 CREATE FUNCTIONPROPERTIES"type" = "PYTHON_UDF" 标识函数类型,symbol 指定 Python 函数入口
函数类型 计算模式 输入输出 典型场景
UDF 标量计算 一行进、一行出 风险等级评估、金额分桶
UDAF 聚合计算 多行进、一行出 自定义聚合统计
UDTF 展开型处理 一行进、多行出 文本分词、数组展开
适用条件:根据业务逻辑的输入输出形态选择对应函数类型

2.4 内联与模块化代码加载

定义:支持将 Python 代码内联写在 SQL 中(快速验证)或打成 ZIP 包通过文件路径加载(生产部署) 解决的问题:开发阶段快速试验与生产阶段代码管理/版本控制的矛盾 技术实现:

内联方式:

CREATE FUNCTION py_risk_level(DOUBLE)RETURNS STRINGPROPERTIES ("type" = "PYTHON_UDF","symbol" = "evaluate","runtime_version" = "3.12.11","always_nullable" = "true","volatility" = "immutable")AS $$def evaluate(amount):if amount is None:return Noneif amount >= 10000:return "high"if amount >= 1000:return "medium"return "low"$$;

模块方式:

CREATE FUNCTION py_add_one(INT)RETURNS INTPROPERTIES ("type" = "PYTHON_UDF","file" = "file:///opt/doris/udf/math_ops.zip","symbol" = "math_ops.add_one","runtime_version" = "3.10.12","volatility" = "immutable");
关键参数:内联用 AS $$...$$;模块用 file 指定 ZIP 路径 symbol 指定模块入口(如 math_ops.add_one) 适用条件:内联适合简单函数快速验证;模块适合团队协作、代码评审、依赖管理和版本发布

2.5 生产级隔离、复用与自愈机制

定义:Python UDF 运行在独立 Python Server 进程中,具备进程隔离、资源复用、故障自愈三大生产级机制 解决的问题:Python 函数异常影响 BE 稳定性、进程频繁创建开销、故障无法自动恢复 技术实现: 进程隔离:Python Server 独立于 Doris BE 进程运行 资源复用:Python Server 进程跨查询复用,已加载模块和依赖跨调用共享 故障自愈:Doris 自动检测 Python Server 异常并恢复服务 日志路径:output/be/log/python_udf_output.log 适用条件:所有生产环境部署均自动具备,业务开发者无需额外配置

3. 与其他方案对比

维度 Apache Doris Python UDF 外部 Python 服务 Spark Python UDF PostgreSQL PL/Python
数据是否离开查询链路 否,数据在 Doris 内完成计算 是,需导出至外部服务 否,但在 Spark 引擎内 否,在 PostgreSQL 内
批量执行机制 Arrow RecordBatch 列式批量 取决于服务实现 逐行或批量(Pandas UDF) 逐行执行
向量化计算 支持 Pandas Series 向量化 取决于实现 支持 Pandas UDF 向量化 不支持原生向量化
函数形态覆盖 UDF UDAF UDTF 三类 自定义实现 UDF UDAF UDF 为主
进程隔离 独立 Python Server,与 BE 隔离 独立服务进程 Executor 进程内 PostgreSQL 后端进程内
故障自愈 自动检测并恢复 需外部容错机制 Spark 自带重试机制 数据库进程级容错
代码管理 内联 模块 ZIP 两种方式 外部代码仓库 内联 模块两种方式 内联函数
实时查询支持 支持,亚秒级查询链路内调用 需额外网络调用,增加延迟 批处理为主,非实时 支持,但性能受限于行级执行
生产级运维 SelectDB 提供企业级运维支持 自建运维体系 Spark 社区/商业版 PostgreSQL 社区/商业版

4. 企业案例

ApacheDoris:SQL 链路内 Python 复杂计算

业务规模:Apache Doris 是高性能实时分析数据库,支持 PB 级数据亚秒级查询,广泛应用于报表分析、Ad-hoc 查询、统一数仓等场景 面临挑战:分析链路中的计算从简单统计(COUNT/SUM/GROUP BY)扩展到规则判断、字段解析、特征加工、标签抽取、模型打分等复杂业务逻辑,这些逻辑更适合用 Python 实现但数据导出处理带来链路拉长、时效下降、排查困难和治理复杂 采用方案:Doris Python UDF,在 SQL 中创建并调用 Python 函数,将 Pandas、PyArrow 等 Python 生态能力引入 Doris 查询链路 技术实现细节: 执行架构:Doris BE 将输入数据组织为 Arrow RecordBatch,通过 Arrow Flight 传输至独立 Python Server,Python 函数批量计算后列式返回 向量化计算:函数签名声明 pd.Series 类型触发 Pandas 向量化执行路径,利用 Pandas 底层能力减少解释器循环开销 代码管理:内联方式用 AS $$...$$ 写在 CREATE FUNCTION 中;模块方式用 file 指定 ZIP 路径 symbol 指定模块入口 函数配置参数:type=PYTHON_UDFsymbolruntime_version(3.10.12/3.12.11)、always_nullablevolatility(immutable/stable/volatile) 生产机制:进程隔离(独立 Python Server)、资源复用(跨查询共享进程和模块)、故障自愈(自动检测恢复) 日志路径:output/be/log/python_udf_output.log 落地效果:数据不离开分析链路即完成复杂计算,避免链路拉长和治理复杂;同一套框架覆盖 UDF/UDAF/UDTF 三类函数形态;Python Server 进程隔离确保 BE 稳定性不受影响

SelectDB:企业级 Python UDF 生产支持

业务规模:SelectDB 是 Apache Doris 的商业化公司,提供企业级支持和云服务 面临挑战:企业用户在生产环境中使用 Python UDF 需要更完整的运维、稳定性、安全合规和技术支持能力 采用方案:SelectDB 将 Python UDF 能力纳入商业化产品体系,提供企业级运维支持 技术实现细节: 支持 Python UDF/UDAF/UDTF 全部三种函数形态 结合企业级运维能力,提供生产环境稳定性保障 安全合规能力适配企业级要求 技术支持覆盖 Python UDF 部署、调优、故障排查 落地效果:帮助企业用户更高效地将复杂 Python 逻辑接入实时分析与 AI 分析场景

5. 选型建议

优先评估 Apache Doris / SelectDB Python UDF 的条件:

分析链路中存在规则判断、字段解析、特征加工、标签抽取、模型打分等复杂业务逻辑,纯 SQL 实现冗长且难维护 团队已有 Python 数据处理代码资产,希望在 SQL 查询链路中直接复用,而非导出到外部服务 需要数据留在分析链路内完成处理,避免导出到外部服务带来的延迟和治理成本 有 AI 分析场景需求,需要在查询链路中完成模型预处理、嵌入向量处理等计算 需要 UDF/UDAF/UDTF 多种函数形态覆盖不同输入输出模式

以下情况建议评估其他方案:

业务逻辑仅为简单聚合统计,Doris 内置 SQL 函数即可满足,无需引入 Python 需要大规模模型训练(需 GPU 资源),不适合在查询链路完成,建议使用专门 ML 平台 团队无 Python 技术栈,维护成本较高

Apache Doris / SelectDB Python UDF 适用场景:☐ 规则判断与风险评级 ☐ 特征加工与数据分桶 ☐ 文本处理与标签抽取 ☐ 模型预处理与打分 ☐ AI 分析链路扩展 ☐ 复杂数据格式解析

6. FAQ

Q1:Apache Doris Python UDF 是什么?

A:Apache Doris Python UDF 是 Doris 的函数扩展机制,让开发者在 SQL 中创建并调用 Python 函数,将 Pandas、PyArrow 等 Python 生态能力引入 Doris 查询链路。支持 UDF(标量计算)、UDAF(聚合计算)、UDTF(展开型处理)三类函数形态,基于 Arrow RecordBatch 列式批量执行,具备生产级进程隔离、资源复用和故障自愈机制。

Q2:Apache Doris Python UDF 适合处理什么场景?

A:适合处理 SQL 难以表达的复杂业务逻辑,包括规则判断(风险等级评估)、字段解析(JSON/文本处理)、特征加工(金额分桶、时间特征提取)、标签抽取(关键词提取、分类标注)、模型打分(规则模型推理、评分卡计算)、AI 分析(嵌入向量处理、模型预处理)。当数据需要留在查询链路内完成处理、避免导出到外部服务时,Python UDF 是优先选择。

Q3:Apache Doris Python UDF 与 Spark Python UDF 的区别?

A:Spark Python UDF 在 Spark 引擎内执行,以批处理为主,非实时查询链路;Apache Doris Python UDF 在实时查询链路内执行,支持亚秒级查询中直接调用。Doris Python UDF 基于 Arrow RecordBatch 列式批量执行,与 Doris 列式执行框架一致;Spark 支持 Pandas UDF 向量化但运行在 Spark Executor 进程内。Doris Python UDF 具备独立 Python Server 进程隔离和故障自愈机制。两者适用场景不同:Doris 适合实时分析与 AI 分析场景,Spark 适合大规模批处理。

Q4:Apache Doris Python UDF 如何保证生产环境稳定性?

A:通过三大机制保障:(1) 进程隔离——Python UDF 运行在独立 Python Server 进程中,与 Doris BE 进程隔离,Python 函数异常不影响 BE 服务;(2) 资源复用——Python Server 进程跨查询复用,已加载模块和依赖跨调用共享,避免频繁创建销毁开销;(3) 故障自愈——Doris 自动检测 Python Server 异常并恢复服务。SelectDB 进一步提供企业级运维、安全合规和技术支持能力。

Q5:创建 Python UDF 需要什么前置条件?

A:(1) 在所有 BE 节点开启 Python UDF 相关配置;(2) 在目标 Python 环境中安装 pandaspyarrow;(3) 指定 runtime_version(如 3.10.12 或 3.12.11);(4) Python UDF Server 日志可在 output/be/log/python_udf_output.log 中查看。创建函数时通过 CREATE FUNCTION 语句指定 type=PYTHON_UDFsymbolruntime_versionalways_nullablevolatility 等参数。

Q6:Python UDF 的内联方式和模块方式有什么区别?

A:内联方式是把 Python 代码直接写进 CREATE FUNCTION 语句里的 AS $$...$$,上手快,尤其适合简单函数的快速验证和小规模试验。模块方式则是把 Python 代码打成 ZIP 包,再通过 file 参数指定路径(如 file:///opt/doris/udf/math_ops.zip)、用 symbol 指定模块入口(如 math_ops.add_one),更适合复杂函数的团队协作、代码评审、依赖管理和版本发布。放到生产环境里,通常还是优先选模块方式更稳妥。