首页
/ Apache DolphinScheduler OpenMLDB 任务节点实战:连接机器学习数据库执行特征工程

Apache DolphinScheduler OpenMLDB 任务节点实战:连接机器学习数据库执行特征工程

2026-09-14 14:04:13作者:申梦珏Efrain

导读

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 任务参数模型,持有 zkzkPathexecuteModesql 字段并做参数校验
OpenmldbTaskChannel 任务通道,负责创建 OpenmldbTask 与解析参数 JSON
OpenmldbTaskChannelFactory 插件工厂,向 DolphinScheduler 注册 OPENMLDB 任务类型

同时,前端在 use-openmldb.ts 中提供了对应的表单渲染逻辑,任务类型标识为 OPENMLDB(见 use-openmldb.ts)。

创建 OpenMLDB 任务

在 DolphinScheduler Web UI 中创建一个 OpenMLDB 任务节点,操作步骤如下:

  1. 进入 项目管理 → 项目名称 → 工作流定义,点击 创建工作流 按钮,进入 DAG 编辑页面;
  2. 从左侧工具栏中拖动 Apache DolphinScheduler OpenMLDB 任务节点实战:连接机器学习数据库执行特征工程 OpenMLDB 任务节点到画板中(图标文件见 docs/img/tasks/icons/openmldb.png);
  3. 双击该节点,在右侧任务配置面板中填写下述任务参数,并连接前置任务、保存并上线工作流。

任务参数详解

除所有任务通用的默认参数外(详见 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)要求 zkzkPathsql 三者均非空,否则 OpenmldbTask.init()(见 OpenmldbTask.java)会抛出 TaskException("openmldb task params is not valid")

前端表单与之一致:use-openmldb.tszkzkPathsql 均为必填项(带 required 校验),executeMode 则是单选按钮,可选项为 offlineonline,默认值 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 做两件事:

  1. \r\n 统一规范为 \n
  2. 通过 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 建立连接,zkzkPath 直接拼入连接 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 openmldbimport sqlalchemy as db 开头,连接串为 openmldb:///?zk=localhost:2181&zkPath=dolphinscheduler,离线模式下包含 set @@sync_job=trueset @@job_timeout=1800000,且空语句(末尾的 ;)被跳过、注释被保留在语句中。该测试同时验证了 buildPythonExecuteCommand("test.py") 输出为 python3 test.py,与上述解释器选择逻辑一致。

任务样例

下面通过两个典型场景演示 OpenMLDB 任务的 SQL 写法。两个样例均选择 offline(离线) 执行模式。

样例一:导入数据(LOAD DATA)

OpenMLDB任务配置:使用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; 切换到目标数据库;
  • 因为执行模式为 offlineLOAD DATA 会把数据导入 离线存储
  • set @@job_timeout=200000; 将本次任务超时调整为 200 秒(单位毫秒);
  • OPTIONS(mode='overwrite') 表示以覆盖方式写入目标表。

样例二:特征抽取(SELECT INTO / INTO OUTFILE)

OpenMLDB任务配置:使用窗口函数离线抽取特征并导出

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 所在主机必须满足:

  1. 安装 Python 环境:默认使用 python3;如使用自定义 Python,可通过设置 PYTHON_LAUNCHER 环境变量来指定,该变量既可以是可执行文件路径,也可以是 Python 安装目录(源码会强制归一为 python3,详见上文"Python 解释器选择逻辑");
  2. 安装 SDK:在 Worker Server 主机上执行 pip install openmldb,确保 openmldbsqlalchemy 模块可被 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 提交到集群。配合离线/在线执行模式切换、自定义参数与超时控制,即可把数据导入、特征抽取乃至在线查询完整地纳入统一的数据编排工作流。

登录后查看全文
热门项目推荐
相关项目推荐