ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

DolphinScheduler 落地手册:四套生产工作流的编排与调优清单

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凌晨三点上游宽表还没跑出来下游三十多个任务在队列里干等值班的人盯着告警群心里没底。这类依赖没到、任务空转的锅往往不在任务本身而在编排层。DolphinScheduler 把 DAG 工作流编排、任务级依赖、失败重试和告警通知收敛到一处让你一次性把谁先谁后、失败了怎么办、通知到谁说清楚。从离线数仓到模型上线四套跑在生产上的工作流下面四套编排覆盖了数据团队最常碰的四个诉求。它们的共同点不是任务本身而是把依赖、资源、校验、通知这些容易失控的东西从脚本里搬进工作流定义。任务类型SPARK、FLINK、DATAX、PYTHON、MLFLOW都来自任务插件目录按需启用即可。T1 离线数仓Spark 批处理怎么配资源最常见的一条线每天凌晨把 ODS 洗到 DWS产出第二天八点前要用的报表底表。你的目标不是跑 Spark而是保证这张表第二天一定是干净的。编排三段先用一个 SHELL 把当天分区拉到工作目录SPARK 任务做真正的计算写表最后挂一个 PYTHON 节点做质量校验。校验节点是守门员卡住整条链。Spark 的资源参数别凭感觉。driverCores 给 2、driverMemory 给 2G 就够driver 主要做调度和元数据、不扛数据真正吃资源的是 executornumExecutors 配 10、executorCores 4、executorMemory 8G是按单分区数据量反推出来的。yarnQueue 单独指一个队列避免和在线服务抢资源。deployMode 用 cluster让 driver 也进 YARNMaster 重启不丢任务。{ name: dws_user_behavior, type: SPARK, dependsOn: [pull_partition], params: { programType: SCALA, mainClass: com.example.DwsProcessor, mainJar: {id: 123, name: dws-processor-1.0.0.jar}, deployMode: cluster, master: yarn, yarnQueue: dw_batch, driverCores: 2, driverMemory: 2G, numExecutors: 10, executorCores: 4, executorMemory: 8G, mainArgs: --input /raw/${system.biz.date} --output /dws/${system.biz.date} } }怎么验证跑对了质量校验节点查四件事总行数、去重用户数、事件时间的最小最大值。行数对上昨天量级、时间范围不穿越到明天才算过不过就失败并触发重试加告警。这一步别省Spark 自己不会告诉你数据错了但没报错。实时风控Flink 流处理 异常告警分支上一类是一天一次这一类是一秒一次。交易流水从 Kafka 进来指标要实时算异常得秒级推送晚一分钟风控就废了。编排上Flink 任务作为常驻节点消费 Kafka结果分两条路正常指标写监控库命中风控规则的异常走一条独立告警分支。把告警从主链路拆出来是为了让发通知慢反压不到算指标。Flink 参数和 Spark 不是一路的。parallelism 按 Kafka 分区数对齐这里给 8jobManagerMemory 给 2G 就行JM 只做调度和 checkpoint 协调taskManagerMemory 4G、taskManagerSlots 4一个 TM 起 4 个 slot 摊薄开销。deployMode 仍是 cluster作业托管给集群Worker 挂了能从 checkpoint 恢复。怎么验证常驻任务没有跑完一说只能靠对账和回放。拿当天几笔已知异常交易回放到测试 topic确认告警分支真推到了群里再核对监控库里的实时指标和离线批算的同口径结果误差在可接受范围内。数据同步管道MySQL 增量搬到 Hive数仓的数据从哪来很大一部分是从业务库搬过来的。典型诉求业务库的某张表每天增量同步到 Hive第二天离线链路直接消费。DataX 任务在这里很合适。reader 用 mysqlreader靠 splitPk 按主键切分、多 channel 并行拉取速度快且对源库压力可控writer 用 hdfswriter落成 text 文件直接进 Hive 外表路径省一次格式转换。channel 给 3源库扛得住再往上加别一上来就拉满。{ job: { content: [ { reader: { name: mysqlreader, parameter: { column: [id, name, age, create_time], splitPk: id, username: ${mysql_user}, password: ${mysql_pass}, connection: [{ jdbcUrl: [jdbc:mysql://biz-host:3306/biz_db], table: [user_table] }] } }, writer: { name: hdfswriter, parameter: { defaultFS: hdfs://cluster:8020, fileType: text, path: /warehouse/user_db.db/user_table/dt${system.biz.date}, fileName: user_sync, writeMode: append, fieldDelimiter: \t, column: [ {name: id, type: BIGINT}, {name: name, type: STRING}, {name: age, type: INT}, {name: create_time, type: TIMESTAMP} ] } } } ], setting: {speed: {channel: 3}} } }怎么验证对账。同步完立刻比对源库和目标分区的行数再按主键抽几行核对字段一致。分区路径里带日期天然支持增量回溯——哪天数据坏了把那天分区重跑一遍就行不用全量重灌。机器学习流水线训练 → 评估 → 部署算法团队最疼的不是模型训不出来而是训完没人管、上线全靠手动。MLflow 插件把这段收编进调度一条工作流串起数据准备、训练、评估、部署跑完自动登记实验和模型版本。MLflow 任务用 PROJECTS 模式mlflowTrackingUri 指向你的 mlflow-server:5000数据路径指向原始数据。要调参就切到搜索模式把超参写成搜索空间交给框架去扫。训练产出模型后评估节点跑一遍测试集指标达标才走部署不达标回炉调参——这个达标才往下走用条件分支卡住。怎么验证看两样一是 MLflow 里这次 run 的实验、参数、指标是不是都登记上了、能回溯二是部署节点起的服务能否用真实请求打通返回的预测和离线评估口径一致。模型上线前的最后一道关是线上返回的值和离线算的是不是一回事。这三套案例能跑通只是及格线真正折磨人的是半夜那次失败的自动重试。怎么让任务不裸奔重试、超时与告警收口白天跑得好好的一到夜里批量并发任务就开始各种挂——超时、OOM、上游延迟。裸奔的任务挂一次就得人肉去翻日志。重试不是万能药。网络抖一下、资源偶发不足值得重试SQL 写错了、数据真的缺重试一百次也是白搭。 failRetryTimes 给 3、failRetryInterval 给 5分钟够覆盖瞬时抖动又不至于把真故障拖成慢性故障。超时要留余量。timeout 给一个比正常耗时高 50% 的值超时策略用 WARN先告警不直接杀。批任务杀一半最坑——表写了一半下游要么拿到脏数据、要么拿到空表。{ name: robust_dws_job, type: SPARK, failRetryTimes: 3, failRetryInterval: 5, timeoutFlag: OPEN, timeout: 3600, timeoutNotifyStrategy: WARN }节点故障交给注册中心兜底。Worker 挂了它身上的任务会被重新调度Master 挂了WatchManager 通过 ZooKeeper 的删除事件发现节点下线接管在途命令。你不用自己写心跳和抢占逻辑。质量校验独立成节点。守门员查四项行数对量级、时间不穿越到明天、去重用户数不异常四条全过才放行。-- 守门员四项不通过就判失败 SELECT COUNT(*) AS total_rows, COUNT(DISTINCT user_id) AS distinct_users, MIN(event_time) AS min_ts, MAX(event_time) AS max_ts FROM dws_user_behavior WHERE dt ${system.biz.date} 告警要收口。一个工作流挂一个告警组节点失败、超时、质量校验失败都走同一个群别每个任务各配一套——半夜你会被几十条告警淹没也分不清哪条该管。任务本身稳了接下来是让它在一个会扩缩、会挂节点、还会被人动配置的集群里继续稳。上生产高可用架构、K8s 部署与监控体系单机跑得再顺上生产第一件烦人的事就是节点会挂、流量会涨、有人要改配置。整体形态是多 Master 加多 WorkerZooKeeper 做服务发现和分布式锁UI 经 API 下发Worker 执行具体任务插件。K8s 部署三个最容易被忽略的参数容器化直接套用 Helm Chart官方部署模板里的 values.yaml 有几个地方最容易翻车。Master 3 副本起步因为单 Master 挂掉时在途命令得有别人接管Worker 按待执行队列扩缩5 只是起点看积压再调externalDatabase 单独指出去别让调度库和业务库共用一个实例互相拖。master: replicas: 3 resources: requests: { memory: 4Gi, cpu: 2 } limits: { memory: 8Gi, cpu: 4 } env: MASTER_EXEC_THREADS: 200 MASTER_DISPATCH_TASK_NUM: 5 worker: replicas: 5 env: WORKER_EXEC_THREADS: 100 externalDatabase: enabled: true type: mysql database: dolphinscheduler externalRegistry: registryPluginName: zookeeper registryServers: zk-0.zookeeper:2181,zk-1.zookeeper:2181生产诉求怎么配为什么Master 高可用3 副本起单节点挂掉在途命令得有别人接管Worker 弹性按待执行队列扩缩积压涨的是 Worker 不够不是 Master调度库独立专用 MySQL 独立账号避免和业务库抢连接和 IO数据库调优索引和历史清理监控页和详情页最常按 state start_time 查没索引就是全表扫任务实例按 process_instance_id 关联加索引后翻页和聚合都顺。成功实例 30 天后清控制表体积留 30 天是为了回溯排障。-- 高频查询路径加索引 ALTER TABLE t_ds_process_instance ADD INDEX idx_state_start_time (state, start_time); ALTER TABLE t_ds_task_instance ADD INDEX idx_process_instance_id (process_instance_id); -- 历史数据定期清 DELETE FROM t_ds_process_instance WHERE state SUCCESS AND start_time DATE_SUB(NOW(), INTERVAL 30 DAY);备份、安全与监控怎么搭备份每天一个全量 mysqldump加 --single-transaction 保证一致性保留 30 天滚动。配置进 Git 管理分 production、staging 目录谁动了配置能追溯。安全上数据库密码、AK 这类别明文写进 ConfigMap丢 Kubernetes Secret 里挂载再叠一层 NetworkPolicy只放本 namespace 和数据库 namespace 的进出。监控走 Prometheus 抓取、Grafana 出图。盯四个Master 命令消费速率、失败任务数、Worker 待执行队列长度、数据库连接数。队列开始积压往往比 CPU 报警更早——该扩 Worker 了。更多部署细节可翻官方文档。落地速查场景对做法你的场景推荐做法T1 离线数仓Spark cluster 模式 独立队列尾部挂质量校验节点实时风控Flink 常驻消费告警从主链路拆成独立分支MySQL → Hive 同步DataX mysqlreader hdfswriter分区带日期支持回溯机器学习流水线MLflow PROJECTS 模式串训练-评估-部署条件分支卡指标怕半夜失败failRetryTimes3 超时 WARN 先告警别让任务裸奔上生产多 Master ZooKeeper 服务发现指标盯队列长度而非 CPU这些不是越多越好而是每个场景挑对应的几招。真正拉开差距的是你愿不愿意在上线前把失败了怎么办这件事先想清楚。【免费下载链接】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),仅供参考
返回列表