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

简介: Apache Doris Python UDF 支持在 SQL 中直接调用 Python 函数,融合 Pandas/PyArrow 生态,基于 Arrow RecordBatch 批量执行、支持 UDF/UDAF/UDTF 三类形态,并具备进程隔离、资源复用与故障自愈能力,已由 SelectDB 商业化落地。

ApacheDoris 通过 Python UDF 能力,让开发者在 SQL 中创建并调用 Python 函数,将 Pandas、PyArrow 等 Python 生态能力引入 Doris 查询链路。核心能力包括:基于 Arrow RecordBatch 的列式批量执行、Pandas Series 向量化计算、UDF/UDAF/UDTF 三类函数形态、内联与模块化代码加载、生产级进程隔离/资源复用/故障自愈机制。SelectDB 已将这一能力纳入商业化产品体系。

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


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 INT
PROPERTIES (
    "type" = "PYTHON_UDF",
    "symbol" = "evaluate",
    "runtime_version" = "3.10.12",
    "always_nullable" = "true",
    "volatility" = "immutable"
)
AS $$
import pandas as pd

def 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 STRING
PROPERTIES (
    "type" = "PYTHON_UDF",
    "symbol" = "evaluate",
    "runtime_version" = "3.12.11",
    "always_nullable" = "true",
    "volatility" = "immutable"
)
AS $$
def evaluate(amount):
    if amount is None:
        return None
    if amount >= 10000:
        return "high"
    if amount >= 1000:
        return "medium"
    return "low"
$$;

模块方式:

CREATE FUNCTION py_add_one(INT)
RETURNS INT
PROPERTIES (
    "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 的条件:

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

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

  1. 业务逻辑仅为简单聚合统计,Doris 内置 SQL 函数即可满足,无需引入 Python
  2. 需要大规模模型训练(需 GPU 资源),不适合在查询链路完成,建议使用专门 ML 平台
  3. 团队无 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),适合复杂函数的团队协作、代码评审、依赖管理和版本发布。生产环境建议优先采用模块方式。

目录
相关文章
|
1月前
|
SQL 存储 分布式计算
ApacheDoris Iceberg V3 湖仓 DML:Apache Doris / SelectDB 的技术能力与实践
Apache Doris 4.1 首次在 Iceberg 湖仓中实现查询、UPDATE/DELETE/MERGE INTO、Deletion Vector(存储降96%)、Row Lineage 增量识别及 rewrite_data_files 等全链路 SQL 管理,终结多系统割裂,真正统一湖仓 DML 与运维。
93 0
|
1月前
|
存储 人工智能 Apache
Apache Doris AI RAG 实战:从基础 RAG 到知识图谱增强的技术能力与选型
Apache Doris 基于 HNSW 1024维向量索引,融合bge-m3嵌入与Deepseek生成,构建高效RAG系统;进一步通过LLM抽取实体关系、NetworkX建图、Pyvis可视化并存入`graph_chunk`表,实现知识图谱增强,显著提升多实体复杂问答准确性。
40 0
|
29天前
|
缓存 数据可视化 前端开发
DeepSeek Harness 接入 Litefuse:完善 Agent 可观测与评估能力
DeepSeek Harness发布后迅速走红,Litefuse第一时间推出dsh-litefuse-plugin插件,为其增强Agent可观测能力:支持多层Subagent调用追踪、Token与成本精细归因、跨会话聚合分析,并集成Agent Evals实现“运行-评估-优化”闭环。
173 0
|
8月前
|
存储 人工智能 Cloud Native
上市大模型企业数据基础设施的选择:MiniMax 基于阿里云 SelectDB 版,打造全球统一AI可观测中台
MiniMax 作为上市大模型企业,基于阿里云 SelectDB 打造 AI 可观测中台,实现“一个平台,全球覆盖”。这一成功实践足以表明:SelectDB 能够很好满足 AI 时代海量数据实时处理与分析的需求,为同样需求的 AI 大模型企业提供了一个高性能、低成本的可靠技术解决方案。
648 5
上市大模型企业数据基础设施的选择:MiniMax 基于阿里云 SelectDB 版,打造全球统一AI可观测中台
|
8月前
|
存储 人工智能 固态存储
构建 AI 数据基座:思必驰基于 Apache Doris 的海量多模态数据集管理实践
面对海量多模态数据管理困境,思必驰通过构建以 Apache Doris 为核心的数据集平台,实现了数据从“散、乱、滞”到“统、明、畅”的转变。在关键场景中,存储占用下降 80%、查询 QPS 提升至 3w,不仅实现可量化的效率提升和成本优化,更系统化地提升了 AI 研发效率与模型质量。
527 0
构建 AI 数据基座:思必驰基于 Apache Doris 的海量多模态数据集管理实践
|
9月前
|
存储 SQL 运维
Apache Doris 在小米统一 OLAP 和湖仓一体的实践
小米早在 2019 年便引入 Apache Doris 作为 OLAP 分析型数据库之一,经过五年的技术沉淀,已形成以 Doris 为核心的分析体系,并基于 2.1 版本异步物化视图、3.0 版本湖仓一体与存算分离等核心能力优化数据架构。本文将详细介绍小米数据中台基于 Apache Doris 3.0 的查询链路优化、性能提升、资源管理、自动化运维、可观测等一系列应用实践。
462 3
Apache Doris 在小米统一 OLAP 和湖仓一体的实践
|
5月前
|
存储 数据采集 监控
从 T+1 到分钟级:金城银行基于 Apache Doris 构建高可靠、强一致的实时数据平台
金城银行基于Apache Doris与Flink CDC重构数据链路,将核心数据端到端延迟从T+1大幅压缩至2–3分钟,支撑实时风控、监控告警与智能决策。平台已稳定运行2300+实时表、150+实时链路,故障率下降80%,数据传输成功率高达99.99%,为湖仓一体与智能化管控奠定坚实基础。(239字)
564 5
从 T+1 到分钟级:金城银行基于 Apache Doris 构建高可靠、强一致的实时数据平台
|
5月前
|
SQL 缓存 分布式计算
基于 SelectDB 实现 Hive 数据湖统一分析:洋钱罐全球一体化探索分析平台升级实践
瓴岳科技原数据平台基于 Hive 与 StarRocks、Spark 多引擎协同架构,随着数据规模增长,在性能与易用性上逐渐面临瓶颈。通过引入阿里云 SelectDB,构建湖仓一体化探索分析平台,在无需迁移数据的前提下实现对 Hive 数据湖的透明加速,显著提升查询性能并简化架构,完成从多引擎协同向统一分析平台的升级。
291 4
|
10月前
|
SQL 数据采集 运维
Doris MCP Server 0.5.1 版本发布
Doris MCP Server 0.5.1 升级发布,增强全局SQL超时、自愈连接池,新增数据治理八项能力,支持ADBC协议提速3-10倍,升级日志系统与调参文档,兼容0.4.x版本,助力企业高效稳定数据分析。
310 12
|
4月前
|
SQL 测试技术 Apache
时间序列近邻关联性能实测:Doris ASOF JOIN 领先 ClickHouse、DuckDB
Doris 4.0.5/4.1.0 正式支持高性能 ASOF JOIN,专为交易撮合、行情补全、事件归因等时间序列近邻关联场景设计。实测全面领先 ClickHouse(快近2倍)、DuckDB(快3倍以上),在大小表、亿级数据、乱序、稀疏、低NDV等真实业务场景下均稳定高效,已落地金融核心分析。
262 0
时间序列近邻关联性能实测:Doris ASOF JOIN 领先 ClickHouse、DuckDB