DolphinScheduler 任务调度实战:从最小工作流到生产三层防线的完整路径
DolphinScheduler 任务调度实战从最小工作流到生产三层防线的完整路径【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler周一下午 2 点上游物流轨迹增量表延迟了 40 分钟下游三张库存汇总表全部空跑。这种上游慢、下游空跑我们用 cron 补了好几年脚本链越补越乱谁都不敢动。最后把整条管道挪进了 DolphinScheduler 的分布式任务调度工作流编排里把依赖、重试、通知全部声明成 DAG。下面用同一条物流轨迹 → 库存预警的数据流转从先跑通讲到稳定跑在生产。先跑通4 个任务的最小工作流与三个常见坑起步阶段别追功能先让整条链路转起来。下面这个 4 节点管道就是我们当时最先跑通的那版抽取 → 聚合 → 校验 → 通知复制到新项目里就能用。{ name: logistics_stock_pipeline, description: 物流轨迹到库存预警的日管道, tasks: [ { name: tracking_extract, type: DATAX, params: { dsType: MYSQL, dataSource: 1, dtType: HIVE, dataTarget: 2, sql: SELECT tracking_no, status, update_time FROM t_tracking WHERE update_time $[yyyy-mm-dd 00:00:00], targetTable: ods.t_tracking_incr, preStatements: [TRUNCATE TABLE ods.t_tracking_incr], jobChannel: 3, jobSpeedByte: 10485760, xms: 2, xmx: 2, batchSize: 2048 } }, { name: stock_spark_agg, type: SPARK, dependsOn: [tracking_extract], params: { programType: SQL, rawScript: INSERT OVERWRITE TABLE dws_stock_snapshot PARTITION(dt$[yyyy-mm-dd-1]) SELECT warehouse, SUM(1) AS cnt FROM ods.t_tracking_incr GROUP BY warehouse } }, { name: quality_gate, type: SQL, dependsOn: [stock_spark_agg], failRetryTimes: 2, failRetryInterval: 2, timeoutFlag: true, timeout: 30, params: { rawScript: SELECT IF(COUNT(1) 100000, 0, 1) FROM dws_stock_snapshot WHERE dt$[yyyy-mm-dd-1] } }, { name: notify_downstream, type: HTTP, dependsOn: [quality_gate], params: { httpMethod: POST, url: http://inventory-svc.internal/api/v1/stock-snapshot/ready, checkCondition: { responseCode: 200, responseContent: ok }, connectionTimeout: 5000, socketTimeout: 5000 } } ] }逐行说明几个关键字段dependsOn声明依赖因为 A 不先成功B 就不会被派发这是替代脚本里 sleep 等上游的根本手段两个数据源 id 必须先在数据源菜单里创建好任务里只引用 id不写密码preStatements里的 TRUNCATE 让重跑幂等否则失败重试会写出重复数据checkCondition让 HTTP 任务按响应码判断成败而不是发出去就算成功。依赖拓扑画出来是这样的一屏能看清整条链新手第一周必踩的坑我们当时的处理动作整理成这张表照着改现象根因一行修复任务状态卡住不动mainJar 没传到资源中心资源中心上传 jar 后重选日期变量原样输出写成 shell 风格的 ${...}改成 $[yyyy-mm-dd-1]DataX 报 dataSource 错数据源被删id 失效重建数据源并选新 id更多任务类型的字段可以参考 docs/docs/zh/guide/ 里的任务文档各插件的参数定义也都能直接在 dolphinscheduler-task-plugin/ 目录下的*Parameters.java里对到源码。跟着数据走接入、计算、校验、交付四段配置按组件分章节容易变成名词堆砌按数据经过的位置来讲才贴手。下面四节各自只回答一个问题这段数据流进来时调度侧要配什么。接入三类数据源三种配置逻辑同一句抽昨天的增量源不同配置逻辑完全不同。如果是关系库DataX 用数据源 id 拼出 JSON连接串不进工作流定义如果是消息队列消费逻辑写进 Flink 或 shell调度侧只管拉起如果是对象存储就是命令式下载靠文件名做增量。下面这段 DataX 抽取 MySQL 到 Hive是三类里最典型的{ customConfig: 0, dsType: MYSQL, dataSource: 1, dtType: HIVE, dataTarget: 3, sql: SELECT tracking_no, carrier, status FROM t_tracking WHERE dt$[yyyy-mm-dd-1], targetTable: ods.t_tracking_incr, preStatements: [TRUNCATE TABLE ods.t_tracking_incr], jobChannel: 5, jobSpeedByte: 10485760, xms: 4, xmx: 4 }把jobChannel和jobSpeedByte按你源库的承受能力调我们第一版 channel 拉到 10直接把 MySQL 主库的连接池打满了。三类接入路径的差异对照接入路径任务类型关键参数限速字段增量方式关系库DATAXdataSourcejobSpeedBytewhere splitPk消息队列FLINKtopic 放 mainArgsparallelism指定起始 offset对象存储SHELLbucket / path无控并发文件名按日分区⚠️customConfig一旦切成 1手写 DataX json连接串和密码就进了工作流定义等于把密码挂在了界面上。默认保持 0一律走数据源引用。计算批和流的配置分叉在三处批任务和流任务在调度侧的分歧本质是跑完就结束和常驻不退出的区别。如果是 T1 批聚合选 SPARK cluster 模式 超时 重试因为批是无状态的杀掉重跑没有代价如果是实时轨迹处理选 FLINK 并且不设任务超时因为流作业必须常驻状态恢复靠 Flink 自己的 checkpoint不靠调度器重试。批任务这边资源参数给到能直接提交的粒度{ name: stock_batch_agg, type: SPARK, timeoutFlag: true, timeout: 120, failRetryTimes: 3, failRetryInterval: 5, params: { programType: SCALA, mainClass: com.example.stock.StockAgg, mainJar: { id: 101, name: warehouse-agg-1.0.jar }, deployMode: cluster, master: yarn, driverCores: 2, driverMemory: 2G, numExecutors: 10, executorCores: 4, executorMemory: 8G, yarnQueue: etl, appName: stock_agg_daily } }改yarnQueue把批作业钉进独立队列这是和流任务做资源隔离的第一道闸timeout单位是分钟120 分钟是我们按历史 P95 执行时间翻倍估的。流任务这边配置上最反直觉的一点是什么都不设{ name: tracking_realtime, type: FLINK, params: { programType: JAVA, mainClass: com.example.tracking.RealtimeProcessor, mainJar: { id: 202, name: tracking-rt-2.1.jar }, deployMode: yarn-cluster, taskManager: 4, taskManagerMemory: 4G, jobManagerMemory: 2G, parallelism: 8, slot: 1, yarnQueue: realtime } }注意它没有timeoutFlag和failRetryTimes给流作业设超时等于定时杀作业设重试等于每次失败都从头消费。改yarnQueue指向实时专用队列别让批作业把它的资源吃光。模型训练也挂在同一个 DAG 里。我们给输送线做故障预测训练任务就是一个 MLFLOW 节点排在日特征产出之后{ name: fault_model_train, type: MLFLOW, params: { mlflowTaskType: PROJECTS, mlflowJobType: BASIC_ALGORITHM, algorithm: lightgbm, experimentName: conveyor_fault_predict, dataPath: /data/sensor/conveyor.csv, searchParams: {\num_leaves\: [20, 40], \learning_rate\: [0.01, 0.05]}, mlflowTrackingUri: http://mlflow-tracker.prod:5000, modelName: conveyor_fault_v3 } }改dataPath和searchParams就能复用到别的预测场景训练完的模型落在 tracking server 里部署是独立的一个任务不在同一节点里混做。校验质量检查的三个位置与重试熔断校验任务放哪里决定了它保护的是什么。前置门禁放在抽取之后检查源数据到了没有适合上游不保证准时的场景旁路采样和计算任务并行按比例抽样验证适合链路长、不能因检查阻塞主流程的场景后置兜底放在计算之后、交付之前做全量核对适合数据最终要对外交付的环节。我们这条管道用的是后置兜底配置如下重试加熔断都在这个节点上{ name: quality_gate, type: SQL, failRetryTimes: 2, failRetryInterval: 2, timeoutFlag: true, timeout: 30, timeoutNotifyStrategy: WARN, params: { rawScript: SELECT IF(COUNT(1) 100000, 0, 1) FROM dws_stock_snapshot WHERE dt$[yyyy-mm-dd-1] } }熔断不是靠一个开关实现的而是三件事的合力failRetryTimes限定最多重试 2 次避免把坏数据反复洗一遍timeout30 分钟防止门禁空转失败后整个工作流实例置失败notify_downstream因为dependsOn不会执行——下游拿不到一个成功信号这就是熔断的效果。改COUNT(1) 100000这个阈值按你业务量的日 P1 设别拍脑袋。⚠️ 一个永远不会失败的检查等于没有检查。检查 SQL 必须写出会失败的出口阈值用历史数据的分布推别用 0。交付成功走通知失败走升级数据落盘后管道并没有结束交付段要回答两个问题下游怎么知道数据好了失败谁来负责。成功分支交给 DAG 里最后一个 HTTP 任务上面notify_downstream下游服务收到回调才拉数这样数据就绪有明确时间点失败分支不靠 DAG 表达靠工作流级的告警配置兜底{ name: logistics_stock_pipeline, alertConfig: { alertGroupIds: [2], warningType: FAILED } }把alertGroupIds指向你的告警组warningType改成ALL就能连成功也通知。双分支的完整链路是这样的分支触发条件通道动作成功管道跑完DAG 内 HTTP 任务下游自动拉数失败重试后仍失败告警组 FAILED钉钉 值班30 分钟无响应电话持续连续 2 天失败电话人工介入排查数据口径⚠️ 告警组里只放会被叫醒的人。我们把整个部门拉进钉钉组的第一周告警全员已读等于没有告警。生产环境的三层防线Master 凌晨挂了任务会怎样高可用副本规划如果 Master 节点在凌晨挂了任务会怎样答案先说结论已经派给 Worker 的任务不受影响还在 Master 手里排队的任务会被存活节点接管。因为 Master 是无状态多副本都在注册中心挂心跳任务派发以数据库为准坏节点的心跳停了存活 Master 就重新接管它名下的调度职责——前提是你真的部署了多副本并且线程参数按并发量估过。组件推荐副本JVM 内存关键参数Master3-Xms2g -Xmx2gexec-threads200Master3-Xms2g -Xmx2gdispatch-task-num5Worker4~8-Xms4g -Xmx4gexec-threads100API2-Xms1g -Xmx1g无特殊要求Alert2-Xms512m -Xmx512m无特殊要求整体拓扑按这个来部署注册中心自身也要 3 节点以上否则它挂了高可用就断在最后一环✅ 我们上线前在预发环境 kill 过一次 Master 进程做演练队列里的任务约 1 分钟后被接管这个体感时间值得你自己测一遍。队列正在堆积你怎么知道指标与告警规则链路你怎么知道队列正在堆积而不是第二天早上看报表才发现靠的是把排队这件事做成一条自动链路组件自带指标端点Prometheus 定期抓取Grafana 画成面板规则超阈值就推给告警通道。面板上盯这四个数就够了阈值是我们跑了一个月取的分位数指标含义参考阈值Master 待派发队列长度排队等 Worker 的任务持续 5 分钟 100Worker 运行中任务数本机并发水位持续 10 分钟 100任务实例失败数近 1 小时失败量 10 个/小时注册中心掉线节点数Master/Worker 离线 0 立即第二条最容易踩Worker 运行数长期贴着exec-threads上限说明是容量问题不是任务问题这时候优化单个任务是没用的加机器才有用。上周的配置变更是不是故障根源备份、版本化与灰度上周的配置变更是不是就是那次故障的根源如果配置散落在各台机器的本地文件里这个问题的答案永远是不知道。我们后来把配置收敛进一个 Git 仓库每次变更一个 commitcommit 信息里带故障单号出事后先翻 diff 再找任务/config-repo ├── production/ │ ├── application-prod.yaml │ ├── datasource.properties │ └── registry.properties ├── staging/ └── change-log/ # 每次变更一个 commit附故障单号改完这个目录结构就能落地版本化注意production下只放差异项公共项放基线diff 才看得懂。元数据库的备份脚本15 行以内日全量加 30 天滚动#!/bin/bash # 元数据日全量保留 30 天 BACKUP_DIR/backup/dolphinscheduler TS$(date %Y%m%d_%H%M%S) mysqldump -h ds-mysql.prod -u ds_backup -p$(cat /etc/ds/backup.pass) \ --single-transaction --routines --triggers \ dolphinscheduler $BACKUP_DIR/full_$TS.sql find $BACKUP_DIR -name full_*.sql -mtime 30 -delete把ds-mysql.prod和保留天数换成你的环境即可密码不要写死在脚本里。灰度发布就三步新配置先发到一台 Worker观察 24 小时派发与失败率无异常扩到一半全部推平后旧配置的回滚包保留一个完整发布周期。我们栽过一次推平当天就发现新参数把超时单位看错了回滚包还在5 分钟退回。回滚包过期再出事就是事故了。收工前自查清单依赖是否全部在 DAG 里声明没有残留的脚本等脚本每个任务是否都有重试 超时没有无限等待质量门禁是否写得出一个会失败的出口失败告警是否区分了通知与升级两级配置变更是否能追溯到一个 commit 并回滚下一步可以把 Worker 池按任务类型做静态分区让 IO 密集的 DataX 和 CPU 密集的 Spark 跑在不同节点组上从根上消除互相挤占。【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考