DolphinScheduler 3.1.8 与 Flink/Spark 集成:5个典型ETL工作流实战

在大数据生态系统中,任务调度系统与计算引擎的高效协同是构建稳定数据管道的核心。本文将深入探讨如何通过DolphinScheduler 3.1.8版本实现与Flink流处理和Spark批处理的无缝集成,并通过五个典型场景展示实际应用方案。

1. 环境配置与基础集成

1.1 计算引擎环境准备

在开始集成前,需确保计算引擎环境满足以下条件:

Flink环境要求

  • 版本兼容性:支持1.13+版本
  • 集群模式:Standalone/YARN/Kubernetes
  • 网络访问:Worker节点需能访问Flink JobManager REST接口

Spark环境要求

  • 版本支持:Spark 2.4+/3.x
  • 部署模式:Cluster/Client模式
  • 资源管理:与YARN或Kubernetes集成时需配置队列资源

注意:生产环境建议将引擎依赖包(如flink-dist、spark-core)预先部署在所有Worker节点的 libs 目录下,路径通常为 /opt/dolphinscheduler/libs

1.2 DolphinScheduler服务配置

修改 common.properties 关键参数:

# Flink配置
flink.home=/opt/flink-1.15
flink.app.status.poll.interval=10s

# Spark配置
spark.home=/opt/spark-3.2
spark.master=yarn
spark.deploy.mode=cluster

通过API创建计算引擎集群配置:

# 创建Spark集群配置示例
curl -X POST \
  http://ds-server:12345/dolphinscheduler/clusters/create \
  -H 'Token: YOUR_ACCESS_TOKEN' \
  -d '{
    "name": "spark_prod",
    "config": "spark.master=yarn\nspark.executor.memory=4g",
    "type": "SPARK"
  }'

2. 实时数据入湖工作流设计

2.1 Kafka到HDFS的实时ETL

该工作流实现从Kafka消费数据,经Flink实时处理写入HDFS的完整流程。

DAG节点组成

  1. Kafka源配置节点 :定义消费的Topic和反序列化方式
  2. Flink SQL处理节点 :执行数据清洗和转换
  3. HDFS Sink节点 :配置写入路径和文件滚动策略

关键Flink任务参数示例:

-- Flink SQL节点配置
CREATE TABLE kafka_source (
    user_id STRING,
    event_time TIMESTAMP(3),
    metadata ROW<ip STRING, device STRING>
) WITH (
    'connector' = 'kafka',
    'topic' = 'user_events',
    'properties.bootstrap.servers' = 'kafka:9092',
    'format' = 'json'
);

CREATE TABLE hdfs_sink (
    dt STRING,
    hour STRING,
    user_count BIGINT
) PARTITIONED BY (dt, hour) WITH (
    'connector' = 'filesystem',
    'path' = 'hdfs://cluster/data/user_stats',
    'format' = 'parquet',
    'sink.rolling-policy.file-size' = '128MB'
);

INSERT INTO hdfs_sink
SELECT 
    DATE_FORMAT(event_time, 'yyyy-MM-dd') AS dt,
    DATE_FORMAT(event_time, 'HH') AS hour,
    COUNT(DISTINCT user_id) AS user_count
FROM kafka_source
GROUP BY 
    TUMBLE(event_time, INTERVAL '1' HOUR),
    DATE_FORMAT(event_time, 'yyyy-MM-dd'),
    DATE_FORMAT(event_time, 'HH')

2.2 工作流参数化设计

通过全局参数实现动态配置:

参数名 示例值 描述
kafka.brokers kafka:9092 Kafka集群地址
sink.path /data/${bizDate} HDFS写入路径
checkpoint.interval 30000 Flink检查点间隔(ms)

3. 离线批处理工作流优化

3.1 多阶段Spark SQL处理

典型数仓分层处理流程,包含ODS→DWD→DWS三层转换:

graph TD
    A[ODS层数据加载] -->|Spark SQL| B[DWD层维度关联]
    B -->|Spark SQL| C[DWS层聚合计算]
    C -->|Spark SQL| D[ADS层结果导出]

关键配置技巧

  1. 资源动态分配 :根据数据量调整Executor数量
spark-submit \
  --conf spark.dynamicAllocation.enabled=true \
  --conf spark.shuffle.service.enabled=true \
  --conf spark.dynamicAllocation.maxExecutors=20
  1. 分区策略优化 :对DWD层表按日期分区
-- 动态分区设置
SET hive.exec.dynamic.partition=true;
SET hive.exec.dynamic.partition.mode=nonstrict;

INSERT OVERWRITE TABLE dwd.user_behavior PARTITION(dt='${bizDate}')
SELECT 
    user_id,
    item_id,
    behavior_type,
    from_unixtime(event_time) AS action_time
FROM ods.kafka_user_events
WHERE dt='${bizDate}';

3.2 依赖管理实践

通过DolphinScheduler的 Dependent节点 实现跨工作流依赖:

  1. 上游工作流结束时写入标记文件
# 在Shell节点中执行
hdfs dfs -touchz /checkpoints/${workflowInstanceId}.success
  1. 下游工作流通过条件依赖检查
# dependent节点Python脚本
import subprocess
result = subprocess.run(
    ['hdfs', 'dfs', '-test', '-e', f'/checkpoints/{upstream_instance}.success'],
    stdout=subprocess.PIPE
)
exit(0 if result.returncode == 0 else 1)

4. 混合计算场景实现

4.1 Flink+Spark混合管道

实时维度表更新结合离线计算的典型案例:

工作流结构

  1. Flink实时任务 :监听MySQL binlog更新Redis维度数据
  2. Spark离线任务 :每日全量刷新HBase维度表
  3. 数据一致性检查 :对比Redis与HBase关键指标

Flink CDC配置示例

// 构建MySQL CDC源
DebeziumSourceFunction<String> sourceFunction = MySQLSource.<String>builder()
    .hostname("mysql-host")
    .port(3306)
    .databaseList("inventory")
    .tableList("inventory.products")
    .username("flinkuser")
    .password("password")
    .deserializer(new JsonDebeziumDeserializationSchema())
    .build();

// 写入Redis的Sink实现
RedisSink<String> redisSink = new RedisSink<>(
    new FlinkJedisPoolConfig.Builder().setHost("redis").build(),
    new RedisProductUpdateMapper()
);

4.2 资源隔离方案

通过多租户实现计算资源隔离:

  1. YARN队列配置
<!-- capacity-scheduler.xml -->
<queue name="ds_flink">
  <minResources>10000 mb,10 vcores</minResources>
  <maxResources>50000 mb,50 vcores</maxResources>
</queue>

<queue name="ds_spark">
  <minResources>20000 mb,20 vcores</minResources>
  <maxResources>100000 mb,100 vcores</maxResources>
</queue>
  1. DolphinScheduler租户绑定
-- 在数据库中添加租户队列映射
INSERT INTO t_ds_tenant_queue_relation 
(tenant_id, queue_name) VALUES 
(1, 'ds_flink'), 
(2, 'ds_spark');

5. 生产环境最佳实践

5.1 高可用保障措施

Master-Worker架构优化

组件 配置项 推荐值 说明
Master master.exec.threads CPU核心数×2 控制并行工作流数
Worker worker.exec.threads CPU核心数×1.5 控制并行任务数
所有节点 heartbeat.interval 10s 心跳检测间隔

ZooKeeper关键配置

# zoo.cfg
tickTime=2000
initLimit=10
syncLimit=5
maxClientCnxns=100
minSessionTimeout=4000
maxSessionTimeout=40000

5.2 监控与告警集成

Prometheus监控指标采集

  1. 配置DolphinScheduler暴露指标
# application.yaml
metrics:
  enabled: true
  exporter:
    type: prometheus
    port: 12346
  1. 关键告警规则示例
# prometheus_rules.yml
- alert: MasterTaskQueueFull
  expr: ds_master_task_queue_size > 1000
  for: 5m
  labels:
    severity: critical
  annotations:
    summary: "Master任务队列积压 (instance {{ $labels.instance }})"
    description: "任务队列持续高位,当前值: {{ $value }}"

邮件告警模板配置

{
  "type": "EMAIL",
  "name": "prod_alert",
  "params": {
    "receivers": "data-team@company.com",
    "serverHost": "smtp.exmail.qq.com",
    "serverPort": "465",
    "sender": "ds-alert@company.com",
    "enableSmtpAuth": "true",
    "user": "ds-alert",
    "password": "xxxxxx",
    "starttlsEnable": "false",
    "sslEnable": "true"
  }
}

6. 性能调优实战

6.1 大规模工作流优化

千万级任务调度优化方案

  1. 数据库优化
-- 建立关键表索引
CREATE INDEX idx_command_create_time ON t_ds_command(create_time);
CREATE INDEX idx_process_instance ON t_ds_process_instance(id, state);

-- 历史数据归档策略
DELETE FROM t_ds_task_instance 
WHERE end_time < DATE_SUB(NOW(), INTERVAL 30 DAY);
  1. 工作流拆分原则
  • 单工作流任务数不超过200个
  • 复杂DAG拆分为子工作流
  • 使用"子流程"节点进行嵌套调用

6.2 计算引擎参数调优

Flink作业典型配置

# flink-conf.yaml
taskmanager.numberOfTaskSlots: 4
parallelism.default: 10
state.backend: rocksdb
state.checkpoints.dir: hdfs://cluster/flink/checkpoints
state.savepoints.dir: hdfs://cluster/flink/savepoints

Spark性能关键参数

spark-submit \
  --conf spark.sql.adaptive.enabled=true \
  --conf spark.sql.shuffle.partitions=200 \
  --conf spark.executor.memoryOverhead=1g \
  --conf spark.memory.fraction=0.7 \
  --conf spark.serializer=org.apache.spark.serializer.KryoSerializer

7. 异常处理与故障恢复

7.1 失败策略配置

DolphinScheduler提供多种失败处理机制:

策略类型 适用场景 配置方式
继续执行 独立子任务失败 工作流定义→失败策略→继续
结束流程 关键路径失败 工作流定义→失败策略→结束
条件分支 选择性重试 使用条件分支节点

自动重试配置示例

{
  "taskInstance": {
    "retryTimes": 3,
    "retryInterval": 300,
    "failureStrategy": "CONTINUE"
  }
}

7.2 数据一致性保障

端到端精确一次方案

  1. Flink+Kafka
// 启用端到端精确一次
env.enableCheckpointing(60000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);

// Kafka生产者配置
properties.setProperty("enable.idempotence", "true");
properties.setProperty("transactional.id", "flink-producer");
  1. Spark+HDFS
// 使用Delta Lake实现ACID
df.write
  .format("delta")
  .mode("overwrite")
  .option("replaceWhere", s"dt='${partitionValue}'")
  .save("/data/events")

8. 扩展开发与API集成

8.1 自定义任务类型开发

开发SparkSQL任务类型的完整流程:

  1. 实现TaskPlugin接口
public class SparkSQLTask extends AbstractTask {
    @Override
    public AbstractParameters getParameters() {
        return new SparkSQLParameters();
    }

    @Override
    public void handle() throws Exception {
        String sql = sparkSQLParameters.getSql();
        String conf = sparkSQLParameters.getSparkConfig();
        
        // 构建SparkSession
        SparkSession spark = SparkSession.builder()
            .config(parseConfig(conf))
            .enableHiveSupport()
            .getOrCreate();
        
        try {
            spark.sql(sql).show();
            setExitStatusCode(Constants.EXIT_CODE_SUCCESS);
        } catch (Exception e) {
            logger.error("Execute SparkSQL failed", e);
            setExitStatusCode(Constants.EXIT_CODE_FAILURE);
        }
    }
}
  1. 注册插件到resources/plugin.properties
sparksql=org.apache.dolphinscheduler.plugin.task.sparksql.SparkSQLTask

8.2 REST API深度集成

常用API操作示例

  1. 创建工作流定义
import requests

url = "http://ds-server:12345/dolphinscheduler/projects/{projectName}/process-definition"

payload = {
    "name": "flink_etl",
    "description": "实时数据清洗流程",
    "globalParams": "{\"bizDate\":\"${system.datetime}\"}",
    "locations": [
        {
            "taskDefinition": {
                "code": "flink_task_1",
                "type": "FLINK",
                "params": {
                    "mainClass": "com.etl.FlinkJob",
                    "programArgs": "--date ${bizDate}"
                }
            },
            "x": 100,
            "y": 200
        }
    ]
}

response = requests.post(
    url,
    json=payload,
    headers={"Token": "your_access_token"}
)
  1. 补数操作API
curl -X POST \
  "http://ds-server:12345/dolphinscheduler/projects/{projectName}/executors/execute" \
  -H "Token: YOUR_TOKEN" \
  -d '{
    "processDefinitionCode": "workflow_123",
    "scheduleTime": "2026-01-01 00:00:00,2026-01-07 23:59:59",
    "failureStrategy": "CONTINUE",
    "warningType": "NONE",
    "execType": "REPEAT_RUNNING"
  }'
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐