SeaTunnel Engine(Zeta)引擎架构深度解析:Master-Worker 协调、DAG 执行与容错恢复全解
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
SeaTunnel Engine(Zeta)是 SeaTunnel 原生自研的分布式执行引擎,面向数据同步(Data Sync)与 CDC 场景设计。本文以 engine-architecture.md 为骨架,结合仓库内seatunnel-engine模块的真实源码,系统讲解其 Master-Worker 总体架构、CoordinatorService / JobMaster / ResourceManager 核心组件职责、DAG 到 PhysicalPlan 的转换与 Pipeline 执行模型、Task 状态机与 FlowLifeCycle 生命周期、基于 Chandy-Lamport 的 Checkpoint 协调机制、Slot 资源管理与标签过滤,以及 Task / Worker / Master 三级故障处理与设计取舍。读完本文,你将能理解 Zeta 引擎从提交一个 HOCON 配置到数据最终写入目标端的完整链路,并具备据此排查调度、资源与容错问题的源码级能力。
上图完整呈现了 Zeta 引擎从客户端组装、Coordinator 调度、Worker 端执行到状态上报的端到端流程,对应下文第 2.2 节的逐步拆解。
1. 概述:Zeta 引擎要解决什么问题
1.1 问题背景
任何数据集成引擎都必须回答以下分布式系统基础问题:
- 分布式执行(Distributed Execution):如何把作业调度到多台机器上执行?
- 资源管理(Resource Management):如何高效地分配与调度任务资源?
- 容错(Fault Tolerance):Worker / Master 宕机后如何恢复?
- 协调(Coordination):如何同步分布式任务(Checkpoint、Commit)?
- 可扩展性(Scalability):如何应对持续增长的作业负载?
1.2 设计目标
SeaTunnel Engine(Zeta)以「原生执行引擎」为定位,其设计目标如下:
- 轻量(Lightweight):最小化依赖、快速启动、低资源开销;
- 高性能(High Performance):针对数据同步负载专项优化;
- 容错(Fault Tolerance):基于 Checkpoint 恢复,提供 exactly-once 语义;
- 资源高效(Resource Efficiency):基于 Slot 的资源管理,细粒度控制;
- 引擎无关(Engine Independence):与 Flink / Spark 翻译层共用同一套 Connector API。
1.3 架构对比
| 特性 | SeaTunnel Zeta | Apache Flink | Apache Spark |
|---|---|---|---|
| 主要用途 | 数据同步、CDC | 流处理 | 批处理 + ML |
| 资源模型 | Slot 制 | Slot 制 | Executor 制 |
| 状态后端 | 可插拔(HDFS/S3/Local) | RocksDB/Heap | 内存/磁盘 |
| Checkpoint | 分布式快照 | Chandy-Lamport | RDD 血缘 |
| 运维复杂度 | 较低(引擎原生) | 较高 | 较高 |
2. 总体架构:Master-Worker 模型
2.1 整体拓扑
Master 节点运行三类核心服务(均由 CoordinatorService.java 统一装配),Worker 节点则通过TaskExecutionService承载具体的 Source / Transform / Sink 执行流。
2.2 作业提交、类加载与任务分发流程
Zeta 引擎将「客户端组装 → Coordinator 调度 → Worker 执行 → 状态上报」串联成一条完整链路,各环节职责如下:
- 客户端解析插件 JAR:客户端在本地解析插件 JAR,提交
JobImmutableInformation——该对象携带逻辑 DAG、插件 JAR URL 与 Connector JAR 标识符; - CoordinatorService 接收提交:接受提交、记录 pending 作业,在
JobMaster.init()完成且作业入队后返回提交确认; - Scheduler 轮询派发:Coordinator 持有的
Scheduler轮询PendingJobQueue,通过preApplyResources()申请资源,待ResourceFuture就绪后触发JobMaster.run(); - JobMaster 展开物理计划:将逻辑作业展开为 Pipeline 化的
ExecutionPlan与PhysicalPlan,构建TaskGroupImmutableInformation,并经PhysicalVertex.deploy()与DeployTaskOperation下发到 Worker; - Worker 侧解析与执行:每个 Worker 解析缺失的 JAR、为每个任务创建 child-first 类加载器、在
TaskExecutionService内反序列化并初始化 TaskGroup、执行任务,并把部署与终止状态回传给JobMaster与CoordinatorService。
2.3 核心组件职责与源码落点
CoordinatorService —— 集群作业总管
集中管理集群内所有作业,职责包括:
- 接收作业提交;
- 为每个作业创建 JobMaster;
- 在分布式 IMap 中维护作业状态;
- 提供作业查询与管理 API;
- 处理作业生命周期事件。
其关键数据结构(底层为 Hazelcast 分布式 IMap)在源码 CoordinatorService.java 中可直接对应:
// 运行中作业状态(Hazelcast 分布式 IMap 承载) IMap<Long, JobInfo> runningJobInfoIMap; IMap<Long, JobStatus> runningJobStateIMap; IMap<Long, Long> runningJobStateTimestampsIMap; // 已结束作业历史 IMap<Long, JobInfo> completedJobInfoIMap;JobMaster —— 单作业生命周期管理者
每个作业对应一个 JobMaster,职责包括:
- 解析配置 → 生成 LogicalDag;
- 由 LogicalDag 生成 PhysicalPlan;
- 向 ResourceManager 申请资源(Slot);
- 将任务部署到 Worker;
- 协调各 Pipeline 的 Checkpoint;
- 处理任务失败并重新调度。
生命周期为Created → Initialized → Scheduled → Running → Finished / Failed / Canceled。三个关键操作在 JobMaster.java 中均有对应实现:
init():生成物理计划、创建 Checkpoint 协调器;run():申请资源、部署任务、启动执行;handleFailure():重启失败任务、从 Checkpoint 恢复。
ResourceManager —— 资源与 Slot 分配中枢
负责管理 Worker 资源与 Slot 分配,职责包括:
- 跟踪 Worker 注册与心跳;
- 维护 Worker 资源画像(CPU、内存);
- 依据策略分配 Slot(随机、Slot 占比、系统负载);
- 任务完成后释放 Slot;
- 处理 Worker 故障。
Slot 分配策略在 ResourceManager.java 与 AbstractResourceManager.java 中体现为三种实现:
// 1. Random:在可用 Worker 中随机选择 // 2. SlotRatio:优先选择可用 Slot 更多的 Worker // 3. SystemLoad:优先选择 CPU/内存占用更低的 WorkerAbstractResourceManager.applyResources(jobId, resourceProfiles, tagFilter)的流程是:先按tagFilter过滤候选 Worker,再按所选策略为每个资源画像挑选 Worker 并从其未分配 Slot 池中取一个 Slot——这是理解后面「标签过滤」与「分配策略」两条配置线的入口。
3. DAG 执行模型:从 HOCON 配置到可运行任务
3.1 执行计划的多层转换
3.2 LogicalDag:引擎无关的用户意图
LogicalDag以引擎无关的方式表达用户意图,对应源码位于 seatunnel-engine/seatunnel-engine-core/src/main/java/org/apache/seatunnel/engine/core/dag/:
public class LogicalDag { private final Map<Long, LogicalVertex> logicalVertexMap; private final Set<LogicalEdge> edges; private final JobConfig jobConfig; } public class LogicalVertex { private final long vertexId; private final Action action; // SourceAction / TransformChainAction / SinkAction private final int parallelism; } public class LogicalEdge { private final long inputVertexId; private final long targetVertexId; }其创建入口为LogicalDagBuilder.build(jobConfig):解析 HOCON 配置中的source/transform/sink段,为每个组件创建Action对象,依据配置结构推断数据流边,并做 schema 兼容性校验。详细的顶点、边与并行度建模可进一步阅读 dag-execution.md。
3.3 PhysicalPlan:带资源分配的物理执行计划
public class PhysicalPlan { private final List<SubPlan> pipelineList; private final JobImmutableInformation jobImmutableInformation; private final CompletableFuture<JobResult> jobEndFuture; } public class SubPlan { private final int pipelineId; private final List<PhysicalVertex> physicalVertexList; private final List<PhysicalVertex> coordinatorVertexList; private final CheckpointCoordinator checkpointCoordinator; } public class PhysicalVertex { private final TaskGroupLocation taskGroupLocation; private final TaskGroupDefaultImpl taskGroup; private final SlotProfile slotProfile; // Assigned slot private final ExecutionState currentExecutionState; }生成逻辑在JobMaster.getPhysicalPlan()内部完成三步:把 LogicalDag 切分为多个 Pipeline;为每个并行实例生成 PhysicalVertex;为每个 Pipeline 创建 CheckpointCoordinator。
3.4 Pipeline 执行:独立执行单元
作业被划分为若干Pipeline(SubPlan)独立执行。以下面多源多汇配置为例:
env { ... } source { MySQL-CDC { table = "orders" } Kafka { topic = "events" } } transform { Sql { query = "SELECT * FROM orders JOIN events ON ..." } } sink { Elasticsearch { index = "orders" } JDBC { table = "events" } }生成的 Pipeline 为:
Pipeline 1:MySQL-CDC → Transform → ElasticsearchPipeline 2:Kafka → Transform → JDBC
需要说明的是,切分规则由当前
PipelineGenerator实现决定:不连通子图会拆成独立 Pipeline;若某个连通子图内存在多输入顶点(如 UNION/JOIN),则沿每条 source→sink 路径拆分并按需克隆顶点;而单纯的多汇(多 sink 分支)并不必然产生多个 Pipeline,无多输入顶点时通常保持单个 Pipeline。
Pipeline 化带来三个收益:独立的 Checkpoint 协调(降低协调开销)、隔离的故障域(一个 Pipeline 失败不影响其他 Pipeline)、Pipeline 并行执行。
3.5 Task Fusion:任务融合优化
多个 Action 可融合进单个 TaskGroup 以提升效率:
| 模式 | 运行时形态 | 权衡 |
|---|---|---|
| 不融合 | Source Task → Network → Transform Task → Network → Sink Task | 阶段边界清晰,但网络序列化开销大 |
| 融合 | TaskGroup: Source → Transform → Sink(单线程内) | 网络成本低、局部性好,但调度灵活性下降 |
融合条件:并行度一致、顺序依赖、无需 shuffle。以Source(4) → Transform(4) → Sink(4)为例:不融合是 12 个独立 Task 且阶段间有网络跳转;融合后为 4 个 TaskGroup,每个组内串行执行SourceTask → TransformTask → SinkTask,减少网络序列化、改善 CPU 缓存局部性并降低内存占用。
4. 任务生命周期与执行
4.1 Task 状态机
关键状态迁移:
- CREATED → INIT:任务创建并初始化运行时资源;
- INIT → WAITING_RESTORE / READY_START:在「恢复路径」与「全新启动」之间抉择;
- WAITING_RESTORE → READY_START:状态恢复完成、Flow 准备 open;
- READY_START → STARTING → RUNNING:任务收到启动信号进入主处理循环;
- RUNNING → PREPARE_CLOSE → CLOSED:屏障处理与清理后的正常完成路径;
- 活跃态 → CANCELLING → CANCELED:外部取消路径,独立于正常完成流程。
关于 FAILED 的说明:FAILED作为运行时结果存在,但任务级重启由更上层的恢复逻辑处理,而非由状态机中的FAILED → ...直接迁移。
4.2 SeaTunnelTask 执行骨架
public abstract class SeaTunnelTask implements Runnable { private final TaskLocation taskLocation; private final TaskExecutionContext executionContext; private ExecutionState executionState; @Override public void run() { try { init(); restoreState(); // If recovering open(); while (isRunning()) { processData(); // Source: read, Transform: process, Sink: write handleBarrier(); // Checkpoint barriers } close(); } catch (Exception e) { handleException(e); } } }三种任务类型:SourceSeaTunnelTask(运行 SourceReader、产出数据)、SinkSeaTunnelTask(运行 SinkWriter、消费数据)、TransformSeaTunnelTask(运行 Transform 链)。
4.3 FlowLifeCycle 组件生命周期管理
每个任务通过 FlowLifeCycle 管理组件生命周期,核心实现对应 seatunnel-engine-server 的 task 包:
// Source 任务 public class SourceFlowLifeCycle<T> implements FlowLifeCycle { private final SourceReader<T, ?> sourceReader; private final SeaTunnelSourceCollector collector; @Override public void open() { sourceReader.open(); } @Override public void collect() { sourceReader.pollNext(collector); } // 读取数据 @Override public void close() { sourceReader.close(); } } // Sink 任务 public class SinkFlowLifeCycle<T> implements FlowLifeCycle { private final SinkWriter<T, ?, ?> sinkWriter; @Override public void collect() { T record = inputQueue.poll(); sinkWriter.write(record); // 写入数据 } }5. Checkpoint 协调:一致快照与 exactly-once
5.1 CheckpointCoordinator(每 Pipeline 一个)
每个 Pipeline 拥有独立的 Checkpoint 协调器(源码见 CheckpointCoordinator.java),职责:
- 周期性触发 Checkpoint;
- 向数据流注入 Checkpoint Barrier;
- 收集任务 ACK;
- 持久化完成的 Checkpoint;
- 清理过期 Checkpoint。
public class CheckpointCoordinator { private final CheckpointIDCounter checkpointIdCounter; private final Map<Long, PendingCheckpoint> pendingCheckpoints; private final ArrayDeque<String> completedCheckpointIds; private final CheckpointStorage checkpointStorage; }Checkpoint 主流程:
- 协调器触发 Checkpoint(周期或手动);
- 向 Pipeline 内所有 Source 任务发送 Barrier;
- Barrier 沿数据流向下游传播;
- 每个任务收到 Barrier 后快照自身状态;
- 任务向协调器发送 ACK;
- 协调器等待全部 ACK;
- 创建 CompletedCheckpoint 并持久化到存储。
补充:完整机制(PendingCheckpoint 的 ACK 聚合、CompletedCheckpoint 结构、Barrier 对齐、恢复流程、两阶段提交)可参阅 checkpoint-mechanism.md。Checkpoint 存储类型在引擎侧(
config/seatunnel.yaml的seatunnel.engine.checkpoint.storage)配置,而非作业级env选项。
5.2 Checkpoint Barrier
Barrier 是随数据流动的特殊控制消息:
public class Barrier { private final long checkpointId; private final long timestamp; private final CheckpointType type; // CHECKPOINT or SAVEPOINT }Barrier 对齐:多输入任务在快照前必须等待所有输入到达同一 checkpointId 的 Barrier,从而保证跨分布式任务的一致快照。
6. 资源管理:Slot 模型、分配策略与标签过滤
6.1 Slot 与资源画像
SlotProfile(资源分配的基本单元)与WorkerProfile(Worker 节点资源与 Slot 库存快照):
public class SlotProfile { private final int slotID; private final Address worker; private final ResourceProfile resourceProfile; // CPU, memory } public class ResourceProfile { private final CPU cpu; private final Memory heapMemory; } public class WorkerProfile { private final Address address; private final ResourceProfile profile; private final ResourceProfile unassignedResource; private final SlotProfile[] assignedSlots; private final SlotProfile[] unassignedSlots; private final Map<String, String> attributes; }WorkerProfile 的生命周期为:启动时向 ResourceManager 注册 → 周期心跳上报资源信息 → 从 unassigned 池分配 Slot → 任务完成释放 Slot 回池 → 离开集群(优雅退出或故障)。
6.2 资源分配流程
当可用 Slot 不足时,ResourceManager 抛出NoEnoughResourceException,JobMaster 以退避方式重试等待资源释放。任务结束后由 JobMaster 调用releaseResources归还 Slot。
6.3 Tag 标签过滤:把任务钉到指定 Worker 组
在作业配置中通过env.tag_filter指定 Worker 属性过滤(key/value 全匹配):
env { # 作业级 Worker 属性过滤(key/value 全匹配) tag_filter = { zone = "db-zone" } }典型应用场景:
- 数据本地性(Data Locality):把任务分配到靠近数据源的 Worker;
- 资源隔离(Resource Isolation):例如将 ML Transform 固定到 GPU Worker;
- 多租户(Multi-Tenancy):不同团队使用不同的 Worker 池。
匹配语义为:env.tag_filter与 Worker 的attributes做 key/value 全匹配;若无任何 Worker 匹配则资源分配失败。更详细的策略选型(Random / SlotRatio / SystemLoad 适用场景)与 Slot 配置参见 resource-management.md。
6.4 引擎侧 Slot 配置
以 config/seatunnel.yaml 为例,引擎侧通过seatunnel.engine.slot-service配置 Slot 服务:
seatunnel: engine: classloader-cache-mode: true history-job-expire-minutes: 1440 backup-count: 1 queue-type: blockingqueue print-execution-info-interval: 60 print-job-metrics-info-interval: 60 slot-service: dynamic-slot: true checkpoint: interval: 10000 timeout: 60000 storage: type: hdfs max-retained: 3 plugin-config: namespace: /tmp/seatunnel/checkpoint_snapshot storage.type: hdfs fs.defaultFS: file:///tmp/ # 确保目录有写权限 telemetry: metric: enabled: false logs: scheduled-deletion-enable: true http: enable-http: true port: 8080 enable-dynamic-port: false其中slot-service下的典型配置项还包括slot-num(每 Worker 的 Slot 数量)与slot-allocate-strategy(RANDOM/SLOT_RATIO/SYSTEM_LOAD)。从源码结构看,三类分配策略实现在 resourcemanager/allocation/ 目录 下,与 JobMaster 的SlotAllocationStrategy引用一一对应。
7. 故障处理:三级容错
7.1 任务级故障
检测:任务向 JobMaster 上报异常;JobMaster 监控任务心跳;心跳超时触发故障检测。
恢复:
- 标记任务为 FAILED;
- 释放任务占用的 Slot;
- 获取最近一次成功 Checkpoint;
- 以恢复后的状态重启任务;
- 重新分配 Split(针对 Source 任务)。
7.2 Worker 级故障
检测:ResourceManager 监控 Worker 心跳;Hazelcast 集群检测成员移除。
恢复:
- 将故障 Worker 上的所有任务标记为 FAILED;
- 触发作业 failover;
- 从最近 Checkpoint 恢复;
- 在健康 Worker 上重新分配 Slot;
- 重新部署任务。
7.3 Master 级故障(高可用)
高可用设计:多 Master 节点组成 Hazelcast 集群;作业状态存于分布式 IMap(副本复制);新 Master 从 IMap 状态接管。
恢复:
- 检测 Master 故障(Hazelcast);
- 选举新 Master;
- 新 Master 从 IMap 读取作业状态;
- 重新连接 Worker;
- 恢复 Checkpoint 协调。
值得注意的细节(来自 resource-management.md):ResourceManager 本身无状态,其 Worker 注册表由心跳重建;在 active-master 切换时,Worker 仍上报为「已分配给恢复作业」的 Slot 会被直接复用而非重新申请——该复用要求 Worker 地址、Slot ID、分配序号与归属作业 ID 全部匹配,其中固定 Slot(
dynamic-slot: false)是最大受益场景。
8. 设计考量与性能优化
8.1 为什么采用 Pipeline 执行?
备选方案:全局单 DAG 执行。
决策:划分为多个 Pipeline。
收益:独立 Checkpoint 协调(协调开销小);清晰的故障边界(一个 Pipeline 失败,其他继续运行);数据流更易推理;支持多源多汇的复杂 DAG。
代价:无法跨 Pipeline 边界融合任务;Pipeline 间存在潜在数据序列化开销。
8.2 为什么选择 Hazelcast 作为协调层?
备选方案:ZooKeeper、etcd、自研 Raft。
决策:Hazelcast IMDG。
收益:内存态分布式数据结构(低延迟);内建集群管理与故障检测;易于内嵌(无外部依赖);API 贴近 Java Collections。
代价:大状态的内存开销;作为协调组件,实战检验程度不及 ZooKeeper。
8.3 性能优化手段
- Task Fusion:降低网络开销、改善 CPU 缓存局部性、减少序列化成本;
- 异步 Checkpoint:快照上传不阻塞数据处理,任务间并行快照;
- 增量 Checkpoint:仅上传变化状态(规划中的增强能力);
- 零拷贝数据传输:共置任务间共享内存、避免非必要序列化。
9. 延伸阅读与源码导航
- 架构总览
- 设计理念
- Checkpoint 机制详解
- 资源管理详解
- DAG 执行模型详解
关键源码文件:
- 引擎核心:seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/(CoordinatorService、JobMaster、TaskExecutionService 等)
- DAG:seatunnel-engine/seatunnel-engine-core/src/main/java/org/apache/seatunnel/engine/core/dag/
- Checkpoint:seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/checkpoint/
- 资源管理:seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/resourcemanager/
Zeta 引擎的设计吸收了经典分布式系统的经验——Checkpoint 借鉴 Chandy-Lamport 分布式快照算法,资源管理思路与 Google Borg 一脉相承。理解这套「Coordinator 调度 + Pipeline 化执行 + Slot 资源管控 + Checkpoint 容错」的组合拳,是驾驭 SeaTunnel 在数据同步与 CDC 场景下稳定运行的基础。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考