
数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载Connector V2 是 SeaTunnel 基于SeaTunnelConnector API定义的新一代连接器体系它通过引擎无关的 API 与翻译层设计让同一套连接器代码可以运行在 Flink、Spark 与 SeaTunnel Zeta 等不同引擎上并统一支持批处理与流处理。本文以官方文档 docs/en/concept/connector-v2-features.md 为骨架结合仓库内seatunnel-api的源码实现系统讲解 Connector V2 与 V1 的核心差异、Source 连接器的七大能力、Sink 连接器的三大能力及其底层原理帮助读者在选型与开发连接器时准确判断每项能力是否可用、如何生效。Connector V2 与 V1 的核心差异自 apache/seatunnel#1608 提出并引入 Connector V2 特性以来SeaTunnel 的插件体系发生了根本性重构。Connector V2 是基于 SeaTunnel Connector API 接口定义的连接器它与 V1 最大的不同在于V2 连接器只依赖一套引擎无关的 API即seatunnel-api模块而引擎相关的适配逻辑全部交由翻译层Translation Layer完成。官方文档总结了 V2 相对 V1 的四大核心差异1. 多引擎支持Multi Engine SupportSeaTunnel Connector API 是引擎无关的 API。基于该 API 开发的连接器可以在多种引擎上运行。当前仓库中Flink 与 Spark 的适配由 seatunnel-translation 模块承载其中包含seatunnel-translation-flinkFlink 13/15与seatunnel-translation-sparkSpark 2.4/3.3两个翻译子模块同时 SeaTunnel 自研的 Zeta 引擎seatunnel-engine也是 Connector V2 的原生运行环境。未来还可以继续扩展其他引擎而无需改动连接器本身。2. 多引擎版本支持Multi Engine Version Support通过翻译层将连接器与底层引擎解耦解决了旧有架构中“为了支持底层引擎的一个新版本绝大多数连接器都需要修改代码”的痛点。连接器开发者只面向seatunnel-api编程引擎版本的升级与适配集中在翻译层完成例如 Flink 13 与 Flink 15 分别对应seatunnel-translation-flink-13与seatunnel-translation-flink-15两个独立模块连接器代码不需要随之改动。3. 统一批流Unified Batch And StreamConnector V2 可以同时胜任批处理或流处理不需要为批和流分别开发两套连接器。这一能力在源码层面由Boundedness枚举体现Boundedness.javaBOUNDED有限数据流对应批作业UNBOUNDED无限数据流对应流作业。每个 Source 连接器通过 SeaTunnelSource#getBoundedness() 声明自己属于哪种模式甚至可以按配置动态决定。例如 Kafka Source 在 KafkaSource.java 中会根据是否配置了起始/结束偏移等条件动态返回BOUNDED或UNBOUNDED同一份代码即可同时服务批读与流读场景。4. JDBC/Log 连接复用Multiplexing JDBC/Log ConnectionConnector V2 支持JDBC 资源复用与共享数据库日志解析。在多表场景下多个写入目标可以复用同一个 JDBC 连接降低连接开销同时多个读取任务可以共享同一份数据库变更日志的解析结果为 CDC 类连接器的高效运行提供支撑。Source 连接器特性Source 连接器具备一组通用核心特性每个具体的 Source 连接器对它们的支持程度各不相同。这些能力在 API 层面大多通过**标记接口Marker Interface**或约定式方法体现开发者可以在 seatunnel-api/src/main/java/org/apache/seatunnel/api/source 目录下逐一查看。exactly-once精确一次如果数据源中的每条数据只会被 Source向下游发送一次我们就认为该 Source 连接器支持 exactly-once。在 SeaTunnel 中实现机制如下做 checkpoint 时将当前已读取的Split及其offsetSplit 内读取位置的描述例如文件的行号、字节大小、偏移量等保存为StateSnapshot任务重启后取出上一次的StateSnapshot依据其中的 Split 与 offset 定位到上次读取的位置从断点继续向下游发送数据。这一套“状态快照 恢复”的流程在 API 中对应 SourceReader#snapshotState(checkpointId)返回当前 split 的 checkpoint 状态、SeaTunnelSource#restoreEnumerator(...)用 checkpoint 状态重建枚举器等方法的配合。文档给出的典型例子是File、Kafka类 Source。需要说明的是Source 端 exactly-once 只保证“读取不重不漏”下游是否最终一致还要看 Sink 端的配合。column projection列投影如果连接器支持只从数据源读取指定列则称其支持 column projection。注意如果先把所有列读出来、再通过 schema 过滤掉不必要的列这种“先全量读、再裁剪”的做法不是真正的列投影。例如JDBCSource可以通过 SQL 定义要读取的列属于真正的列投影KafkaSource会先从 topic 读取全部内容再用schema过滤不必要的列因此不构成column projection。在 API 层面这一能力由标记接口 SupportColumnProjection 标识。例如文件类 Source 的 BaseFileSource 同时实现了SupportParallelism与SupportColumnProjection说明它可以在读取阶段就完成列的裁剪。batch批模式批作业模式下数据是**有界bounded**的作业在读完所有数据后会自动停止。对应Boundedness.BOUNDED且按 SourceSplitEnumerator#snapshotState 的注释有界 Source 不触发 checkpoint。stream流模式流作业模式下数据是**无界unbounded**的作业持续运行、不会自动停止。对应Boundedness.UNBOUNDED。典型如持续监听 Kafka topic 新消息、持续解析数据库 binlog 等场景。parallelism并行度并行 Source 连接器支持配置parallelism每个并行度会创建一个 task 来读取数据。并行读取的调度机制是枚举器Enumerator将 Source切分成多个 Split枚举器负责把 Split分配给各个 SourceReader进行处理。这套“枚举-分配-读取”模型在 API 中体现为运行在 master 侧的 SourceSplitEnumerator 与运行在 worker 侧的 SourceReader 的协同枚举器通过assignSplit分发 split、通过registerReader感知已注册的 readerreader 通过sendSplitRequest()向枚举器请求新的 split。是否支持并行由标记接口 SupportParallelism 标识连接器实现该接口后即可在作业配置中设置并行度。support user-defined split支持用户自定义 Split用户可以自定义 Split 的划分规则例如按文件大小、按分区、按时间范围等维度控制数据被切分成多少个读取单元从而更精细地控制并行读取的粒度。support multiple table read支持多表读取支持在一个 SeaTunnel 作业中读取多张表。这类 Source 通常实现 SeaTunnelSource#getProducedCatalogTables() 返回多张上游 CatalogTable官方文档建议所有连接器优先实现该方法而非旧的getProducedType()因为CatalogTable携带更完整的元数据能帮助下游实现更准确、更完整的同步能力典型的如 MySQL-CDC、Oracle-CDC 等 CDC Source。Sink 连接器特性与 Source 类似Sink 连接器也具备一组通用核心特性。相关 API 定义位于 seatunnel-api/src/main/java/org/apache/seatunnel/api/sink 目录。exactly-once精确一次对 Sink 连接器而言任何一条数据在整个处理过程中只被精确写入目标一次、且处理结果正确即认为满足 exactly-once。判定标准是任意一条数据只写入目标系统一次。文档指出通常有两种实现路径路径一目标数据库支持主键去重如果目标数据库支持基于主键的幂等去重Sink 连接器可以借助“重复写入被去重吸收”来达到最终 only-once 的效果。典型例子为MySQL、Kudu。路径二XA 事务 两阶段提交如果目标支持XA 事务该事务可以跨会话使用——即使创建事务的程序已经结束新启动的程序只需知道上一次事务的 ID即可重新提交或回滚该事务则可以利用**两阶段提交Two-phase Commit**保证 exactly-once。典型例子为File、MySQL。在源码层面JDBC Sink 的实现非常典型JdbcExactlyOnceSinkWriter 内部使用XidGenerator生成事务 XID、以XidInfo记录并持久化待提交的 XA 事务信息并提供recoverAndRollback恢复后回滚未决事务的能力配合 JdbcSinkCommitter / JdbcSinkAggregatedCommitter 完成提交与回滚。另外注意该实现要求maxRetries必须为 0否则无法保证 XA 语义。Sink 端提交/回滚的统一契约由 SeaTunnelSink 提供createWriter() 创建写入器createCommitter() 创建单个提交器SinkCommitter#commit必须实现幂等性createAggregatedCommitter() 创建聚合提交器SinkAggregatedCommitter在单线程中聚合各 worker 的 commit 消息后统一提交并支持restoreCommit恢复重提交与abort回滚。从源码注释可见官方强烈推荐优先实现SinkAggregatedCommitter因为当前版本下它能提供更一致的行为而abort目前在 Spark 引擎上才会被触发。cdc变更数据捕获如果 Sink 连接器支持基于主键写入多种行类型RowKind——即INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE——就认为它支持 CDC。行类型枚举定义于 RowKind.java这使下游能根据上游的增删改事件做对应写入例如将 UPDATE 拆为 DELETEINSERT或直接按主键 upsert从而支撑“数据库到数据库”的实时同步链路。support multiple table write支持多表写入支持在一个 SeaTunnel 作业中写入多张表用户可以通过**配置占位符placeholder**动态指定目标表的标识符。该能力在 API 层面对应 SupportMultiTableSink 等标记接口例如文件类 Sink 的 BaseMultipleTableFileSink 即同时实现了SeaTunnelSink与SupportMultiTableSink。多表写入的占位符机制多表写入依赖 docs/en/concept/sink-options-placeholders.md 介绍的占位符特性通过占位符在上游 CatalogTable 元数据中动态取值在连接器启动前完成替换确保 Sink 参数就绪后才开始使用。该特性在 SeaTunnel Zeta、Flink、Spark 三种引擎上均受支持。支持的占位符表达式如下占位符含义备注${database_name}上游 CatalogTable 中的数据库名支持默认值${database_name:default_my_db}${schema_name}上游 CatalogTable 中的 schema 名支持默认值${schema_name:default_my_schema}${table_name}上游 CatalogTable 中的表名支持默认值${table_name:default_my_table}${schema_full_name}schema 全路径database schema—${table_full_name}表全路径database schema table—${primary_key}上游表的主键字段—${unique_key}上游表的唯一键字段—${field_names}上游表的字段名集合—使用前提所使用的 Sink 连接器必须实现了TableSinkFactoryAPI见 TableSinkFactory.java。配置示例一MySQL-CDC 源 JDBC Sinkenv { // ignore... } source { MySQL-CDC { // ignore... } } transform { // ignore... } sink { jdbc { url jdbc:mysql://localhost:3306 driver com.mysql.cj.jdbc.Driver user root password 123456 database ${database_name}_test table ${table_name}_test primary_keys [${primary_key}] } }配置示例二Oracle-CDC 源 JDBC Sinkenv { // ignore... } source { Oracle-CDC { // ignore... } } transform { // ignore... } sink { jdbc { url jdbc:mysql://localhost:3306 driver com.mysql.cj.jdbc.Driver user root password 123456 database ${schema_name}_test table ${table_name}_test primary_keys [${primary_key}] } }占位符未替换的排查占位符替换发生在连接器启动之前。如果某个变量没有被替换通常是上游表元数据中缺少该选项所致例如mysql类型 Source 的元数据中不包含${schema_name}MySQL 中库与 schema 同层oracle类型 Source 的元数据中不包含${database_name}其他类似情况。此时需要检查所选 Source 是否提供了对应的元数据字段或改用占位符的默认值语法如${table_name:default_my_table}兜底。从 API 到引擎能力如何落地理解上述特性的底层支撑关键是把握seatunnel-api中三个核心接口的职责划分SeaTunnelSourceSource 的“工厂”负责构造SourceSplitEnumerator、SourceReader及对应的序列化器。getBoundedness()决定批/流语义createEnumerator()/restoreEnumerator()分别用于新建与从 checkpoint 恢复枚举器SourceSplitEnumerator运行在 master 侧负责枚举 Split、通过assignSplit分配、通过registerReader/handleSplitRequest与 reader 交互并通过snapshotState保存状态以支持 exactly-once 恢复SourceReader运行在 worker 侧通过pollNext产出数据、snapshotState上报读取位置、addSplits接收分配到的 split、handleNoMoreSplits感知无更多数据。对应地Sink 侧由 SeaTunnelSink 统一组织SinkWriter写数据、SinkCommitter/SinkAggregatedCommitter两阶段提交/回滚的创建为 exactly-once 写入提供统一的编程契约。小结与能力对照Connector V2 通过引擎无关的 API 翻译层架构实现了多引擎、多引擎版本、统一批流、JDBC/Log 复用四大体系性改进Source 侧的 exactly-once、列投影、批/流、并行度、自定义 Split、多表读取以及 Sink 侧的 exactly-once主键去重 / XA 两阶段提交、CDC 行类型写入、多表写入占位符动态指定目标表构成了完整的连接器能力图谱。在实际使用中建议按以下方式快速判断能力是否可用查看连接器工厂类实现了哪些标记接口SupportParallelism、SupportColumnProjection、SupportMultiTableSink等查看 SeaTunnelSource#getBoundedness() 判断批/流语义查看 Sink 是否实现SinkAggregatedCommitter或 XA 相关类如 JDBC 的JdbcExactlyOnceSinkWriter以判断 exactly-once 的落地方式多表写入场景对照 docs/en/concept/sink-options-placeholders.md 的占位符清单与示例配置进行参数拼装。这套“文档 API 标记 实现类”的对照方法既是理解单个连接器能力的捷径也是评估一个连接器能否满足特定同步场景精确一次、多表、CDC、并行吞吐的可靠依据。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel Connector V2 特性详解多引擎适配、批量流统一与 Source/Sink 核心能力体系SeaTunnel Connector V2 特性详解多引擎适配、批量流统一与 Source/Sink 核心能力体系 本文围绕 docs/en/introdu数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Connector V2 功能全景解析多引擎架构与 Source/Sink 核心特性实战指南SeaTunnel Connector V2 功能全景解析多引擎架构与 Source/Sink 核心特性实战指南 本篇技术指南围绕 Apache SeaTun数据集成ETL大数据批处理流处理变更数据捕获革命性数据集成引擎Apache SeaTunnel批流一体架构原理解析革命性数据集成引擎Apache SeaTunnel批流一体架构原理解析 你是否还在为数据同步时批处理与流处理需要两套系统而烦恼是否因数据一致性问题导致业务决数据集成ETL大数据批处理流处理变更数据捕获上一篇终极指南如何可视化Qwen1.5模型的注意力权重机制下一篇Webpacker编译原理深度解析从源码到打包的完整流程揭秘创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考