SeaTunnel Source Connector 开发指南:从用户契约到分片并行读取的完整实现路径
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
Source Connector(数据源连接器)是 SeaTunnel 数据集成作业的入口:它负责从外部数据源读取数据并转换为SeaTunnelRow交给下游 transform 与 sink。本指南面向想要为 SeaTunnel 贡献或自研新 source connector 的开发者,以docs/zh/developer/source-connector-development.md为主线,结合seatunnel-api源码与仓库内参考实现,给出从"定义用户契约"到"打包发布"的完整实现路线图。读完本文,你将掌握一个生产级 source connector 的类结构设计、并行分片与容错机制、常见数据源模式以及测试与打包检查清单。
Source Connector 必须解决什么问题
无论数据源是文件、数据库、消息队列还是 CDC,一个 source connector 至少要解决四件事:
- 定义并校验用户可见的配置参数——用户通过作业配置(HOCON)传入的参数必须被稳定地声明、校验和解析;
- 描述输出 schema——告诉引擎这个 source 会产出什么结构的数据(列、类型),供下游与类型系统对接;
- 支持 batch、streaming 或两者兼有的数据读取——即引擎侧所称的 bounded(有界)与 unbounded(无界)语义;
- 在需要并行时支持 split 分配与状态恢复——把数据切成可独立读取的分片(split),并在故障后能精确恢复。
在 SeaTunnel 中,这通常意味着要实现一组配套的类:
- source factory:用户视角的入口,负责暴露插件标识、声明参数规则、创建 source 实例;
SeaTunnelSource:顶层 source 定义,充当创建 reader、enumerator 与序列化器的工厂;- 一个或多个
SourceReader:在 worker 侧真正读取数据的执行单元; - 如果要并行,还需要 split 和 enumerator:split 是数据的最小可分配单元,enumerator 在协调端负责分片发现与分配。
这些接口的完整契约定义在 seatunnel-api 下,下文会逐一展开。
推荐开发流程
1. 先定义用户契约
在写任何运行时代码之前,先把用户能看到的东西定义清楚。这是整个 connector 的"对外 API",一旦发布就很难随意改动:
- plugin 名称:即用户配置中的
source { ... }块所使用的 plugin 名,运行时通过它找到 factory; - required options:必填参数(如 Kafka 的
bootstrap.servers、JDBC 的url); - optional options:可选参数及其语义;
- default value:每个可选参数的默认值;
- 最小可运行配置:一个只填必填项就能跑通的示例配置。
如果你还说不清这个 connector 最小配置长什么样,通常说明实现也还没有真正想清楚。可以先参考仓库内的配置体系文档:
- 作业配置指南:了解作业文件如何组织 source / transform / sink;
- 配置与 Option 系统:了解
OptionRule如何驱动参数声明、校验与 UI 生成。
2. 实现 Factory
Factory 是用户视角下的入口,它至少要负责:
- 暴露稳定的 identifier:实现 Factory 接口的
factoryIdentifier(),返回一个全局唯一的插件标识。源码注释明确要求"为保持一致性,identifier 应为小写单词(如kafka);若存在多版本 factory,用-追加版本(如elasticsearch-7)"; - 定义
OptionRule:实现optionRule(),声明必填、可选、互斥、条件参数等规则; - 创建 source 实例:通过
createSource(context)返回一个TableSource包装的SeaTunnelSource。
在实际系统中,factory 也是文档、运行时校验、REST 元数据暴露、UI 配置生成之间的桥梁:OptionRule不仅用于校验用户配置,还用于 Web-UI 提示用户配置选项(见 Factory.java 的源码注释)。
以仓库中最简单的参考实现 FakeSourceFactory 为例,可以看到完整的 factory 写法:
@AutoService(Factory.class) public class FakeSourceFactory implements TableSourceFactory, SupportSourceDryRunValidation { @Override public String factoryIdentifier() { return "FakeSource"; } @Override public OptionRule optionRule() { return OptionRule.builder() .exclusive(ConnectorCommonOptions.TABLE_CONFIGS, ConnectorCommonOptions.SCHEMA) .optional( STRING_FAKE_MODE, TINYINT_FAKE_MODE, /* ... */, ROWS, ROW_NUM, SPLIT_NUM, SPLIT_READ_INTERVAL, /* ... */) .conditional(STRING_FAKE_MODE, FakeSourceOptions.FakeMode.TEMPLATE, STRING_TEMPLATE) .conditional(INT_FAKE_MODE, FakeSourceOptions.FakeMode.TEMPLATE, INT_TEMPLATE) .build(); } @Override public <T, SplitT extends SourceSplit, StateT extends Serializable> TableSource<T, SplitT, StateT> createSource(TableSourceFactoryContext context) { return () -> (SeaTunnelSource<T, SplitT, StateT>) new FakeSource(context.getOptions()); } @Override public Class<? extends SeaTunnelSource> getSourceClass() { return FakeSource.class; } }关键点解读:
@AutoService(Factory.class)会在编译期自动生成META-INF/services/org.apache.seatunnel.api.table.factory.Factory注册文件,这正是 SPI 注册的落地方式;OptionRule.builder()提供required(...)、optional(...)、exclusive(...)、conditional(...)等声明方法(见 OptionRule.java 及后续 builder 方法),用于表达参数间的依赖与互斥关系;createSource直接由配置构造 source 实例,配置解析逻辑通常抽到独立的XxxConfig类中。
3. 实现 Source 运行时
简单 source 可能只需要 reader;需要扩展性和容错的 source,一般还需要 split 和 enumerator。典型职责如下:
SeaTunnelSource:顶层 source 定义。查看接口源码 SeaTunnelSource.java,核心契约包括:getBoundedness():声明 BOUNDED / UNBOUNDED;createReader():创建运行在工作节点侧的SourceReader;createEnumerator()/restoreEnumerator():创建/恢复运行在主节点侧的SourceSplitEnumerator(后者在从 checkpoint 恢复时调用,接收checkpointState参数);getProducedCatalogTables():声明输出的表元数据(CatalogTable列表,支持多表与模式信息;接口注释建议所有 connector 优先实现该方法而非已废弃的getProducedType());getSplitSerializer()/getEnumeratorStateSerializer():split 与枚举器状态的序列化器,用于网络传输与 checkpoint 持久化,默认使用 Java 序列化的DefaultSerializer。
SourceSplitEnumerator:在 master 侧发现并分配工作。接口见 SourceSplitEnumerator.java,关键方法包括open()、run()、registerReader(subtaskId)、handleSplitRequest(subtaskId)、addSplitsBack(splits, subtaskId)(reader 失败时回收未完成分片)与snapshotState(checkpointId)。注意源码对调用顺序的约定:首次run()之前,引擎会按open()→addSplitsBack(...)→registerReader(...)的固定顺序非并发地调用;SourceReader:在 worker 侧真正读取数据。接口见 SourceReader.java,核心方法为pollNext(Collector<T> output)(拉取下一批数据)、addSplits(splits)、snapshotState(checkpointId)(返回List<SplitT>形式的分片状态)、handleNoMoreSplits(),并通过Context提供sendSplitRequest()(向 enumerator 请求分片)、signalNoMoreElement()(有界数据读完后通知框架结束)等能力;- serializer:在网络传输和 checkpoint 时持久化 split / enumerator 状态。
运行时交互流程
围绕上述接口,引擎侧的典型流程(详见 Source 架构文档)为:
- 初始启动:协调端调用
createEnumerator(context)→ enumeratoropen()后在run()内完成分片发现;worker 侧创建 reader 后通过context.sendSplitRequest()请求分片,enumerator 在handleSplitRequest(subtaskId)中调用context.assignSplit(subtaskId, splits)下发,reader 收到后addSplits(splits)并进入pollNext(collector)循环产出数据; - 检查点:框架向 reader 触发 barrier,reader 在
snapshotState(checkpointId)中快照剩余分片/进度;enumerator 同时快照自己的分配状态;全部确认后持久化检查点; - 失败恢复:失败 reader 的已分配未完成分片通过
addSplitsBack回收并标记为待处理,新 reader 在恢复restoreState后重新注册,enumerator 将回收的分片重新分配,新 reader 从检查点偏移量继续消费。
4. 补齐打包与发现元数据
一个 connector 不是"代码能编译就完成了"。你还需要补齐:
- SPI 注册:如第 2 步所示,用
@AutoService(Factory.class)或手动在META-INF/services中注册 factory。SeaTunnel 运行时通过 SeaTunnelFactoryDiscovery 使用标准ServiceLoader.load(Factory.class, classLoader)加载所有 factory 实现,插件发现与类加载机制的完整说明见 插件发现与类加载; - plugin mapping:在仓库根目录的 plugin-mapping.properties 中注册插件名到 connector 模块的映射。例如 FakeSource 的映射行为
seatunnel.source.FakeSource = connector-fake; - 分发包打包配置:让 connector jar 真正进入二进制包。需要在 seatunnel-connectors-v2/pom.xml 中加入
<module>connector-xxx</module>,并在 seatunnel-dist/pom.xml 的依赖清单与 assembly-seatunnel-bin.xml 的分发清单中登记该 connector; - 依赖隔离:如果需要依赖隔离,还要补齐 plugin 目录布局。SeaTunnel 会把插件实现 jar 与 connector 专属第三方依赖分开管理(
connectors/connector-xxx.jar与plugins/connector-xxx/*.jar),映射关系同样由plugin-mapping.properties维护,详见 插件发现与类加载。
5. 写文档和测试
一个用户可见的 connector,如果没有完成下面这些,通常不能算真的完成:
- 同步更新
docs/en和docs/zh:中英文档必须对齐,插件名、参数说明保持一致; - 示例配置与代码完全一致:文档里的示例配置必须与
OptionRule声明的参数和默认值严格对应; - 单测或 E2E 覆盖主读取路径:至少覆盖正常读取路径,外部系统依赖场景补 E2E 测试。
设计检查清单
编码前,先把这些问题回答清楚,因为这些答案应该驱动你的类结构,而不是反过来:
- 这个 source 是bounded、unbounded,还是两者都支持(决定
getBoundedness()的返回与 reader 的结束逻辑); - split 的单位是什么:文件、分片、分区、表范围,还是别的(决定
SourceSplit的字段设计); - reader 在没有工作时怎么继续请求任务(决定
pollNext空转时的行为,应主动sendSplitRequest()); - 恢复时需要保存哪些状态(决定 reader state 与 enumerator state 的内容与序列化方式);
- schema 是自动发现还是用户配置(决定是否实现 schema discoverer,还是依赖
schema/table_list等用户参数); - 输出是单表还是多表(决定
getProducedCatalogTables()返回的CatalogTable数量,多表场景通常需要额外的表路由逻辑); - 输出的是 CDC 语义还是 append-only 数据(决定是否引入
SeaTunnelRowType之外的元数据列,如_table_name、CDC 的 op 类型列)。
典型类结构
对于一个支持并行的 source,最常见的最小结构如下:
connector-<name>/ src/main/java/.../source/ <Name>SourceFactory.java <Name>Source.java <Name>SourceReader.java <Name>SourceSplit.java <Name>SourceSplitEnumerator.java <Name>SourceConfig.java对照仓库中 connector-fake 的实现,这套结构一一对应:FakeSourceFactory、FakeSource、FakeSourceReader、FakeSourceSplit、FakeSourceSplitEnumerator,外加FakeSourceState(枚举器状态类)与MultipleTableFakeSourceConfig(多表配置解析类)。
复杂一点的实现通常还会加入:
- dialect 或 client 抽象:屏蔽不同数据库方言或不同版本 client 的差异;
- split serializer:自定义高效的 split 序列化(Kryo、Protobuf 等),替代默认 Java 序列化;
- enumerator state:记录已分配/待分配分片的快照对象;
- reader state 辅助类:跟踪每个分片内部的消费进度(offset / position);
- schema discoverer:从数据源自动发现并推断表结构。
什么时候用哪种设计
什么时候简单 Reader 就够了
适用于:
- 数据源天然单线程(如 socket、HTTP 轮询等);
- 不需要并行;
- 没有明确的 split 模型。
这种情况下可以不实现(或仅实现空实现的)split 与 enumerator,reader 在pollNext中持续产出数据即可,但要注意:即使单线程 source,也建议遵循非阻塞轮询模式,让出工作线程。
什么时候必须引入 Split 和 Enumerator
适用于:
- 数据源可以按分区或范围并行读取;
- 故障后需要回收并重新分配未完成任务(依赖
addSplitsBack与 checkpoint); - 初始发现逻辑与 worker 侧读取逻辑应当分离。
对数据库、文件、队列、CDC 这类可扩展 source 来说,这基本是默认模式。核心权衡在于:枚举器-读取器分离带来清晰的协调/执行职责划分与独立的容错边界,代价是多了一次分片分配的网络通信和更复杂的 API;分片粒度方面,粗粒度(少量大分片)协调开销低但负载均衡差、恢复时间长,细粒度(大量小分片)负载均衡与恢复更快但协调开销更高,需要按数据源特性与作业目标权衡。
常见 Source 模式
文件 / 对象存储 Source
- 常见 split 单位:文件、文件块范围、分区目录;
- 常见关注点:文件发现(是否递归、是否过滤)、schema 推断(从表头或文件元数据)、checkpoint 当前文件位置(在 reader state 中保存 path + offset)。
数据库快照 Source
- 常见 split 单位:主键范围、分区、shard;
- 常见关注点:chunk 大小(每片读取的行数/主键区间宽度)、query pushdown(尽量下推过滤条件到 SQL)、一致性边界(快照一致性读,避免读到中间态)。
消息队列 Source
- 常见 split 单位:topic partition、subscription shard;
- 常见关注点:offset 管理(提交与 checkpoint 的联动)、watermark 或 event time(用于事件时间处理)、动态分区发现(新 partition 出现时通过
SourceEvent通知 reader)。
CDC Source
- 常见 split 单位:snapshot chunk、incremental log split;
- 常见关注点:snapshot 到 incremental 的切换(全量快照完成后无缝切到 binlog/log 消费)、source metadata(表名、主键、操作类型等附加列)、schema evolution(源端表结构变更的处理)。
相关的架构设计细节可进一步阅读 CDC Pipeline 架构概览。
测试策略
至少建议覆盖这些层次:
- option 校验:合法的必填/可选组合能通过,缺失必填项或非法组合(如
exclusive冲突)能报出明确错误; - split 生成或发现逻辑:enumerator 的
run()能按预期产出分片集合; - reader 在正常数据上的行为:
pollNext能正确产出记录并通过 collector 下发; - checkpoint 或 state snapshot 行为:reader 的
snapshotState与 enumerator 的snapshotState返回的状态可序列化且内容正确; - 恢复或 split 回收分配(并行 source):模拟 reader 失败,验证
addSplitsBack回收、重新分配后新 reader 从正确位置继续消费。
仓库中的参考单测可作模板,例如 FakeSourceSplitEnumeratorTest 覆盖了枚举器分片发现与分配逻辑。如果 connector 依赖外部系统(数据库、消息队列等),尽可能补或扩展 E2E 测试——seatunnel-e2e/seatunnel-connector-v2-e2e下每个 connector 都有对应的connector-xxx-e2e模块可参考。
打包检查清单
提交 PR 前,建议确认:
- factory 注册已经存在:
@AutoService(Factory.class)生效,META-INF/services元数据能随 jar 一起被打进包; - connector module 已加入构建与分发:已在 seatunnel-connectors-v2/pom.xml 注册 module,并在 seatunnel-dist/pom.xml 与 assembly-seatunnel-bin.xml 中登记分发;
- 需要时已更新
plugin-mapping.properties:插件名到 connector 模块的映射正确(参考 plugin-mapping.properties 中seatunnel.source.FakeSource = connector-fake的写法); - 文档示例里的 plugin 名与运行时 identifier 完全一致:
factoryIdentifier()返回值必须与作业配置中的插件名、文档示例严格对应; - 中英文文档都已补齐:
docs/en与docs/zh同步更新。
推荐阅读顺序
- 先读本页(
docs/zh/developer/source-connector-development.md),建立实现检查表; - 再读 Source 架构:深入理解 enumerator-reader 分离、checkpoint 与失败恢复流程、性能与可扩展性设计;
- 再读 插件发现与类加载:理解 SPI 注册、jar 定位与依赖隔离;
- 参考
seatunnel-connectors-v2/下一个现有 connector:最简单的起点是 connector-fake,需要真实数据源语义时参考 connector-kafka、connector-jdbc 或 connector-file; - 最后结合 开发自己的 Connector:按场景跳转到 sink 侧、配置系统或 CDC 相关文档,形成完整的 connector 开发知识地图。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考