返回博客

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 驱动工作流定义

#目录

  1. DolphinScheduler 在 coomia-dip 中的定位
  2. Master-Worker 架构详解
  3. DAG 解析与任务调度
  4. 任务类型与插件体系
  5. 容错与重试机制
  6. 多租户与资源隔离
  7. 与 coomia-dip Ontology 的集成
  8. API 驱动的工作流定义
  9. 告警与监控集成
  10. 性能调优与生产实践
  11. Key Takeaways

#1. DolphinScheduler 在 coomia-dip 中的定位

#1.1 DolphinScheduler vs Temporal 的分工

coomia-dip 同时使用 DolphinScheduler 和 Temporal,但职责划分清晰:

维度DolphinSchedulerTemporal
适用场景批处理 ETL、定时调度实时事件驱动、长时间工作流
触发方式Cron / 手动 / 依赖触发API / Event / Schedule
任务粒度粗粒度(Shell/SQL/Flink Job)细粒度(函数级 Activity)
DAG 定义可视化拖拽 + JSON代码定义
用户群体数据工程师开发人员
coomia-dip LayerF(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数据库查询数据质量检查、聚合计算
FLINKFlink JobCDC 管道、流处理
HTTPREST API 调用触发 coomia-dip API
PYTHONPython 脚本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内存磁盘
Master2 (HA)4 核8 GB50 GB
Worker3-104-8 核8-16 GB100 GB
API Server22 核4 GB20 GB
PostgreSQL1 (HA)4 核16 GBSSD 200 GB
ZooKeeper32 核4 GBSSD 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 种角色。