DolphinScheduler 深潜:DAG 调度引擎与数据管道编排
1. [DolphinScheduler 在 coomia-dip 中的定位](#1-dolphinscheduler-在-coomia-dip-中的定位)
Coomia发布于 2025年11月18日13 分钟阅读
分享本文Twitter / X
“系列:S8 技术组件深潜 · 第 10 篇 | 难度:高级 | 阅读时间:20 分钟
DolphinScheduler 深潜:DAG 调度引擎与数据管道编排
#TL;DR
- DolphinScheduler 是 coomia-dip Pipeline & Orchestration(Pipeline & Orchestration Layer)的批处理调度核心,负责 ETL 管道、数据质量检查、定时聚合等 DAG 工作流的编排
- 本文深入分析 DolphinScheduler 的 Master/Worker 架构、DAG 解析与任务分发机制、容错策略、以及与 coomia-dip Ontology 层的集成方案
- 涵盖多租户隔离、资源管理、告警集成、以及 DolphinScheduler 3.x 的 API 驱动工作流定义
#目录
- DolphinScheduler 在 coomia-dip 中的定位
- Master-Worker 架构详解
- DAG 解析与任务调度
- 任务类型与插件体系
- 容错与重试机制
- 多租户与资源隔离
- 与 coomia-dip Ontology 的集成
- API 驱动的工作流定义
- 告警与监控集成
- 性能调优与生产实践
- Key Takeaways
#1. DolphinScheduler 在 coomia-dip 中的定位
#1.1 DolphinScheduler vs Temporal 的分工
coomia-dip 同时使用 DolphinScheduler 和 Temporal,但职责划分清晰:
| 维度 | DolphinScheduler | Temporal |
|---|---|---|
| 适用场景 | 批处理 ETL、定时调度 | 实时事件驱动、长时间工作流 |
| 触发方式 | Cron / 手动 / 依赖触发 | API / Event / Schedule |
| 任务粒度 | 粗粒度(Shell/SQL/Flink Job) | 细粒度(函数级 Activity) |
| DAG 定义 | 可视化拖拽 + JSON | 代码定义 |
| 用户群体 | 数据工程师 | 开发人员 |
| coomia-dip Layer | F(Pipeline) | E(Agent Runtime) |
#1.2 典型使用场景
Code
DolphinScheduler 管理的管道:
1. 铜层数据采集
┌──────────┐ ┌──────────┐ ┌──────────┐
│ MySQL CDC│───→│ Kafka │───→│ Iceberg │
│ 全量同步 │ │ Topic │ │ Bronze │
└──────────┘ └──────────┘ └──────────┘
2. 银层数据清洗
┌──────────┐ ┌──────────┐ ┌──────────┐
│ Iceberg │───→│ Flink │───→│ Iceberg │
│ Bronze │ │ SQL Job │ │ Silver │
└──────────┘ └──────────┘ └──────────┘
3. 金层聚合分析
┌──────────┐ ┌──────────┐ ┌──────────┐
│ Iceberg │───→│ Trino │───→│ Doris │
│ Silver │ │ Query │ │ Gold │
└──────────┘ └──────────┘ └──────────┘
#2. Master-Worker 架构详解
#2.1 核心组件
Code
┌─────────────────────────────────────────────────────┐
│ DolphinScheduler │
│ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ Master │ │ Master │ │ API │ │
│ │ Server │ │ Server │ │ Server │ │
│ └────┬─────┘ └────┬─────┘ └────┬─────┘ │
│ │ │ │ │
│ ┌────┴──────────────┴──────────────┴────┐ │
│ │ ZooKeeper (Registry) │ │
│ └────┬──────────────┬──────────────┬────┘ │
│ │ │ │ │
│ ┌────┴─────┐ ┌────┴─────┐ ┌────┴─────┐ │
│ │ Worker │ │ Worker │ │ Worker │ │
│ │ Server │ │ Server │ │ Server │ │
│ └──────────┘ └──────────┘ └──────────┘ │
│ │
│ ┌──────────────────────────────────────┐ │
│ │ PostgreSQL (Metadata Store) │ │
│ └──────────────────────────────────────┘ │
└─────────────────────────────────────────────────────┘
#2.2 Master Server 职责
- DAG 解析:将工作流定义解析为 DAG 图
- 任务调度:根据依赖关系和优先级分发任务
- 容错管理:检测 Worker 故障,重新分配任务
- Slot 管理:跟踪各 Worker 的可用 Slot
Java
// Master 的调度循环(简化)
public class MasterSchedulerThread implements Runnable {
@Override
public void run() {
while (isRunning) {
// 1. 从 DB 获取待调度的 Command
List<Command> commands = commandService.findCommandPage(masterSlots);
for (Command command : commands) {
// 2. 构建 DAG
ProcessInstance processInstance = createProcessInstance(command);
DAG<String, TaskNode, TaskNodeRelation> dag = buildDAG(processInstance);
// 3. 获取就绪任务(所有前置依赖已完成)
List<TaskNode> readyTasks = getReadyTasks(dag, processInstance);
// 4. 分发到 Worker
for (TaskNode task : readyTasks) {
WorkerInfo selectedWorker = selectWorker(task);
dispatchTask(task, selectedWorker);
}
}
Thread.sleep(100); // 调度间隔
}
}
}
#2.3 Worker Server 职责
- 任务执行:接收并执行来自 Master 的任务
- 日志收集:收集任务日志并上报
- 心跳上报:定期向 Registry 报告存活状态和负载
Java
// Worker 的任务执行流程
public class WorkerTaskExecuteRunnable implements Runnable {
@Override
public void run() {
try {
// 1. 解析任务参数
TaskExecutionContext context = buildContext(taskInstance);
// 2. 创建任务插件实例
AbstractTask task = TaskPluginManager.createTask(
taskInstance.getTaskType(), context
);
// 3. 初始化
task.init();
// 4. 执行
task.handle();
// 5. 上报结果
reportTaskResult(task.getExitStatus());
} catch (Exception e) {
reportTaskResult(TaskExecutionStatus.FAILURE);
}
}
}
#3. DAG 解析与任务调度
#3.1 DAG 数据结构
JSON
{
"tasks": [
{
"id": "task-extract",
"name": "Bronze Layer Extract",
"type": "FLINK",
"dependence": [],
"params": {
"flinkVersion": "1.18",
"jobManagerMemory": "2g",
"taskManagerMemory": "4g",
"mainJar": "coomia-dip-pipeline-1.0.jar",
"mainClass": "com.onto.pipeline.BronzeExtract"
}
},
{
"id": "task-validate",
"name": "Data Quality Check",
"type": "SQL",
"dependence": ["task-extract"],
"params": {
"datasource": "doris-analytics",
"sql": "SELECT COUNT(*) FROM quality_check WHERE status='FAIL'"
}
},
{
"id": "task-transform",
"name": "Silver Layer Transform",
"type": "FLINK",
"dependence": ["task-validate"],
"conditionResult": {
"successNode": ["task-load"],
"failedNode": ["task-alert"]
}
},
{
"id": "task-load",
"name": "Gold Layer Load",
"type": "SQL",
"dependence": ["task-transform"]
},
{
"id": "task-alert",
"name": "Quality Alert",
"type": "HTTP",
"dependence": ["task-transform"]
}
]
}
#3.2 拓扑排序与并行度
Code
DAG 拓扑排序结果:
Level 0: [task-extract] → 并行度 1
Level 1: [task-validate] → 并行度 1
Level 2: [task-transform] → 并行度 1
Level 3: [task-load, task-alert] → 并行度 2(可并行)
#3.3 调度优先级
Java
// 优先级计算公式
int priority = processInstancePriority * 10 + taskInstancePriority;
// 优先级级别
public enum Priority {
HIGHEST(0),
HIGH(1),
MEDIUM(2),
LOW(3),
LOWEST(4);
}
// 调度队列按优先级排序
PriorityQueue<TaskInstance> queue = new PriorityQueue<>(
Comparator.comparingInt(TaskInstance::getProcessInstancePriority)
.thenComparingInt(TaskInstance::getTaskInstancePriority)
.thenComparing(TaskInstance::getSubmitTime)
);
#4. 任务类型与插件体系
#4.1 coomia-dip 使用的任务类型
| 任务类型 | 用途 | coomia-dip 场景 |
|---|---|---|
| SHELL | 脚本执行 | 数据文件处理、环境准备 |
| SQL | 数据库查询 | 数据质量检查、聚合计算 |
| FLINK | Flink Job | CDC 管道、流处理 |
| HTTP | REST API 调用 | 触发 coomia-dip API |
| PYTHON | Python 脚本 | ML 模型训练、数据分析 |
| DEPENDENT | 跨工作流依赖 | 管道依赖链 |
| SUB_PROCESS | 子工作流 | 模块化管道 |
| CONDITIONS | 条件分支 | 数据质量分支 |
#4.2 自定义任务插件
Java
// coomia-dip 自定义任务插件:Ontology Sync
@AutoService(TaskChannelFactory.class)
public class OntologySyncTaskChannelFactory implements TaskChannelFactory {
@Override
public String getName() {
return "ONTOLOGY_SYNC";
}
@Override
public TaskChannel create() {
return new OntologySyncTaskChannel();
}
}
public class OntologySyncTask extends AbstractTask {
private final OntologySyncParameters parameters;
@Override
public void handle() throws TaskException {
// 1. 调用 coomia-dip gRPC API 获取 Schema
OntologySchema schema = ontologyClient.getSchema(
parameters.getWorldId(),
parameters.getObjectType()
);
// 2. 同步到目标存储
syncToTarget(schema, parameters.getTargetDataSource());
// 3. 记录同步结果
setExitStatusCode(TaskConstants.EXIT_CODE_SUCCESS);
}
}
#5. 容错与重试机制
#5.1 故障类型与处理策略
| 故障类型 | 检测方式 | 处理策略 |
|---|---|---|
| Worker 宕机 | ZooKeeper 心跳超时 | 任务重新分配到其他 Worker |
| Master 宕机 | ZooKeeper 选举 | 备 Master 接管 |
| 任务执行失败 | Exit Code 非 0 | 按重试策略重试 |
| 任务超时 | 超时检测线程 | Kill 进程 + 标记失败 |
| 网络分区 | ZooKeeper Session | 等待恢复后重新注册 |
#5.2 重试配置
JSON
{
"task": {
"name": "Bronze Layer Extract",
"retryTimes": 3,
"retryInterval": "1 minute",
"failRetryStrategy": "RETRY_FROM_CURRENT_TASK",
"timeout": {
"strategy": "WARN_AND_FAILED",
"interval": 3600,
"enable": true
}
}
}
#5.3 工作流级别容错
JSON
{
"processDefinition": {
"name": "CDC Bronze Pipeline",
"failureStrategy": "CONTINUE",
"processInstancePriority": "HIGH",
"warningType": "FAILURE",
"warningGroupId": 1
}
}
| 策略 | 行为 |
|---|---|
END | 任何任务失败则终止整个工作流 |
CONTINUE | 失败任务标记失败,无依赖的后续任务继续执行 |
#6. 多租户与资源隔离
#6.1 coomia-dip 的租户模型
Code
DolphinScheduler 租户体系:
Security Level
├── Admin (coomia-dip 管理员)
│ └── 管理所有项目和资源
├── Tenant: coomia-dip-world-001
│ ├── User: data-engineer-1
│ ├── Project: bronze-pipelines
│ ├── Project: silver-pipelines
│ └── Resource: /shared/jars/
├── Tenant: coomia-dip-world-002
│ ├── User: data-engineer-2
│ └── Project: etl-pipelines
└── Tenant: coomia-dip-system
└── Project: system-maintenance
#6.2 Worker Group 资源隔离
YAML
# Worker 分组配置
worker-groups:
- name: "high-memory"
workers: ["worker-1", "worker-2"]
description: "16GB+ 内存,用于 Flink 大作业"
- name: "gpu"
workers: ["worker-3"]
description: "GPU 节点,用于 ML 训练"
- name: "default"
workers: ["worker-4", "worker-5", "worker-6"]
description: "通用任务执行"
#6.3 队列管理
YAML
# YARN 队列映射
queues:
- name: "production"
capacity: 60%
max-capacity: 80%
tenant: "coomia-dip-production"
- name: "development"
capacity: 20%
max-capacity: 40%
tenant: "coomia-dip-dev"
- name: "system"
capacity: 20%
max-capacity: 30%
tenant: "coomia-dip-system"
#7. 与 coomia-dip Ontology 的集成
#7.1 基于 Ontology 的管道模板
Python
# coomia-dip 管道模板生成器
class OntologyPipelineGenerator:
"""根据 Ontology Schema 自动生成 DolphinScheduler 管道"""
def generate_cdc_pipeline(
self,
object_type: str,
source_db: str,
target_lakehouse: str,
) -> dict:
return {
"name": f"cdc-{object_type}-pipeline",
"tasks": [
{
"name": f"extract-{object_type}",
"type": "FLINK",
"params": self._build_flink_cdc_params(
object_type, source_db
),
},
{
"name": f"validate-{object_type}",
"type": "SQL",
"dependence": [f"extract-{object_type}"],
"params": self._build_quality_check_params(object_type),
},
{
"name": f"transform-{object_type}",
"type": "FLINK",
"dependence": [f"validate-{object_type}"],
"params": self._build_transform_params(
object_type, target_lakehouse
),
},
],
}
def _build_flink_cdc_params(self, object_type: str, source_db: str) -> dict:
schema = self.ontology_client.get_schema(object_type)
return {
"flinkVersion": "1.18",
"mainClass": "com.onto.pipeline.CdcExtractor",
"programArguments": (
f"--source-table {schema.source_table} "
f"--source-db {source_db} "
f"--target-table iceberg.bronze.{object_type} "
f"--columns {','.join(p.name for p in schema.properties)}"
),
}
#7.2 管道执行事件回写 Ontology
Python
# DolphinScheduler Webhook → coomia-dip 审计事件
@app.post("/api/v1/dolphinscheduler/webhook")
async def handle_ds_webhook(event: DSWebhookEvent):
if event.type == "PROCESS_INSTANCE_SUCCESS":
await ontology_client.create_audit_event(
event_type="PIPELINE_COMPLETED",
source="DolphinScheduler",
metadata={
"process_id": event.process_instance_id,
"process_name": event.process_name,
"duration_seconds": event.duration,
"task_count": event.task_count,
},
)
elif event.type == "PROCESS_INSTANCE_FAILURE":
await ontology_client.create_audit_event(
event_type="PIPELINE_FAILED",
source="DolphinScheduler",
metadata={
"process_id": event.process_instance_id,
"error": event.error_message,
"failed_task": event.failed_task_name,
},
)
#8. API 驱动的工作流定义
#8.1 DolphinScheduler 3.x REST API
Python
import httpx
class DolphinSchedulerClient:
"""DolphinScheduler REST API 客户端"""
def __init__(self, base_url: str, token: str):
self.base_url = base_url
self.headers = {"token": token}
async def create_process_definition(
self,
project_code: int,
name: str,
task_definition_json: str,
task_relation_json: str,
) -> dict:
async with httpx.AsyncClient() as client:
response = await client.post(
f"{self.base_url}/projects/{project_code}/process-definition",
headers=self.headers,
data={
"name": name,
"taskDefinitionJson": task_definition_json,
"taskRelationJson": task_relation_json,
"tenantCode": "coomia-dip",
"executionType": "PARALLEL",
"timeout": 0,
},
)
return response.json()
async def run_process(
self,
project_code: int,
process_definition_code: int,
schedule_time: str | None = None,
) -> dict:
async with httpx.AsyncClient() as client:
data = {
"processDefinitionCode": process_definition_code,
"failureStrategy": "CONTINUE",
"warningType": "FAILURE",
"scheduleTime": schedule_time,
"startNodeList": "",
"taskDependType": "TASK_POST",
"runMode": "RUN_MODE_SERIAL",
"processInstancePriority": "MEDIUM",
"workerGroup": "default",
}
response = await client.post(
f"{self.base_url}/projects/{project_code}/executors/start-process-instance",
headers=self.headers,
data=data,
)
return response.json()
async def create_schedule(
self,
project_code: int,
process_definition_code: int,
crontab: str,
) -> dict:
async with httpx.AsyncClient() as client:
response = await client.post(
f"{self.base_url}/projects/{project_code}/schedules",
headers=self.headers,
data={
"processDefinitionCode": process_definition_code,
"schedule": json.dumps({
"startTime": "2026-01-01 00:00:00",
"endTime": "2099-12-31 23:59:59",
"crontab": crontab,
"timezoneId": "Asia/Shanghai",
}),
"warningType": "FAILURE",
"failureStrategy": "END",
"processInstancePriority": "MEDIUM",
"workerGroup": "default",
},
)
return response.json()
#9. 告警与监控集成
#9.1 告警插件配置
YAML
# coomia-dip DolphinScheduler 告警配置
alert:
plugins:
- name: "webhook"
config:
url: "http://coomia-dip-api:8080/api/v1/alerts/dolphinscheduler"
headerParams: '{"Authorization": "Bearer ${coomia-dip_TOKEN}"}'
bodyParams: '{"source": "dolphinscheduler", "event": "${msg}"}'
contentField: "msg"
requestType: "POST"
- name: "dingtalk"
config:
webhook: "${DINGTALK_WEBHOOK_URL}"
keyword: "DolphinScheduler"
isAtAll: false
groups:
- name: "pipeline-alerts"
description: "数据管道告警组"
plugins: ["webhook", "dingtalk"]
#9.2 Prometheus 指标暴露
YAML
# DolphinScheduler Prometheus 指标
metrics:
- name: dolphinscheduler_master_running_process_count
type: gauge
help: "当前运行中的流程实例数"
- name: dolphinscheduler_master_task_dispatch_count
type: counter
help: "任务分发总数"
- name: dolphinscheduler_worker_task_execution_count
type: counter
labels: [task_type, status]
help: "Worker 任务执行计数"
- name: dolphinscheduler_worker_task_execution_duration
type: histogram
labels: [task_type]
help: "任务执行耗时分布"
#10. 性能调优与生产实践
#10.1 Master 调优
YAML
# Master 配置优化
master:
max-cpu-load-avg: 0.7 # CPU 负载上限
reserved-memory: 0.3 # 预留内存比例
exec-threads: 100 # 调度线程数
dispatch-task-number: 3 # 每次调度分发的任务数
host-selector: lower-weight # Worker 选择策略
task-commit-interval: 1000 # 任务提交间隔(ms)
#10.2 Worker 调优
YAML
# Worker 配置优化
worker:
exec-threads: 100 # 任务执行线程数
heartbeat-interval: 10 # 心跳间隔(秒)
max-cpu-load-avg: 0.75 # CPU 负载上限
reserved-memory: 0.25 # 预留内存比例
tenant-auto-create: true # 自动创建 Linux 用户
#10.3 数据库优化
SQL
-- 定期清理历史数据(保留 90 天)
DELETE FROM t_ds_process_instance
WHERE start_time < DATE_SUB(NOW(), INTERVAL 90 DAY)
AND state IN (7, 8); -- SUCCESS, FAILURE
-- 创建常用查询索引
CREATE INDEX idx_process_instance_state_start
ON t_ds_process_instance (state, start_time);
CREATE INDEX idx_task_instance_process_state
ON t_ds_task_instance (process_instance_id, state);
#10.4 资源规划
| 组件 | 实例数 | CPU | 内存 | 磁盘 |
|---|---|---|---|---|
| Master | 2 (HA) | 4 核 | 8 GB | 50 GB |
| Worker | 3-10 | 4-8 核 | 8-16 GB | 100 GB |
| API Server | 2 | 2 核 | 4 GB | 20 GB |
| PostgreSQL | 1 (HA) | 4 核 | 16 GB | SSD 200 GB |
| ZooKeeper | 3 | 2 核 | 4 GB | SSD 50 GB |
#11. Key Takeaways
| 主题 | 关键结论 |
|---|---|
| 定位 | 批处理 DAG 调度,与 Temporal 互补 |
| 架构 | Master 调度 + Worker 执行 + ZooKeeper 注册 |
| DAG | 拓扑排序确定执行顺序,同级可并行 |
| 容错 | Worker 故障自动重分配,任务级重试 |
| 多租户 | Tenant + Project + Worker Group 三级隔离 |
| Ontology 集成 | 根据 Schema 自动生成管道定义 |
| API 驱动 | 3.x REST API 支持代码化管道管理 |
| 监控 | Prometheus + 自定义告警插件 |
“下一篇预告:S8-11 已发布,深入分析 Redis 在 coomia-dip 中的 5 种角色。