DolphinScheduler 3.1.8 与 Flink/Spark 集成:5个典型ETL工作流实战
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节点组成 :
- Kafka源配置节点 :定义消费的Topic和反序列化方式
- Flink SQL处理节点 :执行数据清洗和转换
- 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层结果导出]
关键配置技巧 :
- 资源动态分配 :根据数据量调整Executor数量
spark-submit \
--conf spark.dynamicAllocation.enabled=true \
--conf spark.shuffle.service.enabled=true \
--conf spark.dynamicAllocation.maxExecutors=20
- 分区策略优化 :对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节点 实现跨工作流依赖:
- 上游工作流结束时写入标记文件
# 在Shell节点中执行
hdfs dfs -touchz /checkpoints/${workflowInstanceId}.success
- 下游工作流通过条件依赖检查
# 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混合管道
实时维度表更新结合离线计算的典型案例:
工作流结构 :
- Flink实时任务 :监听MySQL binlog更新Redis维度数据
- Spark离线任务 :每日全量刷新HBase维度表
- 数据一致性检查 :对比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 资源隔离方案
通过多租户实现计算资源隔离:
- 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>
- 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监控指标采集 :
- 配置DolphinScheduler暴露指标
# application.yaml
metrics:
enabled: true
exporter:
type: prometheus
port: 12346
- 关键告警规则示例
# 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 大规模工作流优化
千万级任务调度优化方案 :
- 数据库优化 :
-- 建立关键表索引
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);
- 工作流拆分原则 :
- 单工作流任务数不超过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 数据一致性保障
端到端精确一次方案 :
- 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");
- Spark+HDFS :
// 使用Delta Lake实现ACID
df.write
.format("delta")
.mode("overwrite")
.option("replaceWhere", s"dt='${partitionValue}'")
.save("/data/events")
8. 扩展开发与API集成
8.1 自定义任务类型开发
开发SparkSQL任务类型的完整流程:
- 实现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);
}
}
}
- 注册插件到resources/plugin.properties :
sparksql=org.apache.dolphinscheduler.plugin.task.sparksql.SparkSQLTask
8.2 REST API深度集成
常用API操作示例 :
- 创建工作流定义 :
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"}
)
- 补数操作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"
}'
更多推荐




所有评论(0)