Apache DolphinScheduler OpenMLDB 任务节点实战:连接机器学习数据库执行特征工程
导读
OpenMLDB 是一个开源机器学习数据库,提供生产级数据及特征开发全栈解决方案。Apache DolphinScheduler 通过内置的 OpenMLDB 任务组件(Task Plugin)连接 OpenMLDB 集群,以 SQL 方式在 DAG 工作流中完成数据导入与特征抽取。本文将基于 openmldb.md 文档并结合仓库源码,完整讲解 OpenMLDB 任务节点的创建步骤、参数含义、SQL 转 Python 脚本的底层执行原理、典型任务样例与运行环境准备,帮助你在一张 DAG 中把"离线数据加载 → 特征计算 → 在线服务"的特征工程链路编排起来。
综述:OpenMLDB 任务组件在 DolphinScheduler 中的定位
OpenMLDB 任务组件(dolphinscheduler-task-openmldb)是 DolphinScheduler 众多任务插件之一,其作用是连接 OpenMLDB 集群并执行任务。它并不像 Shell、Python 任务那样直接跑一段操作系统脚本,而是把你在任务表单中填写的 SQL 语句 转换成一个 Python 脚本,再调用 OpenMLDB Python SDK 驱动 SQL 执行,从而复用 OpenMLDB 完整的 SQL 能力(建表、导入、特征抽取、在线查询等)。
从源码结构看,插件位于 dolphinscheduler-task-plugin/dolphinscheduler-task-openmldb,主要由以下四个类构成:
| 类 | 职责 |
|---|---|
OpenmldbTask |
任务执行主体,继承 PythonTask,负责把 SQL 渲染为 Python 脚本并执行 |
OpenmldbParameters |
任务参数模型,持有 zk、zkPath、executeMode、sql 字段并做参数校验 |
OpenmldbTaskChannel |
任务通道,负责创建 OpenmldbTask 与解析参数 JSON |
OpenmldbTaskChannelFactory |
插件工厂,向 DolphinScheduler 注册 OPENMLDB 任务类型 |
同时,前端在 use-openmldb.ts 中提供了对应的表单渲染逻辑,任务类型标识为 OPENMLDB(见 use-openmldb.ts)。
创建 OpenMLDB 任务
在 DolphinScheduler Web UI 中创建一个 OpenMLDB 任务节点,操作步骤如下:
- 进入 项目管理 → 项目名称 → 工作流定义,点击 创建工作流 按钮,进入 DAG 编辑页面;
- 从左侧工具栏中拖动
OpenMLDB 任务节点到画板中(图标文件见 docs/img/tasks/icons/openmldb.png); - 双击该节点,在右侧任务配置面板中填写下述任务参数,并连接前置任务、保存并上线工作流。
任务参数详解
除所有任务通用的默认参数外(详见 DolphinScheduler任务参数附录 中的"默认任务参数"一栏,包括任务名称、运行标志、任务优先级、Worker分组、失败重试次数、超时告警等),OpenMLDB 任务节点还有如下四个核心参数:
| 任务参数 | 描述 |
|---|---|
| zookeeper 地址 | OpenMLDB 集群连接地址中的 zookeeper 地址,例如 127.0.0.1:2181 |
| zookeeper 路径 | OpenMLDB 集群连接地址中的 zookeeper 路径,例如 /openmldb |
| 执行模式 | 初始执行模式(离线 offline / 在线 online),可以在 SQL 语句中随时切换 |
| SQL 语句 | 需要在 OpenMLDB 集群上执行的 SQL 语句 |
参数校验与必填约束
从源码可以确认这些字段的校验逻辑。OpenmldbParameters.checkParameters()(见 OpenmldbParameters.java)要求 zk、zkPath、sql 三者均非空,否则 OpenmldbTask.init()(见 OpenmldbTask.java)会抛出 TaskException("openmldb task params is not valid")。
前端表单与之一致:use-openmldb.ts 中 zk、zkPath、sql 均为必填项(带 required 校验),executeMode 则是单选按钮,可选项为 offline 和 online,默认值 offline。除此之外,该任务还支持 自定义参数(Custom parameters / localParams),用于以 ${variable} 形式在 SQL 脚本中做变量替换——这一点与 Python 任务行为一致,因为 OpenmldbParameters 直接继承自 PythonParameters。
任务执行原理:SQL 是如何被执行的
理解 OpenMLDB 任务的执行原理,关键在于知道它不是直接向集群发送 SQL,而是先由 OpenmldbTask 把 SQL 渲染成一段 Python 脚本,再通过 Python 解释器运行。整条链路如下:
SQL 语句 → 参数渲染(替换 ${变量})→ 生成 openmldb_<taskAppId>.py → python3 执行 → OpenMLDB Python SDK → OpenMLDB 集群
生成 Python 脚本
buildPythonScriptContent()(见 OpenmldbTask.java)先对原始 SQL 做两件事:
- 把
\r\n统一规范为\n; - 通过
ParameterUtils.convertParameterPlaceholders把${variable}形式的自定义参数占位符替换为实际值(参数来源为任务的自定义参数与上下文参数合并结果)。
随后调用 buildPythonScriptsFromSql(rawSqlScript)(见 OpenmldbTask.java)逐行拼接 Python 代码:
import openmldb
import sqlalchemy as db
engine = db.create_engine('openmldb:///?zk=<zk>&zkPath=<zkPath>')
con = engine.connect()
con.execute("set @@execute_mode='offline';")
con.execute("set @@sync_job=true")
con.execute("set @@job_timeout=1800000")
# ……逐条 con.execute("<sql>")
几个关键实现细节:
- 连接串:通过 SQLAlchemy 的
openmldb://dialect 建立连接,zk与zkPath直接拼入连接 URL; - 初始执行模式:首条语句
set @@execute_mode='<mode>';根据"执行模式"参数设置,mode 统一转为小写; - 离线模式特殊处理:当执行模式为
offline时,会自动追加set @@sync_job=true(离线任务改为同步等待执行结果)和set @@job_timeout=1800000(任务超时设为 30 分钟,对应 OpenMLDB 服务端server.channel_keep_alive_time),你可以在 SQL 中通过set @@job_timeout=...覆盖更长的超时; - SQL 切分:按
;切分多条语句,跳过只含空格的空语句,保留注释与换行(换行转义为\n内嵌进 Python 字符串),因此一条 SQL 可以包含多条语句与--注释。
生成脚本文件并执行
buildPythonCommandFilePath()(见 OpenmldbTask.java)把脚本写到任务执行目录下,文件名形如 openmldb_<taskAppId>.py。最终执行命令由 buildPythonExecuteCommand() 拼出,即 <python命令> openmldb_<taskAppId>.py。
Python 解释器选择逻辑
OpenmldbTask 继承自 PythonTask,解释器选择逻辑在 getPythonCommand()(见 OpenmldbTask.java):
- 未设置
PYTHON_LAUNCHER环境变量时,默认使用python3; - 若设置了
PYTHON_LAUNCHER,且其路径形如xx/bin/python[数字.],会被强制替换为/bin/python3(因为 OpenMLDB 默认只支持 Python 3); - 若设置的是目录,则拼接为
<pythonHome>/bin/python3。
单元测试印证
OpenmldbTaskTest.java 中的 buildSQLWithComment 测试精确断言了"带注释的多条 SQL → Python 脚本"的渲染结果:脚本以 import openmldb、import sqlalchemy as db 开头,连接串为 openmldb:///?zk=localhost:2181&zkPath=dolphinscheduler,离线模式下包含 set @@sync_job=true 与 set @@job_timeout=1800000,且空语句(末尾的 ;)被跳过、注释被保留在语句中。该测试同时验证了 buildPythonExecuteCommand("test.py") 输出为 python3 test.py,与上述解释器选择逻辑一致。
任务样例
下面通过两个典型场景演示 OpenMLDB 任务的 SQL 写法。两个样例均选择 offline(离线) 执行模式。
样例一:导入数据(LOAD DATA)
USE demo_db;
set @@job_timeout=200000;
LOAD DATA INFILE 'file:///tmp/train_sample.csv'
INTO TABLE talkingdata OPTIONS(mode='overwrite');
说明:
USE demo_db;切换到目标数据库;- 因为执行模式为 offline,
LOAD DATA会把数据导入 离线存储; set @@job_timeout=200000;将本次任务超时调整为 200 秒(单位毫秒);OPTIONS(mode='overwrite')表示以覆盖方式写入目标表。
样例二:特征抽取(SELECT INTO / INTO OUTFILE)
select is_attributed, ip, app, device, os, channel, hour,
count(channel) over w1 as qty,
count(channel) over w2 as ip_app_count,
count(channel) over w3 as ip_app_os_count
from demo_db.talkingdata
window
w1 as (partition by ip order by click_time ROWS BETWEEN ...),
w2 as (partition by ip, app order by click_time ROWS BETWEEN ...),
w3 as (partition by ip, app, os order by click_time ROWS BETWEEN ...)
INTO OUTFILE '/tmp/train_feature';
说明:
- 使用窗口函数(
over w1/w2/w3)基于click_time排序、按ip/app/os分组统计channel计数,生成特征列; - 因为执行模式为 offline,该 SQL 会使用 离线引擎 做特征计算;
INTO OUTFILE '/tmp/train_feature'把特征结果导出到指定路径,供下游训练使用。
在真实工作流中,可将"导入数据"与"特征抽取"两个 OpenMLDB 节点串联成一条 DAG,也可在 SQL 中通过 set @@execute_mode='online'; 动态切换执行模式,实现离线训练与在线服务的一体化编排。
环境准备
启动 OpenMLDB 集群
执行 OpenMLDB 任务之前,必须保证 OpenMLDB 集群已经启动并可连接。生产环境请按照 OpenMLDB 官方部署文档安装部署;快速验证场景下,可以在 Docker 中运行 OpenMLDB 集群完成一键启动。集群启动后,任务参数中的 zookeeper 地址与 zookeeper 路径必须与集群实际注册信息一致(例如 127.0.0.1:2181 与 /openmldb)。
Python 环境
由于 OpenMLDB 任务组件依赖 OpenMLDB Python SDK 连接集群,Worker 所在主机必须满足:
- 安装 Python 环境:默认使用
python3;如使用自定义 Python,可通过设置PYTHON_LAUNCHER环境变量来指定,该变量既可以是可执行文件路径,也可以是 Python 安装目录(源码会强制归一为python3,详见上文"Python 解释器选择逻辑"); - 安装 SDK:在 Worker Server 主机上执行
pip install openmldb,确保openmldb与sqlalchemy模块可被python3正常导入(生成的脚本首行即import openmldb)。
排障提示
- 若任务在
init阶段失败并提示openmldb task params is not valid,请检查 zookeeper 地址、zookeeper 路径与 SQL 语句是否均已填写; - 若报
ModuleNotFoundError: openmldb,说明 Worker 主机上未安装 OpenMLDB Python SDK,或PYTHON_LAUNCHER指向的 Python 环境与安装 SDK 的环境不一致; - 若离线任务长时间不返回,可检查
set @@job_timeout配置是否过小(源码默认注入 1800000ms 超时,可在 SQL 中覆盖)。
小结
OpenMLDB 任务组件将 DolphinScheduler 的 DAG 编排能力与 OpenMLDB 的 SQL 特征工程能力衔接在一起:你只需要在任务节点中填写 zookeeper 连接信息、执行模式与一段 SQL,引擎便会自动生成 Python 脚本并通过 OpenMLDB Python SDK 提交到集群。配合离线/在线执行模式切换、自定义参数与超时控制,即可把数据导入、特征抽取乃至在线查询完整地纳入统一的数据编排工作流。
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 的上下文长度。该模型原生支持图像和文本输入,并以自回归方式生成文本Python650
SlideSCIPPT插件,支持素材库、AI助手、一键添加图片标题,复制粘贴位置、一键图片对齐、一键插入Markdown(加粗、超链接等行内样式、代码块、LaTeX等块级样式)、便捷导出图片!C#180
hello-agents📚 《从零开始构建智能体》——从零开始的智能体原理与实践教程Python52774
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

