Apache DolphinScheduler OpenMLDB 任务节点:SQL 驱动机器学习特征开发的完整实践指南
OpenMLDB 是开源机器学习数据库,提供面向生产环境的全栈 FeatureOps 解决方案。本文聚焦 Apache DolphinScheduler 的 OpenMLDB 任务节点(Task Plugin),系统讲解如何在 DolphinScheduler 工作流中连接 OpenMLDB 集群、通过 SQL 完成数据导入与特征抽取,并深入源码揭示 SQL 到 Python 脚本的转换原理与执行细节。读完本文,你将掌握 OpenMLDB 节点的完整配置方法、离线/在线双模式执行机制,以及基于真实源码的底层实现原理与可复现的验证方法。
概述:什么是 OpenMLDB 节点
OpenMLDB 是一个优秀的开源机器学习数据库,为生产环境提供端到端的 FeatureOps 解决方案,覆盖从特征开发、特征上线到特征存储与计算的完整链路。DolphinScheduler 的 OpenMLDB 任务插件用于在 DolphinScheduler 工作流中直接对 OpenMLDB 集群执行任务,使特征工程作业能够像普通 ETL 任务一样被编排、调度与监控。
从实现上看,OpenMLDB 节点并不是独立实现的一套执行引擎,而是建立在 Python 任务基础之上的 SQL 执行器:节点将用户在界面上填写的 SQL 语句自动转换为一段 Python 脚本,脚本通过 OpenMLDB Python SDK(openmldb)连接集群并逐条执行 SQL。插件源码位于 dolphinscheduler-task-plugin/dolphinscheduler-task-openmldb,其中 OpenmldbTask 继承自 PythonTask(见 PythonTask.java),任务参数类 OpenmldbParameters 同样继承自 PythonParameters(见 OpenmldbParameters.java)。
任务节点通过 OpenmldbTaskChannelFactory 以 SPI 机制注册到 DolphinScheduler,getName() 返回的任务类型名称为 OPENMLDB(见 OpenmldbTaskChannelFactory.java)。
创建 OpenMLDB 任务
在 DolphinScheduler Web UI 中创建 OpenMLDB 任务的步骤如下:
- 进入
项目管理 -> 项目名称 -> 工作流定义,点击 创建工作流 按钮,进入 DAG 编辑页面; - 从左侧工具栏中拖拽 OpenMLDB 任务节点(图标标识)到画布中;
- 在右侧配置面板中填写任务参数并保存。
需要说明的是,创建任务节点的操作路径与通用工作流一致,任务节点在 DAG 工具栏中的注册由前端统一管理,相关任务字段定义可参考 use-openmldb.ts。
任务参数详解
除了 DolphinScheduler 任务参数附录 中 默认任务参数 一节说明的通用参数(如任务名称、运行环境、优先级、失败重试等)之外,OpenMLDB 节点还包含以下专有参数:
| 参数 | 描述 | 是否必填 |
|---|---|---|
| zookeeper 地址 | OpenMLDB 集群连接地址中的 zookeeper 地址,例如 127.0.0.1:2181 |
是 |
| zookeeper 路径 | OpenMLDB 集群连接地址中的 zookeeper 路径,例如 /openmldb |
是 |
| 执行模式 | 初始执行模式(离线 / 在线),可在 SQL 语句中随时切换 | 是 |
| SQL 语句 | 要执行的 SQL 语句,支持多条语句以分号分隔 | 是 |
| 自定义参数 | Python 用户自定义参数,脚本中 \${variable} 占位内容将被替换 |
否 |
其中 zookeeper 地址与路径是连接 OpenMLDB 集群的核心寻址信息,二者共同构成 SQLAlchemy 连接串 openmldb:///?zk=<zk>&zkPath=<zkPath>,用于定位集群的元数据与数据存储位置。前端表单对 zk、zkPath 与 sql 三个字段均设置了必填校验(见 use-openmldb.ts);executeMode 在前端以单选框形式提供 offline / online 两个选项(见 use-openmldb.ts)。
参数校验规则同样体现在服务端:OpenmldbParameters.checkParameters() 要求 zk、zkPath、sql 三者均非空才判定参数有效(见 OpenmldbParameters.java),校验失败时 OpenmldbTask.init() 会抛出 TaskException("openmldb task params is not valid")(见 OpenmldbTask.java)。
执行模式:离线与在线
执行模式是 OpenMLDB 节点的关键概念,决定了任务初始运行在哪个存储/计算引擎上:
- 离线(offline):数据导入至离线存储,特征计算由离线引擎(Spark 等)执行,适合大规模批量特征工程;
- 在线(online):面向在线服务场景,特征计算由在线引擎执行,适合低延迟的实时特征查询。
该模式对应 OpenMLDB 的 SET @@execute_mode='offline' / SET @@execute_mode='online' 会话变量。正如参数表中强调的,执行模式只是节点的初始模式——由于模式是会话级设置,你完全可以在同一段 SQL 中通过 SET @@execute_mode=... 语句随时切换,例如在离线模式导入数据后再切到在线模式执行在线查询。
SQL 到 Python 脚本的转换原理
OpenMLDB 节点的核心工作机制是将用户编写的 SQL 语句转换为 Python 脚本,再由 Python 解释器执行。理解这一过程有助于排查问题与编写高效的任务 SQL。
参数渲染
任务启动后,OpenmldbTask.init() 将任务参数 JSON 反序列化为 OpenmldbParameters。随后在 buildPythonScriptContent() 中(见 OpenmldbTask.java):
- 将 SQL 中的
\r\n统一规范为\n,保证脚本内容一致性; - 调用
mergeParamsWithContext()合并任务上下文中的参数(继承自 Python 任务的变量机制),再通过ParameterUtils.convertParameterPlaceholders()将 SQL 文本中的\${variable}占位符替换为实际参数值。
脚本模板生成
替换完成后的 SQL 会经 buildPythonScriptsFromSql() 生成最终的 Python 脚本(见 OpenmldbTask.java),生成逻辑如下:
import openmldb
import sqlalchemy as db
engine = db.create_engine('openmldb:///?zk=<zk>&zkPath=<zkPath>')
con = engine.connect()
con.execute("set @@execute_mode='<mode>';")
若执行模式为 offline,还会额外注入两条设置,确保离线任务同步执行并设置合理的超时时间:
con.execute("set @@sync_job=true")
con.execute("set @@job_timeout=1800000")
这两条语句的语义是:离线任务必须以同步方式执行,且 job 超时时间预设为 30 分钟(1800000 毫秒),与 OpenMLDB 服务端 server.channel_keep_alive_time 保持一致。如果你预期离线任务执行时间会超过 30 分钟,可以在 SQL 语句中通过 SET @@job_timeout=... 自行调大(源码注释已明确提示这一点,见 OpenmldbTask.java)。
随后 SQL 按分号 ; 切分为多条语句,跳过仅含空白字符的片段,每条非空语句都会被转换为一行 con.execute("<sql>")(SQL 内部的换行会被转义为 \n 以保持单行字符串形式)。
执行与文件落盘
生成好的 Python 脚本内容与脚本路径(<执行路径>/openmldb_<taskAppId>.py,见 OpenmldbTask.java)会交给继承自 PythonTask 的 handle() 流程:先写入带 #-*- encoding=utf8 -*- 编码声明的 .py 文件,再以 python3 <脚本路径> 命令通过 ShellCommandExecutor 执行(见 PythonTask.java)。执行模式与任务结果由 ShellCommandExecutor 统一管理,退出码、进程号与输出参数会在任务完成后回写。
任务示例:以 SQL 完成特征工程
导入数据(Load Data)
LOAD DATA INFILE 'hdfs:///data/training_data.csv'
INTO TABLE t1
OPTIONS(mode='overwrite');
示例中我们使用 LOAD DATA 语句将数据导入 OpenMLDB 集群。由于节点执行模式选择了 offline,数据会被导入离线存储,供后续批量特征计算使用。
说明:示例图片展示的是真实运行界面截图(见 openmldb-load-data.png),界面中
LOAD DATA语句、zookeeper 地址与执行模式等参数均与上表一一对应。
特征抽取(Feature Extraction)
SELECT cid, window_start, COUNT(*) AS cnt
FROM t1
WINDOW (ORDER BY ts ROWS_RANGE BETWEEN 10s PRECEDING AND CURRENT ROW)
INTO OUTFILE 'hdfs:///tmp/feature_output';
示例中我们使用 SELECT INTO(SELECT ... INTO OUTFILE 语法)进行特征抽取,将计算结果写出到指定路径。因为执行模式为 offline,该 SQL 会由 OpenMLDB 的离线引擎执行计算,适合对全量历史数据进行滑动窗口等复杂特征加工。
说明:示例界面截图见 openmldb-feature-extraction.png,展示了特征抽取任务在 DAG 编辑器中的参数配置。
环境准备
在执行 OpenMLDB 任务之前,需要准备以下两部分环境。
启动 OpenMLDB 集群
首先必须有一个可用的 OpenMLDB 集群。生产环境请参考 OpenMLDB 官方部署文档完成集群安装与配置;如需快速体验,可以按照 OpenMLDB 官方 Docker 快速启动方案,在 Docker 中拉起一套最小集群用于验证任务链路。
集群启动后,确认 zookeeper 地址与路径与任务参数中填写的一致(即 zk 与 zkPath 必须精确匹配集群实际注册信息)。
Python 环境与 OpenMLDB SDK
OpenMLDB 任务通过 OpenMLDB Python SDK 连接集群,因此需要满足:
- Python 版本:任务默认使用
python3执行。源码中OPENMLDB_PYTHON = "python3"硬编码了默认解释器(见 OpenmldbTask.java),OpenMLDB SDK 默认只支持 Python 3; - SDK 安装:在 Worker 服务器所在主机上执行
pip install openmldb,确保 SDK 安装到了python3对应的解释器环境中(因为 Python 脚本最终由 Worker 节点拉起执行); - 自定义 Python 环境(可选):可以通过设置环境变量
PYTHON_LAUNCHER指定自定义 Python 环境。getPythonCommand()的处理逻辑为(见 OpenmldbTask.java):- 未设置
PYTHON_LAUNCHER时,直接使用python3; - 若
PYTHON_LAUNCHER指向xx/bin/python[版本号]形式的路径,会强制替换为xx/bin/python3; - 其余情况则在
PYTHON_LAUNCHER目录下拼接/bin/python3。
- 未设置
该设计保证了无论何种自定义环境,最终执行的始终是 python3 解释器。
深入验证:单元测试中的转换逻辑
插件自带的单元测试 OpenmldbTaskTest.java 直接验证了上述转换逻辑,是理解节点行为的绝佳参考:
buildPythonExecuteCommand()用例断言默认执行命令为python3 test.py,验证了默认解释器逻辑;buildSQLWithComment()用例构造了一段包含注释、空语句与多条查询的混合 SQL,并断言生成脚本内容完全符合预期:包括连接串openmldb:///?zk=localhost:2181&zkPath=dolphinscheduler、offline模式下的set @@execute_mode='offline';、set @@sync_job=true、set @@job_timeout=1800000三条会话设置,以及逐条转换后的con.execute(...)语句。
同时该测试也揭示了一个值得注意的边界行为:SQL 中的 -- 注释不会被剥离,而是作为 SQL 语句内容的一部分原样传给 OpenMLDB 执行;空语句(只有空白字符)会被跳过。
常见问题与排查建议
- 参数校验失败(task params is not valid):确认
zookeeper、zookeeper path、SQL statement三个参数均已填写。前端必填校验与后端checkParameters()双重把关,任一缺失都会直接报错; - 连接集群失败:核对
zk与zkPath是否与 OpenMLDB 集群实际注册信息一致;确保 Worker 主机与集群网络互通; - Python 相关报错(如找不到 openmldb 模块):确认
pip install openmldb已安装到 Worker 主机上,且与python3/PYTHON_LAUNCHER指向的解释器环境一致; - 离线任务长时间不返回:离线模式默认
job_timeout为 30 分钟,长耗时作业请在 SQL 中通过SET @@job_timeout=<毫秒>调大超时; - SQL 书写规范:多条语句用分号
;分隔;需要切换模式时直接写入SET @@execute_mode='online';即可在任务中途切换执行引擎。
小结
OpenMLDB 节点将 DolphinScheduler 的工作流编排能力与 OpenMLDB 的特征工程能力无缝衔接:你只需在 DAG 中拖入节点、填写 zookeeper 地址与 SQL 语句,即可将数据导入、特征抽取等 FeatureOps 作业纳入统一的调度与运维体系。通过源码可以看到,节点的实现本质上是"SQL 到 Python 脚本"的自动转换器,其执行模式(离线/在线)、同步离线任务与超时控制、参数占位符替换等行为均有清晰的代码实现与单元测试支撑,为在生产环境中的深度使用与问题排查提供了可靠依据。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust4.24 K638- DDeepSeek-V4.1-FlashDeepSeek-V4.1-Flash 是一个多模态混合专家(MoE)模型,拥有 5520 亿骨干参数,并支持最多一百万 token 的上下文长度。该模型原生支持图像和文本输入,并以自回归方式生成文本Python670
SlideSCIPPT插件,支持素材库、AI助手、一键添加图片标题,复制粘贴位置、一键图片对齐、一键插入Markdown(加粗、超链接等行内样式、代码块、LaTeX等块级样式)、便捷导出图片!C#230
hello-agents📚 《从零开始构建智能体》——从零开始的智能体原理与实践教程Python52874
new-apiAI模型聚合管理中转分发系统,一个应用管理您的所有AI模型,支持将多种大模型转为统一格式调用,支持OpenAI、Claude、Gemini等格式,可供个人或者企业内部管理与分发渠道使用。🍥 A Unified AI Model Management & Distribution System. Aggregate all your LLMs into one app and access them via an OpenAI-compatible API, with native support for Claude (Messages) and Gemini formats.Go22545
JeecgBoot🔥企业级低代码平台集成了AI应用平台,帮助企业快速实现低代码开发和构建AI应用!前后端分离架构 SpringBoot,SpringCloud、Mybatis,Ant Design4、 Vue3.0、TS+vite!强大的代码生成器让前后端代码一键生成,无需写任何代码! 引领AI低代码开发模式: AI生成->OnlineCoding-> 代码生成-> 手工MERGE,显著的提高效率,又不失灵活~Java36351