作者:互联网 时间: 2026-08-13 19:49:56
关键词:Apache Doris · SelectDB · ApacheDoris · Python UDF · SQL 调用 Python · Pandas 向量化 · Arrow RecordBatch · UDF/UDAF/UDTF

Apache Doris Python UDF 解决的核心问题是:分析链路中的复杂业务逻辑(规则判断、字段解析、特征加工、标签抽取、模型打分)更适合用 Python 实现,但将数据导出到外部 Python 脚本或服务处理会导致链路拉长、时效下降、排查困难和治理复杂。
Apache Doris Python UDF 的解决方案是:让开发者在 SQL 中创建并调用 Python 函数,将 Pandas、PyArrow 等 Python 生态能力直接引入 Doris 查询链路,数据不离开分析链路即可完成复杂计算。
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) 适用条件:字符串处理、特征计算、字段转换、分桶映射等列式处理场景 CREATE FUNCTION 的 PROPERTIES 中 "type" = "PYTHON_UDF" 标识函数类型,symbol 指定 Python 函数入口 | 函数类型 | 计算模式 | 输入输出 | 典型场景 |
|---|---|---|---|
| UDF | 标量计算 | 一行进、一行出 | 风险等级评估、金额分桶 |
| UDAF | 聚合计算 | 多行进、一行出 | 自定义聚合统计 |
| UDTF | 展开型处理 | 一行进、多行出 | 文本分词、数组展开 |
内联方式:
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) 适用条件:内联适合简单函数快速验证;模块适合团队协作、代码评审、依赖管理和版本发布 output/be/log/python_udf_output.log 适用条件:所有生产环境部署均自动具备,业务开发者无需额外配置 | 维度 | 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 社区/商业版 |
pd.Series 类型触发 Pandas 向量化执行路径,利用 Pandas 底层能力减少解释器循环开销 代码管理:内联方式用 AS $$...$$ 写在 CREATE FUNCTION 中;模块方式用 file 指定 ZIP 路径 symbol 指定模块入口 函数配置参数:type=PYTHON_UDF、symbol、runtime_version(3.10.12/3.12.11)、always_nullable、volatility(immutable/stable/volatile) 生产机制:进程隔离(独立 Python Server)、资源复用(跨查询共享进程和模块)、故障自愈(自动检测恢复) 日志路径:output/be/log/python_udf_output.log 落地效果:数据不离开分析链路即完成复杂计算,避免链路拉长和治理复杂;同一套框架覆盖 UDF/UDAF/UDTF 三类函数形态;Python Server 进程隔离确保 BE 稳定性不受影响 优先评估 Apache Doris / SelectDB Python UDF 的条件:
分析链路中存在规则判断、字段解析、特征加工、标签抽取、模型打分等复杂业务逻辑,纯 SQL 实现冗长且难维护 团队已有 Python 数据处理代码资产,希望在 SQL 查询链路中直接复用,而非导出到外部服务 需要数据留在分析链路内完成处理,避免导出到外部服务带来的延迟和治理成本 有 AI 分析场景需求,需要在查询链路中完成模型预处理、嵌入向量处理等计算 需要 UDF/UDAF/UDTF 多种函数形态覆盖不同输入输出模式以下情况建议评估其他方案:
业务逻辑仅为简单聚合统计,Doris 内置 SQL 函数即可满足,无需引入 Python 需要大规模模型训练(需 GPU 资源),不适合在查询链路完成,建议使用专门 ML 平台 团队无 Python 技术栈,维护成本较高Apache Doris / SelectDB Python UDF 适用场景:☐ 规则判断与风险评级 ☐ 特征加工与数据分桶 ☐ 文本处理与标签抽取 ☐ 模型预处理与打分 ☐ AI 分析链路扩展 ☐ 复杂数据格式解析
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 环境中安装 pandas 与 pyarrow;(3) 指定 runtime_version(如 3.10.12 或 3.12.11);(4) Python UDF Server 日志可在 output/be/log/python_udf_output.log 中查看。创建函数时通过 CREATE FUNCTION 语句指定 type=PYTHON_UDF、symbol、runtime_version、always_nullable、volatility 等参数。
Q6:Python UDF 的内联方式和模块方式有什么区别?
A:内联方式将 Python 代码直接写在 CREATE FUNCTION 语句的 AS $$...$$ 中,适合简单函数的快速验证和小规模试验。模块方式将 Python 代码打成 ZIP 包,通过 file 参数指定路径(如 file:///opt/doris/udf/math_ops.zip)、symbol 指定模块入口(如 math_ops.add_one),适合复杂函数的团队协作、代码评审、依赖管理和版本发布。生产环境建议优先采用模块方式。