
SeaTunnel Source Connector 开发指南从用户契约到分片并行读取的完整实现路径【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSource 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 和 enumeratorsplit 是数据的最小可分配单元enumerator 在协调端负责分片发现与分配。这些接口的完整契约定义在 seatunnel-api 下下文会逐一展开。推荐开发流程1. 先定义用户契约在写任何运行时代码之前先把用户能看到的东西定义清楚。这是整个 connector 的对外 API一旦发布就很难随意改动plugin 名称即用户配置中的source { ... }块所使用的 plugin 名运行时通过它找到 factoryrequired options必填参数如 Kafka 的bootstrap.servers、JDBC 的urloptional options可选参数及其语义default value每个可选参数的默认值最小可运行配置一个只填必填项就能跑通的示例配置。如果你还说不清这个 connector 最小配置长什么样通常说明实现也还没有真正想清楚。可以先参考仓库内的配置体系文档作业配置指南了解作业文件如何组织 source / transform / sink配置与 Option 系统了解OptionRule如何驱动参数声明、校验与 UI 生成。2. 实现 FactoryFactory 是用户视角下的入口它至少要负责暴露稳定的 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 TableSourceT, SplitT, StateT createSource(TableSourceFactoryContext context) { return () - (SeaTunnelSourceT, 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 / UNBOUNDEDcreateReader()创建运行在工作节点侧的SourceReadercreateEnumerator()/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(CollectorT output)拉取下一批数据、addSplits(splits)、snapshotState(checkpointId)返回ListSplitT形式的分片状态、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 触发 barrierreader 在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 中加入moduleconnector-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/ NameSourceFactory.java NameSource.java NameSourceReader.java NameSourceSplit.java NameSourceSplitEnumerator.java NameSourceConfig.java对照仓库中 connector-fake 的实现这套结构一一对应FakeSourceFactory、FakeSource、FakeSourceReader、FakeSourceSplit、FakeSourceSplitEnumerator外加FakeSourceState枚举器状态类与MultipleTableFakeSourceConfig多表配置解析类。复杂一点的实现通常还会加入dialect 或 client 抽象屏蔽不同数据库方言或不同版本 client 的差异split serializer自定义高效的 split 序列化Kryo、Protobuf 等替代默认 Java 序列化enumerator state记录已分配/待分配分片的快照对象reader state 辅助类跟踪每个分片内部的消费进度offset / positionschema discoverer从数据源自动发现并推断表结构。什么时候用哪种设计什么时候简单 Reader 就够了适用于数据源天然单线程如 socket、HTTP 轮询等不需要并行没有明确的 split 模型。这种情况下可以不实现或仅实现空实现的split 与 enumeratorreader 在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),仅供参考