SeaTunnel 架构概览:从分层设计到 Zeta 引擎的分布式数据集成全解析
SeaTunnel 架构概览从分层设计到 Zeta 引擎的分布式数据集成全解析【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 是一个多模态、高性能、分布式的海量数据集成工具本文以 架构概览文档 为主体结合仓库源码seatunnel-api、seatunnel-engine、seatunnel-translation 等模块与 配置模板系统讲解其分层架构、核心组件Source/Sink/Transform/Table API、数据流与执行模型、作业状态机、容错与精确一次语义并给出模块结构与源码级佐证。读完本文你将掌握 SeaTunnel 从配置提交 → LogicalDag → 物理计划 → 管道执行 → 两阶段提交的完整链路理解连接器如何在多引擎上无泄漏地运行以及如何基于其插件体系扩展新的连接器。1. 引言SeaTunnel 的设计目标与目标场景1.1 设计目标SeaTunnel 定位为分布式多模态数据集成工具其设计目标贯穿整个架构引擎独立性连接器逻辑尽量与执行引擎解耦连接器可通过转换层Translation Layer适配到不同引擎。具体可用性以连接器能力与引擎支持为准。超高性能支持高吞吐、低延迟的大规模数据同步。容错性在启用 checkpoint 且外部系统支持事务/幂等提交等前提下通过分布式快照与提交协议提供可验证的一致性语义。易用性提供简单的配置方式HOCON/SQL和丰富的连接器生态系统。可扩展性基于插件的架构便于添加新的连接器和转换组件。其中引擎独立性直接体现在代码结构上连接器实现只依赖 seatunnel-api 中定义的标准接口如SeaTunnelSource、SeaTunnelSink而非任何引擎 API引擎适配由 seatunnel-translation 模块统一完成。1.2 目标场景批量数据同步异构数据源之间的大规模批量数据迁移。实时数据集成支持 CDC 的流式数据捕获和同步。数据湖/仓入库高效加载数据到数据湖Iceberg、Hudi、Delta Lake和数据仓库。多表同步在单个作业中同步多个表支持模式演化Schema Evolution。1.3 推荐阅读路径架构章节用于建立整体认知建议按以下顺序阅读先读本页建立分层视图再看 配置与 Option 系统理解插件配置如何定义、校验和暴露再看 Transform 插件体系理解 transform 如何位于 source、sink、schema 与引擎适配之间再看 表模型与类型系统理解表元数据和可移植类型如何贯穿整条链路关注 CDC 链路可继续看 CDC Pipeline 架构概览再看 Checkpoint 机制 与 Exactly-Once理解一致性语义再看 资源管理理解 slot 分配与 worker 协调再看 插件发现与类加载理解插件打包、发现与依赖隔离如需理解多引擎适配再看 转换层。2. 整体架构五层分层模型SeaTunnel 采用分层架构实现关注点分离和灵活性2.1 层级职责层级职责核心组件配置层作业定义、参数配置HOCON 解析器、SQL 解析器、配置验证API 层连接器的统一抽象数据源/数据 Sink /转换接口、CatalogTable连接器层数据源/Sink 实现连接器实现JDBC、Kafka、CDC 等转换层引擎特定适配Flink/Spark 适配器、上下文包装器引擎层作业执行和资源管理调度、容错、状态管理源码佐证配置层在仓库中以 config/v2.batch.config.template 等模板文件呈现采用env {}、source {}、sink {}的 HOCON 结构转换层由 seatunnel-translationFlink 适配器与 seatunnel-translationSpark 适配器承载引擎层由 seatunnel-engine 承载。3. 核心组件3.1 SeaTunnel API引擎独立的抽象层API 层提供引擎无关的抽象连接器只依赖这一层开发。数据源 Source APISeaTunnelSource创建读取器和枚举器的工厂接口同时负责提供分片序列化器、枚举器状态序列化器以及源数据的 CatalogTable 描述SourceSplitEnumerator主节点侧组件负责分片Split生成和分配SourceReader工作节点侧组件负责从分片读取数据SourceSplit表示数据分区的最小可序列化单元。关键设计协调Enumerator与执行Reader分离实现高效的并行处理和容错。源码佐证SeaTunnelSource.java 定义了createReader、createEnumerator、restoreEnumerator三个核心工厂方法其中restoreEnumerator专门用于从 checkpoint 状态恢复枚举器——这正是故障恢复能力在 API 契约层面的体现SourceSplitEnumerator.java 中的addSplitsBack(splits, subtaskId)方法用于读取器失败后把上次成功 checkpoint 之后分配给它的分片收回并重新分配SourceReader.java 通过pollNext(Collector)拉取数据、snapshotState(checkpointId)快照分片状态、handleNoMoreSplits()处理分片耗尽信号。数据 Sink APISeaTunnelSink创建写入器和可选提交策略的工厂接口SinkWriter工作节点侧组件负责写入数据SinkCommitter工作节点侧的可选提交器负责独立提交单个 writer 的变更SinkAggregatedCommitter协调端聚合提交路径上的可选全局提交器。关键设计两阶段提交协议prepareCommit → commit在外部系统支持事务/幂等提交且启用 checkpoint 的前提下可提供一致性语义。源码佐证SeaTunnelSink.java 中createCommitter()、createAggregatedCommitter()均返回Optional说明提交策略是可选能力——不实现提交的 Sink 无法获得精确一次保证SinkWriter.java 的prepareCommit(checkpointId)在 checkpoint 前生成提交信息注释明确指出两阶段提交2PC的用法SinkAggregatedCommitter.java 中commit()与restoreCommit()的 javadoc 均强调方法需要实现幂等性idempotencycombine()方法负责把各 writer 的 CommitInfo 合并为一次全局提交。转换 APISeaTunnelTransform数据转换接口SeaTunnelMapTransform1:1 转换SeaTunnelFlatMapTransform1:N 转换。表 APICatalogTable完整的表元数据模式、分区键、选项TableSchema模式定义列、主键、约束SchemaChangeEvent表示模式演化的 DDL 变更。3.2 SeaTunnel Engine (Zeta)原生执行引擎原生执行引擎 seatunnel-engine 提供作业调度、分布式快照与资源管理能力。主节点组件CoordinatorService管理所有运行中的 JobMasterJobMaster管理单个作业生命周期、生成物理计划、协调检查点CheckpointCoordinator每个管道协调分布式快照ResourceManager管理工作节点资源和槽位分配。工作节点组件TaskExecutionService部署和执行任务SeaTunnelTask执行数据源 Source/转换/数据 Sink 逻辑FlowLifeCycle管理数据源 Source/转换/数据 Sink 组件的生命周期。源码佐证上述组件均可在 seatunnel-engine/seatunnel-engine-server 下找到对应实现例如 CoordinatorService.java、JobMaster.java、CheckpointCoordinator.java。执行模型从LogicalDag用户意图到PhysicalPlan执行细节的分离是逻辑 vs 物理设计原则的直接体现具体执行与调度细节可参考 引擎架构 与 DAG 执行。3.3 转换层多引擎可移植性通过适配器模式实现引擎可移植性FlinkSource/FlinkSink将 SeaTunnel API 适配到 Flink 的数据源/Sink 接口对应 seatunnel-translation-flink-common 下的FlinkSource、FlinkSinkSparkSource/SparkSink将 SeaTunnel API 适配到 Spark 的 RDD/Dataset 接口对应 seatunnel-translation-spark-common上下文适配器包装引擎特定的上下文SourceReaderContext、SinkWriterContext序列化适配器桥接 SeaTunnel 和引擎序列化机制。3.4 连接器生态系统标准化结构 SPI 发现所有连接器遵循标准化结构区域典型文件职责Source 入口[Name]Source.java、[Name]SourceReader.java、[Name]SourceSplitEnumerator.java、[Name]SourceSplit.java读取数据、切分任务并暴露统一的 Source 契约Sink 入口[Name]Sink.java、[Name]SinkWriter.java缓冲、写入并向目标系统提交数据配置定义config/[Name]Config.java定义连接器参数、校验规则和默认值SPI 注册META-INF/services/TableSourceFactory、META-INF/services/TableSinkFactory注册工厂供运行时发现和装载发现机制Java SPI服务提供者接口用于动态连接器加载。以 connector-kafka、connector-cdc-mysql、connector-jdbc 等为例它们均位于 seatunnel-connectors-v2 下命名、结构与上述表格一致。4. 数据流模型4.1 数据读取 Source 端数据流链路中的行数据统一为SeaTunnelRow转换链可选位于 Source 与 Sink 之间可对行与 schema 进行改写。4.2 基于分片的并行度数据源被划分为分片如文件块、数据库分区、Kafka 分区每个SourceReader独立处理一个或多个分片动态分片分配实现负载均衡和故障恢复分片状态被检查点化以实现精确一次处理。源码佐证SourceReader.java 的snapshotState(checkpointId)返回分片 checkpoint 状态列表且注释明确说明如果 source 是有界bounded的则不触发 checkpointSourceSplitEnumerator.java 的addSplitsBack则承担故障后分片回收重分配职责。4.3 管道执行作业被划分为 SubPlan下图表示同一个作业中的两个独立子计划它们之间不存在直接的数据记录流转每个管道具有独立的并行度配置维护自己的检查点协调器可以并发或顺序执行。5. 作业执行流程5.1 提交阶段5.2 执行阶段任务初始化将任务部署到分配的槽位初始化数据 Source/转换/数据 Sink 组件从检查点恢复状态如果在恢复中。数据处理SourceReader 从分片拉取数据数据流经转换链SinkWriter 缓冲和写入数据。检查点协调CheckpointCoordinator 触发检查点检查点屏障流经数据管道任务快照其状态协调器收集确认。提交阶段SinkWriter 准备提交信息默认由工作节点侧SinkCommitter独立提交各 writer 的变更如果启用聚合提交则改由协调端执行一次全局提交状态持久化到检查点存储。5.3 状态机任务状态转换失败说明FAILED是运行时对不可恢复错误的结果标记但失败后是否重启由更高层的恢复逻辑决定不应在任务状态机图中画成FAILED → ...的直接边。作业状态转换6. 关键特性6.1 容错检查点机制受 Chandy-Lamport 算法启发的分布式快照检查点屏障在数据流中传播状态存储在可插拔的检查点存储中HDFS、S3、本地对应 seatunnel-engine-storage 模块下的 checkpoint-storage 系列子模块从最新成功的检查点自动恢复。故障转移策略任务级故障转移重启失败的任务和相关管道基于区域的故障转移最小化对未受影响任务的影响分片重新分配失败的分片重新分配给健康的工作节点。6.2 精确一次语义两阶段提交协议准备阶段SinkWriter 在检查点期间准备提交信息提交阶段默认由工作节点侧SinkCommitter独立提交各 writer 的变更如果启用聚合提交则在 checkpoint 成功后由协调端执行一次全局提交中止处理在提交前失败时回滚。幂等性SinkCommitter与SinkAggregatedCommitter的提交操作都必须保持幂等以便在重试场景下不重复生效——这一点在 SinkAggregatedCommitter.java 的接口 javadoc 中已作为硬性约束写明。需要注意精确一次语义是有前提条件的——必须启用 checkpoint且外部系统支持事务或幂等提交。这也是 Checkpoint 机制 与 Exactly-Once 两篇文档反复强调的边界。6.3 动态资源管理基于槽位的分配细粒度的资源管理基于标签的过滤将任务分配到特定的工作节点组负载均衡多种策略随机、槽位比率、系统负载动态扩缩容无需重启作业即可添加/移除工作节点未来特性。资源分配与 worker 协调机制的详细说明见 资源管理。6.4 模式演化DDL 传播从数据源捕获模式变更ADD/DROP/MODIFY 列模式映射通过管道转换模式变更动态应用将模式变更应用到数据 Sink 表兼容性检查在应用前验证模式变更。6.5 多表支持单作业多表在一个作业中同步数百个表表路由根据 TablePath 将记录路由到正确的数据 Sink独立模式每个表维护自己的模式副本支持每个表多个写入器副本以获得更高吞吐量。多表能力在 API 层也有对应支撑例如 seatunnel-api 下的SupportMultiTableSink、SupportMultiTableSinkWriter等标记接口以及multitablesink子包。7. 模块结构模块代表子目录职责seatunnel-apisource、sink、transform、table定义核心 API、表模型与跨引擎抽象seatunnel-connectors-v2connector-jdbc、connector-kafka、connector-cdc-mysql提供各类数据源与目标端连接器实现seatunnel-transforms-v2srcSQL、Filter 等提供通用转换能力seatunnel-engineseatunnel-engine-core、seatunnel-engine-server、seatunnel-engine-storage承载 Zeta 执行、调度和检查点存储seatunnel-translationseatunnel-translation-flink、seatunnel-translation-spark负责多引擎适配层seatunnel-formatsseatunnel-format-json、seatunnel-format-avro处理不同数据格式seatunnel-coreCLI 与提交入口负责作业提交和命令行能力seatunnel-e2e端到端测试套件保障关键链路回归以 config/v2.batch.config.template 为例一个最简单的批量作业配置如下它演示了配置层与连接器层的配合方式env { parallelism 2 job.mode BATCH checkpoint.interval 10000 } source { FakeSource { parallelism 2 plugin_output fake row.num 16 schema { fields { name string age int } } } } sink { Console { } }env段定义作业级参数并行度、运行模式、检查点间隔source/sink段按连接器名声明插件并给出各自参数——解析、校验与暴露这些 Option 的机制详见 配置与 Option 系统。8. 设计原则8.1 关注点分离API vs 实现清晰的 API 边界支持多种实现协调 vs 执行枚举器与聚合提交编排负责协调读取器与写入器负责工作节点上的实际执行逻辑 vs 物理LogicalDag用户意图与 PhysicalPlan执行细节分离。8.2 插件架构基于 SPI 的发现连接器通过 Java SPI 动态加载类加载器隔离每个连接器使用隔离的类加载器热插拔无需重新构建核心即可添加连接器。8.3 引擎独立性统一 API相同的连接器代码在任何引擎上运行转换层将 API 适配到引擎特定细节无引擎泄漏连接器开发人员无需了解引擎知识。8.4 可扩展性水平扩展添加工作节点以提高吞吐量基于分片的并行度细粒度并行处理无状态工作节点工作节点可以动态添加/移除。8.5 可靠性分布式检查点跨分布式任务的一致性快照增量状态优化大状态的检查点大小精确一次保证端到端一致性。9. 下一步深入特定架构组件设计理念 —— 核心设计原则和权衡Transform 插件体系 —— 理解 transform 插件如何组织、发现并承担行数据与 schema 改写职责表模型与类型系统 —— 理解 schema、元数据和可移植类型如何在 connector 与引擎之间流动CDC Pipeline 架构概览 —— 理解快照、增量变更捕获与 sink 落地如何协同数据 Source 架构 —— 数据源 API 设计深入探讨数据 Sink 架构 —— 数据 Sink API 设计深入探讨插件发现与类加载 —— 理解 factory、jar 与依赖隔离在运行时如何被解析引擎架构 —— SeaTunnel Engine 内部机制检查点机制 —— 容错实现。实践指南如何创建您的连接器快速入门。10. 参考资料10.1 相关概念Apache Flink—— 检查点和状态管理的灵感来源Apache Kafka—— 消费者组模型影响了分片分配Chandy-Lamport 算法—— 分布式快照算法。10.2 仓库源码索引按阅读顺序SeaTunnelSource.javaSourceSplitEnumerator.javaSourceReader.javaSeaTunnelSink.javaSinkWriter.javaSinkAggregatedCommitter.javaCoordinatorService.javaJobMaster.javaCheckpointCoordinator.java【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考