Article

任务调度 Dolphin Scheduler

更新于:2026-07-13

第一章:DolphinScheduler 概述

1.1 什么是 DolphinScheduler

概念名称说明注意事项
DolphinScheduler一个分布式、易扩展的可视化工作流任务调度系统,原名 EasyScheduler,由 Apache 孵化并开源。专注于大数据任务编排,支持多种任务类型(Shell、SQL、Spark 等),通过 DAG(有向无环图)方式定义任务依赖关系。不是传统 Cron 工具的替代品,而是面向复杂数据流水线的调度平台,适合企业级批量任务管理。
开源协议Apache License 2.0,允许自由使用、修改和分发。商用需遵守开源协议,注意衍生作品的版权标注。
核心目标解决任务依赖复杂、调度不稳定、运维困难等问题,提供高可用、可监控、可追溯的任务执行能力。强调”任务编排”而非”定时触发”,更注重流程完整性与容错机制。

1.2 核心特性与优势

特性名称说明注意事项
分布式调度支持 Master/Worker 架构,任务在 Worker 节点执行,Master 负责调度协调,可水平扩展。需依赖 ZooKeeper 实现 Master 高可用和 Worker 注册发现。
可视化 DAG 编排通过 Web UI 拖拽方式构建任务依赖关系,直观展示执行流程。初学者需理解 DAG 原理,避免循环依赖。
多种任务类型支持内置 Shell、SQL、Spark、Flink、Python、HTTP、DataX 等任务类型,满足大数据生态需求。某些任务需额外配置环境路径或客户端(如 Spark-submit)。
高可用(HA)Master 和 Worker 均支持多节点部署,故障自动转移。需正确配置 ZooKeeper 集群,避免脑裂问题。
灵活调度策略支持 Cron 表达式、手动触发、补数(backfill)、依赖触发等多种调度方式。Cron 表达式需符合 Quartz 格式(6 或 7 位)。
参数化执行支持全局参数、局部参数、运行时参数传递,提升工作流复用性。参数类型分为 VALUE 和 DATE,后者支持日期函数(如 $[yyyyMMdd])。
告警通知机制支持邮件、短信、钉钉、企业微信、飞书等多种通知方式,任务失败可及时告警。需提前配置告警组和通知模板,测试连通性。
多租户资源隔离通过租户(Tenant)机制隔离操作系统用户和资源,确保任务运行安全。租户需绑定系统用户,Worker 节点上该用户必须存在且有执行权限。

1.3 适用场景与典型用户

场景/用户类型说明注意事项
数据仓库 ETL 流程定时抽取、清洗、加载数据,如从 MySQL 同步到 Hive,再跑 Spark 任务。需结合 DataX 或 Sqoop 任务类型实现数据迁移。
机器学习流水线训练前数据准备 → 模型训练(Python/Spark)→ 模型评估 → 结果写入数据库。可通过 Python 任务调用 sklearn/pytorch 等框架。
日报/周报自动化每日凌晨自动执行 SQL 统计任务,生成报表并邮件发送。需配置邮件告警模板,结合 SQL 任务的”查询结果通知”功能。
运维自动化脚本调度定时清理日志、备份文件、检查服务状态等 Shell 脚本统一管理。注意脚本退出码(exit code),非 0 视为失败。
互联网公司数据中台作为统一调度平台,整合多个业务线的数据任务,实现集中监控。需合理划分项目、租户和权限,避免权限混乱。
典型用户数据工程师、ETL 开发者、运维工程师、数据分析师、AI 工程师。非技术人员可通过查看任务日志了解执行状态,但不能修改工作流。

1.4 架构概览(Master/Worker/Alert/Api Server等角色)

组件名称说明注意事项
Master Server负责工作流解析、任务调度、DAG 编排、任务分发,是调度核心。多 Master 通过 ZooKeeper 选举主节点,避免单点故障。
Worker Server接收 Master 分配的任务,执行具体任务脚本或程序,上报执行状态。Worker 按租户(Tenant)分配任务,需确保系统用户存在且有权限。
API Server提供 RESTful API 接口,供前端 Web UI 或外部系统调用。所有操作(如启动工作流)最终都通过 API Server 处理。
Alert Server处理告警事件,根据配置将任务失败等信息发送到邮件、钉钉等渠道。需单独配置告警插件和通知组,支持自定义告警脚本。
Frontend (Web UI)基于 Vue 的前端界面,用于可视化操作和监控任务。可独立部署,通过 Nginx 反向代理连接后端 API。
ZooKeeper用于 Master/Worker 的注册与发现、高可用选举、任务状态协调。生产环境建议部署 3 节点以上 ZooKeeper 集群。
Database (MySQL/PostgreSQL)存储元数据,如用户、项目、工作流定义、执行实例、日志等。初始化需执行 dolphinscheduler-*-create.sql 脚本建表。
Logger Server收集 Worker 执行日志,写入本地文件或 HDFS(可选)。日志路径通常为 logs/workflow/ 目录下,按实例 ID 存储。

第二章:环境准备与安装部署

2.1 环境依赖要求(JDK、MySQL、ZooKeeper等)

依赖项版本要求说明注意事项
JDK1.8+(推荐 OpenJDK 或 Oracle JDK)DolphinScheduler 使用 Java 开发,所有节点均需安装 JDK。必须配置 JAVA_HOME 环境变量,且版本为 64 位。
MySQL5.7 或 8.0作为元数据库存储系统信息,API Server 和 Master 需连接。需创建独立数据库(如 dolphinscheduler)并授权用户远程访问。
PostgreSQL9.4+可替代 MySQL 作为元数据库(社区版支持)。配置时需修改 application-dao.yml 中的数据库类型和驱动。
ZooKeeper3.4.6+(推荐 3.5+)用于 Master/Worker 的注册发现与高可用协调。建议独立部署集群,避免与 Hadoop 共用导致性能瓶颈。
传输工具rsync、ssh(免密登录)用于节点间文件同步和远程命令执行。所有 Worker 节点需配置与 Master 的 SSH 免密登录。
操作系统Linux(CentOS、Ubuntu 等)支持主流 Linux 发行版,不推荐 Windows 部署生产环境。用户需有 sudo 权限或能切换到目标租户用户。

2.2 单机模式部署步骤

步骤操作命令/说明注意事项
1. 下载安装包wget https://downloads.apache.org/dolphinscheduler/1.3.9/apache-dolphinscheduler-1.3.9-bin.tar.gz选择稳定版本,建议从 Apache 官网下载。
2. 解压文件tar -zxvf apache-dolphinscheduler-1.3.9-bin.tar.gz -C /opt/解压路径避免中文和空格。
3. 创建数据库CREATE DATABASE dolphinscheduler DEFAULT CHARACTER SET utf8mb4;推荐使用 utf8mb4 支持 emoji。
4. 修改数据库配置编辑 conf/application-dao.yml,设置 datasource.url, username, password确保网络可达,测试连接可用。
5. 初始化数据库sh script/create-dolphinscheduler.sh自动执行建表脚本,输出无错误即成功。
6. 修改运行用户编辑 conf/config/install_config.conf,设置 ips="127.0.0.1" masters="127.0.0.1" workers="127.0.0.1" installPath="/opt/dolphinscheduler" deployUser="ds"deployUser 必须是系统已存在用户,且有写权限。
7. 执行部署脚本sh script/deploy.sh脚本自动复制文件、创建软链接、启动服务。
8. 验证服务jps 查看是否有 MasterServer、WorkerServer、ApiApplicationServer若无进程,检查 logs/ 目录下对应日志文件。

2.3 集群模式部署指南

配置项说明注意事项
节点规划Master 节点:2 台(高可用);Worker 节点:≥2 台(按负载分配);公共节点:ZooKeeper、MySQL、API Server(可共用)API Server 可部署在 Master 节点上,也可独立部署。
SSH 免密配置在部署机上执行:ssh-keygen -t rsassh-copy-id user@master1ssh-copy-id user@worker1 ...所有节点之间需互相免密,包括自己到自己。
install_config.conf 配置示例ips="master1,master2,worker1,worker2"masters="master1,master2"workers="worker1:default,worker2:default"zooQuorum="zk1:2181,zk2:2181,zk3:2181"dbtype="mysql"dbname="dolphinscheduler"workers 后可指定租户分组,default 为默认组。
部署命令sh script/deploy.sh脚本会通过 SSH 将文件分发到各节点并启动服务。
高可用验证停止一台 Master,观察另一台是否接管调度任务查看 ZooKeeper 中 /dolphinscheduler/masters 节点变化。
资源隔离通过租户(Tenant)绑定不同系统用户,实现 Worker 上的任务隔离确保每个租户对应的系统用户在所有 Worker 上存在。

2.4 初始化数据库与配置文件详解

配置文件关键配置项说明注意事项
conf/application-dao.ymlspring.datasource.urlspring.datasource.usernamespring.datasource.password数据库连接信息生产环境建议使用专用账号,限制权限。
conf/application-master.ymlmaster.exec.threads=100master.exec.task.num=20master.dispatch.task.num=10Master 线程池与任务调度参数根据 CPU 核数调整,避免过载。
conf/application-worker.ymlworker.exec.threads=100worker.max.cpu.load.avg=-1worker.reserved.memory=0.3Worker 执行线程与资源限制reserved.memory 表示保留内存比例(GB),-1 表示不限制。
conf/application-api.ymlserver.port=12345logging.level.org.apache.dolphinscheduler=INFOAPI Server 端口与日志级别可通过 Nginx 代理对外暴露 80/443 端口。
conf/config/install_config.confresource.storage.type=NONEdata.basedir.path=/tmp/dolphinschedulertenant.auto.create=true安装时的全局配置data.basedir.path 是任务工作目录,需保证磁盘空间充足。
conf/worker.properties(旧版本)worker.group=defaultworker.max.cpuload.avg=-1Worker 分组与负载控制新版本使用 application-worker.yml 替代。

2.5 启动与验证服务状态

操作命令说明注意事项
启动所有服务sh bin/start-all.sh启动 Master、Worker、API Server仅适用于单机或已部署的集群。
停止所有服务sh bin/stop-all.sh安全停止所有进程建议在维护时使用。
单独启动 Mastersh bin/dolphinscheduler-daemon.sh start master-server用于故障恢复或调试查看 logs/master-server/*.log 确认启动成功。
单独启动 Workersh bin/dolphinscheduler-daemon.sh start worker-server确保 ZooKeeper 已连接。
单独启动 API Serversh bin/dolphinscheduler-daemon.sh start api-server启动后可通过 curl http://localhost:12345/dolphinscheduler/doc.html 查看 API 文档。
查看进程jps输出应包含:MasterServer、WorkerServer、ApiApplicationServer若缺失某进程,检查对应日志文件。
查看日志tail -f logs/master-server/master-server.log实时观察错误信息常见错误:数据库连接失败、ZooKeeper 超时、端口占用。
Web UI 登录浏览器访问 http://<ip>:12345/dolphinscheduler/ui,默认账号:admin / 密码:dolphinscheduler首次登录建议修改密码若无法访问,检查防火墙或 Nginx 配置。

第三章:Web UI 快速入门

3.1 登录与初始界面介绍

元素名称说明注意事项
登录地址默认为 http://<api-server-ip>:12345/dolphinscheduler/ui确保 API Server 正常运行且端口开放
默认账号密码用户名:admin;密码:dolphinscheduler首次登录后建议立即修改密码
主界面布局 - 顶部导航栏包含”项目管理”、“资源中心”、“安全中心”、“监控中心”、“个人中心”等模块入口不同角色权限显示不同菜单项
主界面布局 - 左侧菜单在”项目管理”中,可切换不同项目,查看工作流定义与实例仅当前用户有权限的项目可见
主界面布局 - 中央区域展示任务 DAG 图、执行日志、调度记录等核心信息支持缩放、拖拽、节点右键操作
多语言支持支持中文和英文界面切换(右上角语言选择)国际化仍在完善中,部分提示仍为英文
退出登录右上角点击头像 → “退出登录”推荐操作,避免会话泄露

3.2 用户、项目、租户的基本概念

概念说明注意事项
用户(User)系统登录者,拥有唯一用户名和密码,可分配角色(Admin、普通用户等)Admin 用户可管理其他用户,普通用户只能操作所属项目
项目(Project)工作流的容器,用于组织和隔离任务流程。每个项目包含工作流定义、任务实例、UDF 函数等用户需被授权才能加入项目,权限分为:管理员、项目经理、开发员、运维员、访客
租户(Tenant)对应操作系统用户,用于 Worker 节点上执行任务的身份隔离。任务在哪个租户下运行,就以该系统用户身份执行租户必须提前在所有 Worker 节点创建,否则任务无法启动
角色(Role)控制用户权限,如安全管理、项目创建、告警配置等自定义角色需谨慎赋权,避免越权操作
资源中心存储脚本、JAR 包等共享资源的目录,可绑定到项目资源上传后可在 Shell、Spark 等任务中引用

3.3 创建第一个工作流任务

操作步骤说明注意事项
1. 进入项目点击”项目管理” → 选择或创建一个项目若无项目,需先由 Admin 或自己创建
2. 创建工作流点击”工作流定义” → “创建工作流”按钮弹出 DAG 编辑器界面
3. 添加任务节点从左侧任务类型栏拖拽”Shell”任务到画布可重命名节点,如”echo_hello”
4. 编辑任务内容双击节点 → 在”脚本”框输入:echo "Hello DolphinScheduler"脚本支持变量,如 ${start_time}
5. 设置租户在”运行标志”下方选择已有租户(如 tenant_ds)必须选择,否则提交时报错
6. 连接节点(单节点无需连接)使用鼠标从一个节点拖线到另一个节点建立依赖DAG 不允许环路
7. 保存工作流点击右上角”保存”图标,填写名称(如”first_workflow”)保存后可在”工作流定义”列表看到

3.4 手动执行与查看日志

操作说明注意事项
手动启动工作流在”工作流定义”列表中,找到目标工作流 → 点击”上线” → “运行”按钮必须先上线才能运行
查看执行实例跳转至”工作流实例”页签,查看刚触发的实例状态状态包括:正在运行、成功、失败、停止等
查看任务日志点击实例中的具体任务节点 → “查看日志”日志实时刷新,包含标准输出和错误信息
日志关键字搜索在日志窗口使用 Ctrl+F 搜索关键词(如”error”)建议开启”自动滚动”跟踪最新输出
终止任务右键任务节点 → “停止”仅对正在运行的任务有效,底层调用 kill -9
查看重试情况若配置了失败重试,可在日志中看到多次尝试记录重试间隔默认 1 分钟
下载日志文件点击”下载日志”按钮,保存本地分析日志路径通常为 /tmp/dolphinscheduler/logs/

3.5 常见操作快捷方式

操作快捷方式 / 技巧说明注意事项
刷新页面数据F5 或 Ctrl+R适用于实例状态未更新时手动刷新避免频繁刷新影响性能
全屏 DAG 图点击 DAG 编辑器右上角”全屏”图标更清晰查看复杂流程ESC 键退出全屏
节点复制粘贴选中节点 → Ctrl+C → Ctrl+V提高重复结构构建效率粘贴后需重新配置脚本和参数
批量删除节点Shift+拖拽框选多个节点 → Delete 键快速清理无效节点删除不可撤销,请确认
快速查找节点在 DAG 图上方搜索框输入节点名称适用于大型工作流支持模糊匹配
查看上下游依赖右键节点 → “查看上游” / “查看下游”分析任务影响范围便于调试依赖逻辑
导出/导入工作流”更多操作” → “导出” JSON 文件实现跨环境迁移或备份导入时注意数据源和资源路径一致性

第四章:任务类型详解

4.1 Shell 任务

参数名称语法用途示例注意事项
脚本内容直接输入 Shell 命令或脚本片段执行 Linux 命令或调用本地脚本#!/bin/bashdateecho "Running as ${USER}"使用 #!/bin/bash 明确解释器;支持变量注入
运行标志NORMAL / FORBIDDEN / DISABLED控制任务是否参与调度NORMAL(正常执行)FORBIDDEN 表示跳过,但仍计入 DAG
失败策略继续 / 结束流程定义任务失败后是否中断整个工作流建议关键任务设为”结束流程""继续”可用于非关键通知类任务
超时设置数值 + 单位(分/小时)防止任务无限挂起超时时间:30 分钟超时后自动 kill 进程并标记失败
自定义参数KEY=VALUE 形式添加注入环境变量或脚本参数key: input_path, value: /data/input/${yyyyMMdd}可在脚本中通过 ${input_path} 引用

4.2 SQL 任务(支持多数据源)

参数名称语法用途示例注意事项
数据源类型MySQL / PostgreSQL / Oracle / SQLServer / Hive / Spark 等选择目标数据库类型选择:MySQL需提前在”数据源中心”配置连接信息
数据源下拉选择已注册的数据源指定连接实例datasource_mysql_01测试连通性确保可用
SQL 脚本输入标准 SQL 语句执行查询或更新操作SELECT count(*) FROM logs WHERE dt='${dt}'支持多条语句,分号分隔
查询结果处理不处理 / 写入表 / 发送告警定义 SELECT 结果的后续动作发送告警:将结果通过邮件通知仅 SELECT 语句产生结果
分页记录数整数(如 100)控制返回结果行数100避免大数据量阻塞
自定义参数KEY=VALUE传递动态值dt=$[yyyy-MM-dd]支持日期函数解析
PreparedStatement开启/关闭是否预编译 SQL开启可防止注入UPDATE/DELETE 推荐开启

4.3 Spark 任务

参数名称语法用途示例注意事项
主程序包选择已上传的 JAR 文件(资源中心)指定 Spark 应用主 Jar 包spark-etl-1.0.jar必须先上传至资源中心
主类名com.example.MainClass指定 main 方法所在类com.etl.DataProcessor类路径必须正确
部署模式local / client / clusterSpark 运行模式cluster(生产推荐)cluster 模式 Driver 在集群运行
Spark 版本1.x / 2.x / 3.x匹配集群实际版本3.1.2版本不一致可能导致兼容问题
命令行参数--input hdfs://... --output /out传递给 main 方法的 args--mode=prod --batchId=123多参数空格分隔
Driver/Core 资源driver-memory, executor-cores 等设置资源配额--driver-memory 2g --num-executors 4根据集群资源合理配置
自定义参数spark.conf.key=value添加 SparkConf 配置spark.sql.adaptive.enabled=true高级调优选项
参数名称语法用途示例注意事项
Flink 脚本类型Java/Scala/PYTHON选择程序语言JAVAPYTHON 需启用 PyFlink
程序包上传或选择已有的 JAR/Python 文件指定作业代码flink-job-1.0.jar支持远程 HDFS 路径
主类名org.apache.flink.WordCountJava/Scala 程序入口类com.streaming.RealTimeJob必须是 public class
部署模式local / remote / yarn-per-job / yarn-session运行模式yarn-per-job(推荐)生产环境避免 local 模式
Flink 配置参数parallelism, checkpointInterval 等调优参数-p 4 --checkpointing 60000参考 Flink CLI 参数格式
JobManager 地址host:portremote 模式需指定jobmanager-host:8081仅 remote 模式需要
YARN 队列default / prod_queue指定资源队列prod_queue需 YARN 集群支持队列划分

4.5 Python 任务

参数名称语法用途示例注意事项
Python 脚本直接编写或引用外部 .py 文件执行 Python 逻辑print("Hello")import pandas as pd确保 Worker 节点安装所需库
脚本类型内联脚本 / 脚本文件选择输入方式内联脚本适合简单逻辑文件方式便于复用
Python 命令路径python / python3 / /usr/bin/python3.8指定解释器路径/opt/venv/bin/python若使用虚拟环境需写完整路径
依赖包管理手动安装或打包 venv确保第三方库可用pip install pandas建议将依赖打包进 zip 并上传
命令行参数--arg1 val1 --arg2 val2传参给脚本sys.argv 获取参数使用 argparse 解析更规范
自定义参数KEY=VALUE注入环境变量env: ENV=prod脚本中通过 os.getenv('ENV') 读取

4.6 HTTP 任务

参数名称语法用途示例注意事项
请求方法GET / POST / PUT / DELETE指定 HTTP 动作POST根据接口文档选择
URLhttps://api.example.com/v1/data目标接口地址https://hooks.slack.com/services/...支持变量替换,如 ${token}
请求头Content-Type: application/json添加 HeaderAuthorization: Bearer ${access_token}多个头换行输入
请求体JSON 或 Form-dataPOST/PUT 数据内容{"name": "test"}JSON 需合法格式
超时时间毫秒数(如 5000)设置连接与读取超时3000 ms防止长时间阻塞
成功状态码200 / 2xx / 自定义判断请求是否成功200,201多个码逗号分隔
认证方式NONE / BASIC / BEARER身份验证机制BEARER token=abc123安全敏感信息建议用参数传递

4.7 Sub-Process 任务

参数名称语法用途示例注意事项
子工作流从下拉框选择已有工作流定义调用另一个完整工作流workflow_data_import子流程独立调度与记录
等待子流程是 / 否是否阻塞等待完成是(常用)“否”表示异步触发
传递参数KEY=VALUE向子流程传递上下文parent_id=${processId}子流程需定义对应全局参数
失败策略继续 / 结束流程子流程失败后的处理关键流程建议”结束流程”可实现异常传播
超时控制数值 + 单位防止子流程长期不结束2 小时超时后自动终止子实例

4.8 Condition 分支任务

条件表达式语法用途示例注意事项
SUCCESSnode_name == 'SUCCESS'判断前驱任务是否成功A == 'SUCCESS'支持节点别名
FAILUREnode_name == 'FAILURE'判断是否失败B == 'FAILURE'可触发告警分支
ALL_FAILUREALL_STATUS == 'FAILURE'所有上游都失败才走此分支ALL_STATUS == 'FAILURE'用于兜底处理
ELSEELSE默认分支,当前面都不满足时执行ELSE必须放在最后
多条件组合&& / ||逻辑与或(A=='SUCCESS') || (B=='SUCCESS')括号控制优先级
输出分支Left / Right / Else分支连接不同下游任务Left → 发送成功通知;Right → 发送失败告警DAG 必须显式连线

4.9 Depend 依赖任务

参数名称语法用途示例注意事项
依赖工作流选择项目内的其他工作流定义跨工作流依赖project_etl_daily必须在同一项目
依赖周期TODAY / YESTERDAY / LAST_1_DAYS 等时间维度匹配YESTERDAY用于按天补数场景
依赖状态SUCCESS / ANY要求的状态SUCCESS建议关键依赖设为 SUCCESS
自定义表达式${dag.offset(1)}==SUCCESS高级依赖判断${dag.offset(-1)}==SUCCESSoffset(-1) 表示前一天实例
检查间隔1m / 5m / 10m轮询频率5 分钟频繁检查增加 DB 压力
超时时间最长等待时间防止无限等待2 小时超时后可配置失败策略

4.10 Stored Procedure 存储过程任务

参数名称语法用途示例注意事项
数据源选择支持存储过程的数据库如 MySQL、Oracleds_oracle_proc必须支持 CALL 语法
存储过程名procedure_name(?, ?)调用带参过程calc_monthly_report(?, ?)? 为占位符
参数类型IN / OUT / INOUT定义参数方向IN: date_str, OUT: result_code多数为 IN 类型
参数值具体值或变量传入实际参数$[yyyy-MM-dd], 100支持表达式解析
调用语法{call proc_name(?,?)}标准 JDBC 调用格式{call backup_data(?,?)}必须符合 JDBC 规范
返回处理忽略 / 写日志 / 告警如何处理 OUT 参数写日志:result_code=${result_code}OUT 参数可在后续任务引用

4.11 MR(MapReduce)任务

参数名称语法用途示例注意事项
主类名com.hadoop.WordCountMapReduce 程序入口com.etl.UserBehaviorMR必须打包进 Jar
程序包上传 Hadoop Jar 包包含 Mapper/Reducer 类etl-job-2.0.jar支持 HDFS 路径
Hadoop 配置core-site.xml, hdfs-site.xml指定集群配置文件自动加载 $HADOOP_CONF_DIR确保 Worker 节点配置正确
命令行参数/input /output传递给 ToolRunner 的 args/raw/logs/$[yyyyMMdd] /dw/fact/user路径需存在且可读写
队列名称default / etl_queueYARN 队列分配etl_queue需集群支持资源队列
任务优先级LOW / NORMAL / HIGH调度优先级NORMALHIGH 可能抢占资源
故障重试1~4 次任务失败自动重试次数2 次Hadoop 本身也有重试机制,避免叠加过多

4.12 DataX 数据同步任务

参数名称语法用途示例注意事项
DataX JSON 模板标准 DataX 作业配置定义读写插件与字段映射{ "job": { "content": [...] } }必须符合 DataX Schema
reader.namemysqlreader / hdfsfiler源端读取插件mysqlreader插件名必须正确
writer.namemysqlwriter / hdfswriter目标写入插件hdfswriter注意目标格式(text/orc)
数据源配置jdbcUrl, username, password数据库连接信息jdbc:mysql://host:3306/db建议使用数据源管理功能
字段映射column: ["id", "name"]指定同步字段支持常量和表达式{value: "fixed", type: "string"}
通道数channel: 3并发读写线程数channel: 5根据源库负载调整
脏数据限制record: 100, percentage: 0.05容忍错误记录数record: 10, percentage: 0.1超过则任务失败
Pre/Post SQL执行前置或后置 SQL清理或统计preSql: ["truncate table tmp"]可用于事务控制

第五章:工作流设计与调度机制

5.1 DAG 有向无环图原理

概念名称说明注意事项
DAG(Directed Acyclic Graph)有向无环图,DolphinScheduler 的核心任务编排模型,用节点表示任务,边表示依赖关系不允许存在循环依赖(如 A→B→C→A),否则无法调度
节点(Node)代表一个具体任务,如 Shell、SQL、Spark 等类型每个节点有唯一名称和运行配置
边(Edge)表示任务之间的执行顺序依赖,前驱任务成功后,后继任务才能启动依赖关系是”成功才执行”,失败默认中断流程(可配置)
入度与出度入度:指向该节点的边数;出度:从该节点出发的边数入度为 0 的节点为起始节点,出度为 0 的为终止节点
拓扑排序系统根据 DAG 自动计算任务执行顺序,确保依赖满足排序结果决定任务分发时机
并行分支多个任务无依赖关系时可并行执行,提升效率受 Worker 资源和任务队列限制
子 DAG通过 Sub-Process 任务实现嵌套流程,提高复用性子流程独立调度、记录日志

5.2 节点间的依赖关系设置

依赖类型配置方式用途示例注意事项
顺序依赖鼠标从任务 A 拖线到任务 BB 在 A 成功后执行A → B → C最常见模式
多前驱依赖多个任务连接到同一任务所有前驱必须成功,后继才执行A→C, B→C,则 C 等 A 和 B 都成功用于合并分支
条件分支使用 Condition 任务判断状态根据前驱结果走不同路径A 成功 → 发邮件;A 失败 → 发告警需配合 Condition 节点
手动依赖不连线,通过 Depend 任务跨工作流依赖实现项目间或周期性依赖工作流 B 依赖工作流 A 昨天的成功实例用于跨日调度场景
广播依赖一个任务成功触发多个下游一对多通知或分发A → B, A → C, A → D可并行执行 B/C/D
可选依赖设置”失败继续”策略前驱失败不影响后继执行B 不依赖 A 的状态用于非关键任务

5.3 全局参数与局部参数

参数类型配置位置语法用途示例注意事项
全局参数工作流定义页面 → “全局参数”按钮${param_name}整个工作流共享的变量${bizdate}, ${region}在任务中通过 ${} 引用
局部参数单个任务节点 → “自定义参数”${param_name}仅当前任务使用的变量${inputpath}, ${thread_num}优先级高于全局参数
参数类型 - VALUE手动输入固定值类型选择 VALUE静态配置VALUE: prod不支持函数解析
参数类型 - DATE使用日期函数动态生成类型选择 DATE动态时间戳DATE: $[yyyy-MM-dd]支持偏移,如 $[yyyyMMdd-1] 表示昨天
内置参数系统自动提供${startTime}${endTime}获取运行上下文${scheduleTime}常用于日志标记
参数传递子流程通过 Sub-Process 任务传参KEY=VALUE 形式上下文传递parent_id=${processId}子流程需预先定义同名参数
参数优先级局部 > 全局 > 内置相同名称时覆盖规则若局部和全局都有 ${env},取局部值建议命名区分作用域

5.4 工作流定时调度配置(Cron 表达式)

配置项说明示例注意事项
Cron 表达式格式秒 分 时 日 月 周 年(年可选)0 0 2 * * ?DolphinScheduler 使用 Quartz 格式
秒字段0-590/30 表示每30秒一次通常设为 0
分字段0-590/15 表示每15分钟支持范围(10-20)和通配符(*)
时字段0-232 表示凌晨2点2 表示 02:00
日字段1-31? 表示不指定(与周互斥)避免与周同时指定具体值
月字段1-12 或 JAN-DEC* 表示每月JAN 表示1月
周字段1-7 或 SUN-SAT? 表示不指定(与日互斥)1=SUN, 2=MON…7=SAT
年字段可选,如 2025* 表示每年通常省略
常见表达式每日凌晨2点0 0 2 * * ?推荐用于 ETL 任务
常见表达式每小时整点0 0 * * * ?适用于实时性要求高的任务
常见表达式工作日9点0 0 9 ? * MON-FRIMON-FRI 表示周一至周五
验证工具Web UI 提供”下一次触发时间”预览输入后自动计算避免错误调度
时区问题默认使用服务器时区(建议 UTC+8)避免跨时区混乱生产环境统一设置时区

5.5 手动触发与补数操作

操作类型配置方式用途示例注意事项
手动运行”工作流定义” → 选择工作流 → “上线” → “运行”临时触发一次执行调试或紧急处理必须先上线才能运行
补数(Backfill)“补数”按钮 → 选择日期范围重新执行历史时间段的任务补 2025-09-01 至 2025-09-10 的数据用于修复数据或修复失败任务
补数模式串行 / 并行控制补数实例的执行方式并行:最多同时运行 N 个并行可能压垮数据库
补数参数覆盖可修改全局参数值覆盖原始调度参数${env} 从 prod 改为 test用于测试修复逻辑
强制运行忽略依赖直接启动跳过前驱任务仅用于调试生产慎用
查看补数实例”工作流实例”中筛选”补数”类型监控补数进度支持暂停、停止补数实例独立记录
补数并发控制在”补数设置”中限制并发数防止资源过载最大并发:3根据集群负载调整

5.6 并行执行与串行执行控制

控制方式配置方法说明示例注意事项
DAG 结构控制多任务无依赖则并行系统自动并行调度A → C, B → C,则 A 和 B 并行受 Worker 资源限制
任务组(Task Group)设置任务组名称和并发数限制同类任务并发组名:etl_group,并发:2需提前在系统配置启用
并行分支多个下游任务同时启动提升吞吐A → B, A → C, A → DB/C/D 并行执行
串行化依赖显式添加顺序依赖强制按序执行A → B → C用于资源竞争场景
Worker 资源限制worker.exec.threads 配置控制单 Worker 最大并发默认 100根据 CPU 核数调整
Master 分发策略master.dispatch.task.num控制任务分发频率避免瞬时高峰通常保持默认
失败重试间隔重试间隔时间(分钟)避免密集重试间隔:5 分钟可防止雪崩
手动暂停/恢复在”任务实例”中暂停队列临时控制执行节奏运维窗口期暂停恢复后继续调度

第六章:系统管理与安全

6.1 用户与角色权限管理

概念说明注意事项
用户(User)系统登录实体,有用户名、密码、邮箱、手机号支持 LDAP/AD 集成(需配置)
角色(Role)权限集合,分为系统角色和项目角色系统角色:Admin、普通用户等;项目角色:管理员、开发员、运维员、访客
系统角色 - Admin拥有所有权限,可管理用户、租户、告警、项目等仅分配给运维人员
系统角色 - 普通用户可创建项目、上传资源、定义工作流默认权限较低,需授权访问项目
项目角色 - 管理员可管理项目内所有资源和成员可添加/删除成员
项目角色 - 开发员可创建、编辑工作流定义不能上线或运行
项目角色 - 运维员可上线、运行、补数、查看日志不能修改工作流定义
项目角色 - 访客只读权限,仅能查看适用于审计或监控人员
权限分配用户 → 角色 → 项目一个用户可在多个项目有不同角色
用户锁定连续失败登录多次后自动锁定默认 5 次

6.2 租户管理与资源隔离

配置项说明示例注意事项
租户名称唯一标识符,对应系统用户tenant_ds建议命名清晰
系统用户操作系统级用户名ds_user必须在所有 Worker 节点存在
创建租户安全中心 → 租户管理 → 创建输入租户名和系统用户系统自动创建用户(若未启用自动创建需手动添加)
自动创建用户tenant.auto.create=true自动调用 useradd 创建系统用户需部署机有 sudo 权限
资源隔离不同租户任务以不同系统用户运行ds_user1 和 ds_user2 互不影响防止越权访问文件
文件权限任务日志、临时文件按租户隔离/tmp/dolphinscheduler/tenant_ds/确保目录可写
资源组(Worker Group)租户可绑定到特定 Worker 组default / high_priority实现物理资源隔离
删除租户必须先解除所有任务关联否则无法删除建议归档后删除

6.3 数据源管理(MySQL、PostgreSQL、Hive 等)

数据源类型配置参数示例注意事项
MySQL主机、端口、数据库名、用户名、密码、JDBC URLjdbc:mysql://192.168.1.10:3306/dw建议使用专用账号,限制 IP
PostgreSQL同上,驱动不同jdbc:postgresql://host:5432/db支持 SSL 连接
OracleSID 或 Service Namejdbc:oracle:thin:@host:1521:ORCL注意驱动版本兼容性
SQL Server使用 jtds 或 mssql-jdbcjdbc:sqlserver://host:1433;DatabaseName=test需开启 TCP/IP
Hive支持 HiveServer2jdbc:hive2://host:10000/default需 Kerberos 认证(如启用)
SparkThrift Server 连接jdbc:hive2://host:10001/default实际走 Hive 协议
测试连接点击”测试连接”按钮验证网络和凭证必须成功才能保存
数据源共享可设置为”公共”或”私有”公共:所有项目可用私有仅创建者所在项目可用
密码加密存储时自动加密AES 或 SM3 算法不以明文存储

6.4 告警组与通知渠道配置(Email、SMS、DingTalk、WeChat)

通知类型配置参数示例注意事项
EmailSMTP 服务器、端口、发件人、用户名、密码smtp.gmail.com:587需开启 SMTP 服务
短信(SMS)第三方平台 API 密钥阿里云 SMS、腾讯云 SMS需申请模板和签名
钉钉(DingTalk)Webhook URL + 自定义关键词https://oapi.dingtalk.com/robot/send?access_token=xxx必须设置”加签”或关键词
企业微信(WeChat Work)应用 Webhookhttps://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=xxx支持 Markdown
飞书(Lark)飞书群机器人 Webhookhttps://open.feishu.cn/open-apis/bot/v2/hook/xxx支持富文本
告警组包含多个通知方式的集合group_prod_alert可绑定到工作流或任务
告警模板自定义消息内容任务失败:${taskName} at ${time}支持变量注入
失败告警触发在任务或工作流中启用”失败告警”勾选”失败时通知”可指定告警组
告警去重避免短时间内重复发送配置冷却时间(如 30 分钟)防止告警风暴

6.5 安全配置(HTTPS、访问控制、Token 管理)

配置项配置方式说明注意事项
HTTPS 启用配置 Nginx 或 Tomcat SSL前端反向代理启用 HTTPS推荐生产环境使用
访问控制(IP 白名单)通过 Nginx 或防火墙限制allow 192.168.1.0/24限制 API 和 Web 访问来源
Token 认证用于 API 调用身份验证请求头添加 X-Access-Token: xxxToken 在”安全中心”生成
Token 过期时间默认 30 天可配置 token.expire.time建议定期轮换
用户密码策略最小长度、复杂度、过期时间默认无强制策略可通过 LDAP 统一管理
操作审计日志记录用户关键操作(如删除工作流)日志位于 logs/api-server/用于安全审计
敏感信息加密数据库密码、API Key 等使用 AES 加密存储密钥管理需安全
CSRF 防护启用 anti-forgery token默认开启防止跨站请求伪造

第七章:高可用与监控运维

7.1 Master/Worker 高可用机制

组件实现机制配置要点注意事项
Master 高可用基于 ZooKeeper 选举主节点,多个 Master 实例中仅一个 Active,其余 Standbyconf/master.properties 中配置 zk.quorum=zk1:2181,zk2:2181,zk3:2181所有 Master 节点需能连接 ZooKeeper 集群
Worker 高可用无主从之分,所有 Worker 注册到 ZooKeeper,Master 随机或按负载分发任务conf/worker.properties 中设置 worker.group=default故障 Worker 自动下线,任务由其他 Worker 接管
故障转移(Failover)Master 宕机后,ZooKeeper 触发重新选举,新 Master 恢复调度确保 master.server.max.idle.time 设置合理(默认 30s)网络抖动可能导致误判,建议调大超时时间
任务容错任务失败可配置重试次数,或由其他 Worker 重新执行在任务节点设置”失败重试”策略重试间隔避免过短导致雪崩
多租户隔离不同租户任务在不同系统用户下运行,防止资源争抢租户绑定操作系统用户确保 Worker 节点上用户存在且权限正确
心跳检测Master/Worker 每隔一定时间向 ZooKeeper 发送心跳heartbeat.interval 默认 5s网络延迟过高可能导致假死

7.2 ZooKeeper 在集群中的作用

功能说明配置文件注意事项
Master 选举利用 ZNode 临时节点和 Watcher 机制实现主节点选举conf/master.propertieszk.quorum=...建议部署奇数个节点(3/5/7)
Worker 注册发现Worker 启动时在 /workers 路径创建临时节点conf/worker.propertiesMaster 通过监听该路径感知 Worker 状态
分布式锁协调多个 Master 对任务调度的并发访问内部机制,无需手动配置避免长时间持有锁
配置管理可集中存储部分动态配置(如队列状态)非主要用途,DolphinScheduler 主要用 DB 存储配置建议仍以数据库为准
状态协调记录任务实例状态变更,确保一致性用于 Master 故障恢复时重建上下文日志与 ZK 状态应一致
会话超时sessionTimeout 参数控制连接有效性zookeeper.session.timeout=60000(单位 ms)设置过短易误判宕机,过长恢复慢
监控命令echo stat | nc zk_host 2181命令行工具用于排查连接问题

7.3 日志路径与日志分析技巧

日志类型默认路径内容说明分析技巧注意事项
Master 日志logs/master-server/master-server.log调度决策、任务分发、ZK 连接状态搜索关键词:Scheduling, Dispatch, Failed关注调度延迟和分发失败
Worker 日志logs/worker-server/worker-server.log任务拉取、执行启动、资源分配搜索:Executing task, Process start failed检查脚本路径、权限、环境变量
API Server 日志logs/api-server/api-server.log用户请求、认证、REST 接口调用搜索:HTTP 500, Unauthorized, SQLException排查登录失败或接口错误
任务实例日志logs/task/{taskInstanceId}.log具体任务的标准输出和错误输出搜索:ERROR, Exception, exit codeShell 任务非 0 退出码即失败
前端日志浏览器 F12 ConsoleWeb UI 交互错误、JS 异常查看网络请求是否 404/500常见于跨域或 Token 过期
日志轮转按天分割,保留 30 天使用 logback 配置可通过 logback-spring.xml 修改策略避免磁盘写满
日志级别调整INFO / DEBUG / WARN修改 application-*.ymllogging.levelDEBUG 级别日志量大,仅调试时开启

7.4 系统监控指标(CPU、内存、队列积压等)

指标类别监控项正常范围异常表现建议监控方式
CPU 使用率Master/Worker 进程 CPU 占用< 70%持续 > 90% 可能导致调度延迟Prometheus + Grafana
内存使用JVM Heap 使用(Xmx)< 80%OOM 或频繁 GCjstat -gc pid,或 JMX 导出
线程池积压Master 待处理任务队列长度< 100队列持续增长表示消费不及查看 master.dispatch.task.num 相关日志
ZooKeeper 连接ZK 会话数、延迟延迟 < 100ms超时或断连影响 HAecho stat | nc zk 2181
数据库连接MySQL 活跃连接数< 最大连接数 80%连接耗尽导致 API 失败show processlist;
任务执行延迟从调度时间到实际启动的时间差< 5s显著延迟可能 Worker 资源不足查看任务实例”开始时间”与”计划时间”
磁盘空间日志目录所在分区> 20% 剩余写入失败导致任务异常df -h 定期检查
网络带宽节点间数据传输无持续打满影响大文件同步或日志上传sar -n DEV 1 3

7.5 故障排查常见问题清单

问题现象可能原因排查步骤解决方案
Master 无法启动数据库连接失败、端口占用、ZK 不可达1. 查看 master-server.log
2. telnet 检查 DB/ZK 连通性
3. netstat 检查 5678 端口
修复网络、修改配置、释放端口
Worker 未注册SSH 免密失败、租户用户不存在、ZK 问题1. 查看 worker-server.log
2. 手动 su - tenant_user 测试
3. 检查 sshd 服务
配置免密、创建系统用户、重启 sshd
任务卡在”正在运行”脚本死循环、超时设置过大、进程未上报状态1. 查看任务日志是否有输出
2. ps 查找对应进程
3. kill 后观察是否恢复
设置合理超时、优化脚本逻辑
工作流不触发Cron 表达式错误、未上线、补数冲突1. 检查”下一次执行时间”预览
2. 确认工作流状态为”上线”
3. 查看 master 调度日志
修正 Cron、上线工作流
SQL 任务连接失败数据源配置错误、网络不通、驱动缺失1. 在”数据源中心”测试连接
2. telnet host port
3. 检查 lib 目录是否有驱动 jar
修正 IP/端口、添加 jdbc 驱动
日志无法查看Logger Server 未启动、路径权限不足1. 查看 logger-server.log
2. 检查 logs/task/ 目录权限
启动 logger 服务、chmod 755
Web UI 加载慢网络延迟、API 响应慢、浏览器缓存1. F12 查看 Network 请求耗时
2. 检查 api-server.log 是否有慢查询
优化网络、升级硬件、清理缓存

第八章:API 接口编程与集成

8.1 REST API 基础认证方式(Token)

参数说明示例注意事项
认证方式Bearer Token请求头:X-Access-Token: abcdefghijklmnopqrstuvwxToken 在”安全中心”生成
获取 Token通过用户密码调用登录接口获取POST /login{"userName": "admin", "userPassword": "dolphinscheduler"}返回 JSON 包含 token 字段
Token 有效期默认 30 天可配置 token.expire.time=2592000(秒)过期需重新登录获取
使用方式所有 API 请求必须携带 Tokencurl -H "X-Access-Token: abc..." http://api:12345/projects否则返回 401 Unauthorized
权限控制Token 绑定用户角色,决定可访问资源admin 的 token 权限最高避免泄露
多租户支持Token 自动关联用户所属租户无需额外传参任务提交自动使用用户默认租户
安全建议HTTPS 传输、定期轮换、最小权限原则生产环境禁用明文 HTTP可结合 LDAP 统一认证

8.2 项目管理相关 API

API 接口请求方式参数用途示例
创建项目POST /projectsprojectName, description新建一个项目容器POST /projects?projectName=my_project&description=ETL
查询项目列表GET /projectspageSize, pageNo获取用户有权限的项目GET /projects?pageSize=10&pageNo=1
删除项目DELETE /projects/{projectId}projectId删除指定项目(需管理员权限)DELETE /projects/5
项目详情GET /projects/{projectId}projectId查看项目基本信息GET /projects/5
添加项目成员POST /project-usersprojectId, userId, perm将用户加入项目并赋权POST /project-users?projectId=5&userId=10&perm=3
移除成员DELETE /project-users/{relationId}relationId解除用户与项目的关联DELETE /project-users/20

8.3 工作流定义操作 API

API 接口请求方式参数用途示例
创建工作流POST /projects/{projectId}/workflow-definitionname, json (DAG 结构)提交一个新的 DAG 定义POST /projects/5/workflow-definition?name=test_wf&json={...}
上线工作流PUT /projects/{projectId}/wf-instance/publishworkflowDefinitionCode, state=publish将工作流设为可调度状态PUT /projects/5/wf-instance/publish?workflowDefinitionCode=100&state=publish
下线工作流PUT ... state=unpublish同上停止调度,禁止手动运行state=unpublish
查询工作流列表GET /projects/{pid}/workflow-definitionsearchVal, pageSize, pageNo模糊查找工作流GET /projects/5/workflow-definition?searchVal=etl
获取工作流详情GET /projects/{pid}/workflow-definition/{code}code查看 DAG 结构和参数GET /projects/5/workflow-definition/100
导出工作流GET /projects/{pid}/export-workflow-defworkflowDefinitionCode下载 JSON 文件备份GET /projects/5/export-workflow-def?workflowDefinitionCode=100
导入工作流POST /projects/{pid}/import-workflow-deffile (JSON)从文件恢复工作流POST /projects/5/import-workflow-def

8.4 工作流实例控制 API

API 接口请求方式参数用途示例
手动启动工作流POST /projects/{pid}/executors/startworkflowDefinitionCode触发一次执行POST /projects/5/executors/start?workflowDefinitionCode=100
停止工作流实例POST /projects/{pid}/executors/stopworkflowInstanceId终止正在运行的实例POST /projects/5/executors/stop?workflowInstanceId=200
补数操作POST /projects/{pid}/batch-executionmode=backfill, start, end, codes批量重跑历史实例mode=backfill&start=20250901&end=20250910&codes=100
查询实例列表GET /projects/{pid}/executors/querystartDate, endDate, state按条件筛选实例GET /projects/5/executors/query?state=SUCCESS&startDate=2025-09-01
实例详情GET /projects/{pid}/executors/{instanceId}instanceId查看 DAG 执行状态和耗时GET /projects/5/executors/300
查看实例日志GET /projects/{pid}/executors/logtaskInstanceId, skipLineNum, limit分页获取任务日志GET /projects/5/executors/log?taskInstanceId=400&skipLineNum=0&limit=1000

8.5 任务实例查询 API

API 接口请求方式参数用途示例
查询任务实例GET /projects/{pid}/task-instancesworkflowInstanceId, taskName, state查找特定任务实例GET /projects/5/task-instances?workflowInstanceId=300&state=RUNNING
任务实例详情GET /tasks/{taskInstanceId}taskInstanceId获取任务配置和执行信息GET /tasks/400
重跑任务POST /projects/{pid}/task-instances/{id}/repeatid, executeType重新执行失败任务POST /projects/5/task-instances/400/repeat?executeType=REPEAT_RUNNING
停止单个任务POST /projects/{pid}/task-instances/{id}/stopid终止正在运行的任务POST /projects/5/task-instances/400/stop
替代运行POST .../replace-by-nodeid, node用指定节点替代当前任务继续调试图形分支时使用

8.6 数据源操作 API

API 接口请求方式参数用途示例
创建数据源POST /datasource/createname, type, connectionParams(json)添加新的数据库连接POST /datasource/create?name=mysql_prod&type=MYSQL&connectionParams={"host":"h","port":3306,"database":"db"}
查询数据源列表GET /datasource/listtype, searchVal, pageNo, pageSize获取可用数据源GET /datasource/list?type=MYSQL&pageNo=1&pageSize=10
测试连接POST /datasource/connecttype, connectionParams验证连接可用性用于前端”测试连接”功能
更新数据源PUT /datasource/updateid, name, connectionParams修改已有数据源配置PUT /datasource/update?id=10&name=new_name
删除数据源DELETE /datasource/deleteid移除数据源(需无任务引用)DELETE /datasource/delete?id=10
获取数据源详情GET /datasource/{id}id查看具体配置信息GET /datasource/10

8.7 用户与权限管理 API

API 接口请求方式参数用途示例
创建用户POST /usersuserName, userPassword, email, phone新增系统用户POST /users?userName=john&userPassword=123&email=john@company.com
查询用户列表GET /userssearchVal, pageNo, pageSize查找用户GET /users?searchVal=admin&pageNo=1
更新用户信息PUT /users/{userId}email, phone, state修改邮箱、手机号、启用/禁用PUT /users/5?email=new@com.com&state=1
重置密码PUT /users/{userId}/passwordoldPassword, newPassword用户自行修改密码需提供旧密码验证
获取当前用户GET /users/get-user-info(无)获取 Token 对应的用户信息GET /users/get-user-info
分配角色POST /user/role/saveuserId, roleId将角色赋予用户POST /user/role/save?userId=5&roleId=2

8.8 告警组与通知 API

API 接口请求方式参数用途示例
创建告警组POST /alert-group/creategroupName, alertInstanceIds定义通知接收组POST /alert-group/create?groupName=prod_alert&alertInstanceIds=1,2,3
查询告警组GET /alert-group/listpageNo, pageSize获取所有告警组GET /alert-group/list?pageNo=1&pageSize=10
更新告警组PUT /alert-group/updateid, groupName, alertInstanceIds修改组内通知方式PUT /alert-group/update?id=5&groupName=new_grp
删除告警组DELETE /alert-group/deleteid移除告警组DELETE /alert-group/delete?id=5
创建告警实例POST /alert-plugin-instance/createpluginType, instanceParams(json)添加钉钉、邮件等通知渠道pluginType=EMAIL&instanceParams={"smtpHost":"smtp.ex.com"}
查询告警实例GET /alert-plugin-instance/listpluginType获取已配置的通知方式GET /alert-plugin-instance/list?pluginType=DINGTALK

第九章:高级特性与最佳实践

9.1 参数传递与运行时变量解析

参数类型语法格式解析时机示例注意事项
全局参数${param_name}工作流启动时解析${biz_date}, ${env}可在任务中直接使用
局部参数${param_name}任务执行前解析${input_path=/data/in}优先级高于全局参数
内置系统参数${system_param}调度系统自动注入${scheduleTime}, ${startTime}, ${endTime}常用于日志标记和分区写入
日期函数表达式$[yyyy-MM-dd]运行时动态计算$[yyyy-MM-dd-7] 表示 7 天前支持 +/- 偏移
日期偏移$[yyyyMMdd-1]运行时动态计算$[yyyyMMdd-1] 表示昨天
Shell 脚本传参使用 param=xxx 形式传递子流程或命令行调用时sh etl.sh date=${biz_date}需脚本内接收 $1
上游任务输出捕获通过 set variable=value 输出任务结束时写入上下文在 Shell 中 echo "SET_OUTPUT:region=shanghai"下游用 ${region} 引用
动态 SQL 参数在 SQL 节点中使用 ${}执行前替换INSERT INTO t VALUES('${user}', ${id});防止 SQL 注入需校验输入
子流程参数映射Sub-Process 任务中设置 KEY=VALUE子流程启动前绑定parent_id=${processId}子流程需预定义同名参数

9.2 工作流版本管理策略

策略说明实现方式优点缺点
手动导出备份定期将工作流导出为 JSON 文件Web UI → 导出功能简单直观,便于归档易遗漏,无差异对比
Git 版本控制将导出的 JSON 提交到 Git 仓库使用 CI/CD 脚本自动提交支持 diff、回滚、分支管理需集成外部系统
命名版本号在工作流名称后加 -v1, -v2etl_user_data-v3快速识别版本不支持自动切换
灰度上线新版本先复制为测试工作流,验证后再替换主流程复制 → 修改 → 上线 → 切流量降低生产风险需人工操作
API 自动化管理使用 /export-workflow-def + /import-workflow-def 接口脚本定期备份或发布可集成进 DevOps 流程需编写维护脚本
版本快照DolphinScheduler 社区版无内置版本,企业版可能支持建议自行实现快照机制
回滚机制从 Git 或备份文件重新导入旧版本curl -X POST ... /import故障时快速恢复需确保参数一致性

9.3 失败重试机制与容错设计

配置项说明示例值注意事项
重试次数任务失败后自动重试的次数0, 1, 3, 5I/O 类任务建议设 2~3 次
重试间隔每次重试之间的等待时间(分钟)5, 10, 30避免密集重试压垮服务
重试策略条件触发重试失败重试 / 异常重试可结合告警通知
容错任务设置”失败继续”标志在任务节点勾选”失败继续”用于非关键路径任务
条件分支容错使用 Condition 节点判断状态A 成功 → B;A 失败 → C(补偿任务)实现事务性语义
超时中断设置任务最大执行时间3600 秒(1小时)防止长时间卡住资源
补偿任务失败后执行清理或通知发送告警、回滚数据提高系统健壮性
断路器模式(需自定义)连续失败 N 次暂停调度结合外部监控系统实现防止雪崩效应

9.4 跨项目任务调用与共享

方式实现机制配置方法适用场景注意事项
Depend 依赖任务跨项目依赖上一个周期的成功实例添加 Depend 节点,选择目标项目和工作流项目间 ETL 依赖仅支持按调度周期依赖
Sub-Process 调用在当前工作流中嵌套执行其他项目的子流程添加 Sub-Process 节点,选择目标项目和工作流复用通用逻辑(如清洗)目标工作流必须上线
公共数据源创建”公共”类型的数据源数据源中心 → 创建 → 设为”公共”多项目共享数据库连接需统一权限管理
资源文件共享上传 UDF 或脚本为”公共”资源资源中心 → 上传 → 公共资源共享 Hive UDF、Shell 工具需注意版本兼容
API 远程触发使用 REST API 启动其他项目的工作流POST /projects/{pid}/executors/start动态触发、条件调用需 Token 认证
全局参数传递通过 Sub-Process 显式传递参数KEY=VALUE 映射上下文传递(如日期)接收方需定义同名参数
权限控制用户需有目标项目的访问权限在目标项目中添加该用户为成员否则无法查看或调用建议最小权限原则

9.5 性能调优建议(线程池、任务队列等)

调优项配置参数推荐值作用注意事项
Master 分发线程数master.dispatch.task.num50 ~ 200控制任务分发速度过大会导致 DB 压力高
Worker 执行线程数worker.exec.threadsCPU 核数 × 2并行执行本地任务避免过多线程争抢资源
ZooKeeper 会话超时zookeeper.session.timeout60000 ms(1分钟)防止网络抖动误判宕机过短易频繁切换主节点
任务心跳间隔task.executor.heartbeat.interval10s保持任务活跃状态上报过长可能导致假死
日志批量刷盘logback.appender.file.bufferSize8KB ~ 64KB提升日志写入性能生产环境建议开启缓冲
数据库连接池spring.datasource.hikari.maximum-pool-size20 ~ 50提高并发查询能力需匹配 MySQL max_connections
缓存启用spring.redis.open=truetrue加速元数据读取建议搭配 Redis 使用
JVM 堆内存-Xms4g -Xmx4g根据机器内存设置避免频繁 GCMaster/Worker 建议独立部署
磁盘 IO 优化使用 SSD 存储日志和临时文件减少 I/O 等待特别是 Spark/Hive 任务

9.6 与 CI/CD 集成自动化发布

阶段工具/方式实现方式示例注意事项
代码管理Git(GitHub/GitLab)将工作流 JSON 文件纳入版本控制etl_workflow_v1.json建议按项目/模块分类存储
构建触发Jenkins / GitLab CI监听 Git Push 事件webhook 触发构建需配置安全令牌
自动化测试Shell 脚本 / Python验证工作流结构合法性JSON Schema 校验可模拟参数运行
自动发布REST API 调用使用 /import-workflow-def 接口curl -X POST ... -F file=@wf.json需携带有效 Token
环境隔离多套 DolphinScheduler 环境dev → test → prod不同环境对应不同集群配置参数差异化(如 ${env}
回滚机制Git revert + 重新导入回退到上一版本并发布git revert HEAD && make deploy需记录发布历史
审批流程Jenkins Manual Step / MR Review人工确认后才发布到生产”是否继续部署生产?“关键任务建议人工审核
发布报告邮件 / 钉钉通知发送成功/失败消息”工作流 etl_user 已部署至 PROD”包含版本号和变更内容

第十章:扩展开发与源码解析(可选)

10.1 自定义任务类型开发

步骤说明关键类/接口注意事项
1. 继承 TaskExecutionContext封装任务运行上下文org.apache.dolphinscheduler.plugin.task.api.TaskExecutionContext获取参数、环境信息
2. 实现 TaskPluginDelegate定义任务执行入口org.apache.dolphinscheduler.plugin.task.api.TaskPluginDelegate核心执行逻辑
3. 创建 TaskDefinition定义任务配置模型TaskDefinition包含参数、资源、超时等
4. 编写 TaskProcessor处理任务生命周期TaskProcessor启动、监控、终止
5. 打包为 JAR 插件放入 plugins/task/ 目录dolphinscheduler-task-custom-1.0.jar需符合 SPI 规范
6. 重启 Worker加载新任务类型Master 不需要重启
7. Web UI 支持前端添加任务图标和表单Vue 组件开发非必须,可用 API 调用
示例任务Python 脚本任务、Flink SQL 任务、Kafka 生产者任务等可复用现有进程执行器

10.2 插件机制与扩展点说明

扩展点用途实现方式示例
任务插件(Task Plugin)支持新类型任务(如 Flink, Kubernetes)实现 TaskPluginDelegate 接口自定义 AI 训练任务
告警插件(Alert Plugin)新增通知渠道(如 Feishu, Slack)实现 AlertChannelPlugin企业微信机器人
数据源插件支持新型数据库连接实现 DataSourceChannelDoris、StarRocks
资源存储插件切换资源存储后端(HDFS, S3)实现 ResourceStoragePluginAWS S3 存储脚本
认证插件集成 LDAP、OAuth2、JWT实现 AuthenticationProvider单点登录(SSO)
日志插件自定义日志收集方式(ELK, Loki)实现 LoggerPlugin推送到 Kafka 进行分析
SPI 机制Java Service Provider InterfaceMETA-INF/services/ 下注册实现类工业标准插件机制
热加载插件放入目录后自动识别无需重启核心服务仅 Worker 需重启

10.3 核心模块源码结构分析

模块主要功能关键包路径说明
API Server提供 REST 接口、用户认证、元数据管理org.apache.dolphinscheduler.apiSpring Boot 应用,处理 HTTP 请求
Master Server调度决策、DAG 解析、任务分发org.apache.dolphinscheduler.server.master基于 Quartz 和 ZooKeeper 实现 HA
Worker Server任务执行、资源调度、状态上报org.apache.dolphinscheduler.server.worker实际 fork 子进程运行 Shell/SQL 等
Common公共工具类、实体、常量org.apache.dolphinscheduler.common跨模块共享代码
Alert告警发送、渠道管理、事件监听org.apache.dolphinscheduler.alert支持邮件、短信、IM
Plugin插件框架与各类插件实现org.apache.dolphinscheduler.plugin任务、告警、数据源等扩展
DAO数据访问层,操作 PostgreSQL/MySQLorg.apache.dolphinscheduler.daoMyBatis 实现
RPC节点间通信协议(Netty)org.apache.dolphinscheduler.remoteMaster 与 Worker 通信基础

10.4 编译与调试环境搭建

步骤操作命令/工具注意事项
1. 获取源码Clone 官方仓库git clone https://github.com/apache/dolphinscheduler.git建议使用最新稳定分支
2. 安装 JDKJava 8 或 11java -versionMaven 编译依赖
3. 安装 Maven构建工具mvn -v用于打包
4. 编译项目执行构建mvn clean install -Dmaven.test.skip=true跳过测试加快编译
5. 导入 IDEIntelliJ IDEA / EclipseFile → Open → pom.xml等待依赖下载完成
6. 配置数据库初始化 MySQL/PostgreSQLsh script/create-dolphinscheduler.sql修改 application-dao.yml 连接信息
7. 启动服务分别运行各模块./bin/start.sh master-server./bin/start.sh worker-server可单节点调试
8. 调试模式添加远程调试参数-agentlib:jdwp=transport=dt_socket,server=y,suspend=n,address=5005使用 IDEA 远程调试
9. 日志定位查看 logs/ 目录tail -f logs/master-server.log快速发现问题
10. 单元测试运行测试用例mvn test验证修改正确性