ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

SeaTunnel Engine(Zeta)引擎架构深度解析:Master-Worker 协调、DAG 执行与容错恢复全解

SeaTunnel Engine(Zeta)引擎架构深度解析:Master-Worker 协调、DAG 执行与容错恢复全解 SeaTunnel EngineZeta引擎架构深度解析Master-Worker 协调、DAG 执行与容错恢复全解【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel EngineZeta是 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 ToleranceWorker / Master 宕机后如何恢复协调Coordination如何同步分布式任务Checkpoint、Commit可扩展性Scalability如何应对持续增长的作业负载1.2 设计目标SeaTunnel EngineZeta以「原生执行引擎」为定位其设计目标如下轻量Lightweight最小化依赖、快速启动、低资源开销高性能High Performance针对数据同步负载专项优化容错Fault Tolerance基于 Checkpoint 恢复提供 exactly-once 语义资源高效Resource Efficiency基于 Slot 的资源管理细粒度控制引擎无关Engine Independence与 Flink / Spark 翻译层共用同一套 Connector API。1.3 架构对比特性SeaTunnel ZetaApache FlinkApache Spark主要用途数据同步、CDC流处理批处理 ML资源模型Slot 制Slot 制Executor 制状态后端可插拔HDFS/S3/LocalRocksDB/Heap内存/磁盘Checkpoint分布式快照Chandy-LamportRDD 血缘运维复杂度较低引擎原生较高较高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下发到 WorkerWorker 侧解析与执行每个 Worker 解析缺失的 JAR、为每个任务创建 child-first 类加载器、在TaskExecutionService内反序列化并初始化 TaskGroup、执行任务并把部署与终止状态回传给JobMaster与CoordinatorService。2.3 核心组件职责与源码落点CoordinatorService —— 集群作业总管集中管理集群内所有作业职责包括接收作业提交为每个作业创建 JobMaster在分布式 IMap 中维护作业状态提供作业查询与管理 API处理作业生命周期事件。其关键数据结构底层为 Hazelcast 分布式 IMap在源码 CoordinatorService.java 中可直接对应// 运行中作业状态Hazelcast 分布式 IMap 承载 IMapLong, JobInfo runningJobInfoIMap; IMapLong, JobStatus runningJobStateIMap; IMapLong, Long runningJobStateTimestampsIMap; // 已结束作业历史 IMapLong, 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 MapLong, LogicalVertex logicalVertexMap; private final SetLogicalEdge 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 ListSubPlan pipelineList; private final JobImmutableInformation jobImmutableInformation; private final CompletableFutureJobResult jobEndFuture; } public class SubPlan { private final int pipelineId; private final ListPhysicalVertex physicalVertexList; private final ListPhysicalVertex 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 执行独立执行单元作业被划分为若干PipelineSubPlan独立执行。以下面多源多汇配置为例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 1MySQL-CDC → Transform → ElasticsearchPipeline 2Kafka → 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 准备 openREADY_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 SourceFlowLifeCycleT implements FlowLifeCycle { private final SourceReaderT, ? 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 SinkFlowLifeCycleT implements FlowLifeCycle { private final SinkWriterT, ?, ? sinkWriter; Override public void collect() { T record inputQueue.poll(); sinkWriter.write(record); // 写入数据 } }5. Checkpoint 协调一致快照与 exactly-once5.1 CheckpointCoordinator每 Pipeline 一个每个 Pipeline 拥有独立的 Checkpoint 协调器源码见 CheckpointCoordinator.java职责周期性触发 Checkpoint向数据流注入 Checkpoint Barrier收集任务 ACK持久化完成的 Checkpoint清理过期 Checkpoint。public class CheckpointCoordinator { private final CheckpointIDCounter checkpointIdCounter; private final MapLong, PendingCheckpoint pendingCheckpoints; private final ArrayDequeString completedCheckpointIds; private final CheckpointStorage checkpointStorage; }Checkpoint 主流程协调器触发 Checkpoint周期或手动向 Pipeline 内所有 Source 任务发送 BarrierBarrier 沿数据流向下游传播每个任务收到 Barrier 后快照自身状态任务向协调器发送 ACK协调器等待全部 ACK创建 CompletedCheckpoint 并持久化到存储。补充完整机制PendingCheckpoint 的 ACK 聚合、CompletedCheckpoint 结构、Barrier 对齐、恢复流程、两阶段提交可参阅 checkpoint-mechanism.md。Checkpoint 存储类型在引擎侧config/seatunnel.yaml的seatunnel.engine.checkpoint.storage配置而非作业级env选项。5.2 Checkpoint BarrierBarrier 是随数据流动的特殊控制消息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资源分配的基本单元与WorkerProfileWorker 节点资源与 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 MapString, String attributes; }WorkerProfile 的生命周期为启动时向 ResourceManager 注册 → 周期心跳上报资源信息 → 从 unassigned 池分配 Slot → 任务完成释放 Slot 回池 → 离开集群优雅退出或故障。6.2 资源分配流程当可用 Slot 不足时ResourceManager 抛出NoEnoughResourceExceptionJobMaster 以退避方式重试等待资源释放。任务结束后由 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-strategyRANDOM/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.mdResourceManager 本身无状态其 Worker 注册表由心跳重建在 active-master 切换时Worker 仍上报为「已分配给恢复作业」的 Slot 会被直接复用而非重新申请——该复用要求 Worker 地址、Slot ID、分配序号与归属作业 ID 全部匹配其中固定 Slotdynamic-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 等DAGseatunnel-engine/seatunnel-engine-core/src/main/java/org/apache/seatunnel/engine/core/dag/Checkpointseatunnel-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),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进