ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SeaTunnel Connector-V2 插件贡献开发指南:新架构、核心 API 与源码级实现

SeaTunnel Connector-V2 插件贡献开发指南:新架构、核心 API 与源码级实现 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载SeaTunnelApache SeaTunnel新一代 Connector-V2 体系为连接器定义了与旧版完全不同的开发范式通过seatunnel-api中一套与引擎无关的 Source/Sink 接口配合翻译层Translation Layer实现 Flink、Spark、Zeta 多引擎的平滑迁移。本文以仓库中 贡献 Connector-V2 插件入口文档 及其指向的 Connector-V2 贡献指南 为主体结合seatunnel-api的源码实现完整讲解从工程结构、本地调试、新建插件的完整步骤到 Source/Sink/Factory/Option 四大 API 的设计原理与实战写法。读完本文你将具备在 SeaTunnel 仓库中独立开发、调试并合入一个全新 Connector-V2 插件的完整能力。一、从贡献入口说起Connector-V2 指南的位置在 docs/en/contribution/contribute-plugin.md以及对应的 中文版中官方明确将「贡献 Connector-V2 插件」的完整指南指向seatunnel-connectors-v2/README.md——这份 README 不是泛泛的模块介绍而是新接口、新代码结构、调试方法与 API 设计的权威说明书。其开篇即点明目的本文介绍基于新设计 API 的 Connector 新接口与新代码结构帮助开发者快速理解 API 与转换层transformation layer的改进同时指导贡献者如何使用新 API 开发新连接器。这意味着贡献 Connector-V2 插件本质上是实现一套与引擎解耦的插件 API。Connector 不再为某个引擎单独编写而是通过seatunnel-api定义的标准接口 seatunnel-translation翻译层一套代码适配多引擎运行。二、新架构下的工程结构Connector 代码该放哪里为了与旧代码隔离、降低并行开发时的合并难度SeaTunnel 将 Connector-V2 相关代码拆分到独立模块。指南给出的工程结构对应到本仓库如下模块职责仓库相对路径seatunnel-connectors-v2所有 Connector-V2 连接器实现seatunnel-connectors-v2seatunnel-translationConnector-V2 的翻译层多引擎适配seatunnel-translationseatunnel-transforms-v2Transform V2 转换器实现seatunnel-transforms-v2seatunnel-e2e/seatunnel-connector-v2-e2eConnector V2 端到端E2E测试seatunnel-e2e/seatunnel-connector-v2-e2eseatunnel-core/seatunnel-flink-starterFlink 引擎下的新启动模块seatunnel-core/seatunnel-flink-starterseatunnel-core/seatunnel-spark-starterSpark 引擎下的新启动模块seatunnel-core/seatunnel-spark-starterseatunnel-core/seatunnel-starterSeaTunnel Zeta 引擎启动模块seatunnel-core/seatunnel-starter指南特别说明seatunnel-translation通过适配不同引擎的接口完成 SeaTunnel API 与引擎 API 之间的「翻译」从而让一套 Connector 能在多个不同引擎上运行——这是 Connector-V2 与旧版连接器最本质的架构差异。三、本地调试 ConnectorSeaTunnelApiExample 五步法指南推荐使用seatunnel-examples下两个可在本地直接运行的示例程序来调试 Connector一个运行在 Flink 引擎一个运行在 Spark 引擎。这也是 Connector 本地开发中最常用的调试方式——直接调试示例的main方法可以直观理解程序运行逻辑。示例所用的配置文件保存在resources/examples目录。为你的新 Connector 添加本地调试示例按以下 5 步操作添加被测 Connector 依赖将待测 connector 的groupId、artifactId、version加入 Flink 示例的pom.xml若要在 Spark 引擎运行则加入 Spark 示例的pom.xml作为依赖。补充 test/provided 依赖检查你 Connector 的pom.xml中scope为test或provided的依赖同样加入示例模块的pom.xml并将其scope改为compile保证运行时类路径完整。添加任务配置文件在resources/examples下新增任务配置文件即 SeaTunnel 的source/transform/sink三段式配置。配置主方法在示例类如SeaTunnelApiExample的main方法中指向该配置文件。直接运行 main 方法即可在本地 IDE 中单步调试你的 Connector。如果你更倾向于自动化验证本仓库还提供了完整的 E2E 体系seatunnel-e2e/seatunnel-connector-v2-e2e 下为每个主流连接器如 connector-kafka-e2e、connector-jdbc-e2e建立了独立的端到端测试模块新插件合入时同样可以在该目录下补充 e2e 用例。四、新建 Connector-V2 插件的完整步骤指南给出了创建新连接器的标准流程共 5 个关键动作第 1 步创建模块。在seatunnel-connectors-v2目录下新建模块命名为connector-{连接器名}例如connector-kafka、connector-console。第 2 步编写 pom。直接参考已有连接器的pom.xml例如 connector-console/pom.xml并把当前子模块添加到父模块seatunnel-connectors-v2/pom.xml的modules中。第 3 步创建 source/sink 两个包。分别建立org.apache.seatunnel.connectors.seatunnel.{连接器名}.source org.apache.seatunnel.connectors.seatunnel.{连接器名}.sink第 4 步注册插件映射。在仓库根目录的 plugin-mapping.properties 中添加连接器信息。该文件用于在用户配置中根据插件名解析对应的 Jar 包artifactId。文件顶部有重要警告注释seatunnel.source.XXX中的XXX必须是SeaTunnelSource::getPluginName和TableSinkFactory::factoryIdentifier返回的字符串值。仓库中已有的映射示例如seatunnel.source.FakeSource connector-fake seatunnel.sink.Console connector-console seatunnel.sink.Assert connector-assert seatunnel.source.Kafka connector-kafka seatunnel.sink.Kafka connector-kafka第 5 步加入发布包。在 seatunnel-dist/pom.xml 中添加上述连接器依赖这样连接器 Jar 才能被打进二进制发行包用户在安装目录的plugins下找到它。五、SeaTunnel Source API 源码解析Source 侧共四个核心接口全部位于 seatunnel-api/src/main/java/org/apache/seatunnel/api/source。SeaTunnel 的 API 设计借鉴了 Flink 的设计理念。5.1 SeaTunnelSource驱动端的工厂类SeaTunnelSource.java 是 Source 的总入口注释明确说明它像工厂类一样帮助构建 SourceSplitEnumerator、SourceReader 及对应的序列化器在驱动端driver/master执行。其关键方法getBoundedness()L49返回Boundedness决定当前 Source 是流式UNBOUNDED还是批式BOUNDED。SeaTunnel 的 Source 采用流批一体设计可通过配置动态指定即同一个 Source 既可作为流也可作为批。getProducedCatalogTables()L67返回输出的CatalogTable列表包含比旧getProducedType()更完整的 schema 信息是官方推荐的 schema 获取方式。Connector 可以硬编码固定 schema也可以让用户通过 config 自定义 schema后者被推荐。createReader(Context)L79创建SourceReader。getSplitSerializer()L87Split 序列化器用于跨进程传输 Enumerator 生成的 Split。createEnumerator(Context)L99与restoreEnumerator(Context, checkpointState)L111分别用于启动时创建 Enumerator 与从 checkpoint 恢复时重建 Enumerator。getEnumeratorStateSerializer()L120Enumerator 状态序列化器。当前 SeaTunnelSource 支持的数据类型必须是SeaTunnelRow。5.2 SourceSplitEnumerator切分与调度者SourceSplitEnumerator.java 运行在 master 端负责获取数据读取分片SourceSplit并分配到不同 SourceReader。它继承CheckpointListener因而天然具备 checkpoint 回调能力。关键方法run()L42由引擎只执行一次用于生成 SourceSplit并通过Context.assignSplit(subtaskId, splits)L89把分片分发给对应的 SourceReader。addSplitsBack(splits, subtaskId)L58当 SourceReader 异常或重启导致 Split 无法正常处理时把这些 Split 收回并重新分配。registerReader(subtaskId)L64run()之后注册进来的新 Reader。如果此时没有已分配的 Split可以分配给这些新 Reader——因此大多数情况下需要在 Enumerator 中自行维护 Split 分配记录。handleSplitRequest(subtaskId)L62当某个 Reader 主动向 Enumerator 请求 Split 时被调用可配合Context.assignSplit将分片发给对应 Reader。snapshotState(checkpointId)L67流处理中周期性返回需要保存的当前状态恢复时会调用SeaTunnelSource.restoreEnumerator重建 Enumerator 并注入已保存状态。notifyCheckpointComplete状态成功保存后的后续处理回调可用于在三方存储中保存状态或打标记。5.3 SourceSplit分片信息的载体SourceSplit.java 是保存分片的接口只定义了一个方法splitId()L30。不同的分片需要定义不同的 splitId。你可以实现该接口来保存分片需要携带的信息例如 Kafka 的 partition 与 topic、HBase 的 column family 等这些信息被 SourceReader 用来确定该读取总数据中的哪一部分。5.4 SourceReader与数据源直接交互的读取者SourceReader.java 是与数据源直接交互的接口运行在 worker 端从数据源读取数据的动作由实现该接口的类完成。关键方法pollNext(CollectorT)L52Reader 的核心。实现从数据源读取数据并返回给 SeaTunnel 的过程。每当准备把数据交给 SeaTunnel 时调用参数中的Collector.collect方法该方法可被无限次调用以完成大批量数据读取。由于 Source 是流批一体设计批模式下 Connector 必须自行决定何时结束读取——读完一批比如 100 条后需调用Context.signalNoMoreElement()L97通知 SeaTunnel 没有更多数据这批数据即可用于批处理。流处理没有此要求因此大多数流批一体的 SourceReader 都会有如下代码if (Boundedness.BOUNDED.equals(context.getBoundedness())) { // signal to the source that we have reached the end of the data. context.signalNoMoreElement(); break; }即仅在批模式下通知 SeaTunnel。addSplits(splits)L70框架用它把 SourceSplit 分配给不同的 SourceReader。Reader 应保存拿到的分片随后在pollNext中读取对应分片数据但也存在 Reader 暂时没有分片的情况Split 尚未生成或确实未分配此时pollNext应做相应处理例如继续等待。handleNoMoreSplits()L79被触发时表示不再有分片Connector Source 可以选择性地做出反馈。snapshotState(checkpointId)L63流处理中周期性返回需要保存的状态即分片信息SeaTunnel 将分片信息与状态一起保存以实现动态分配。notifyCheckpointComplete/notifyCheckpointAborted对应 checkpoint 不同状态的回调。小结 Source 数据流SeaTunnelSource驱动端→createEnumerator生成SourceSplitEnumeratormaster→ Enumerator 通过assignSplit把SourceSplit分发给各并行子任务的SourceReaderworker→ Reader 在pollNext中消费分片数据并通过Collector.collect上抛。整个流程与 Flink 的 Source 体系一一对应。六、SeaTunnel Sink API 源码解析Sink 侧的核心接口位于 seatunnel-api/src/main/java/org/apache/seatunnel/api/sink。6.1 SeaTunnelSink写入目标的定义入口SeaTunnelSink.java 用于定义数据写入目标端的方式并从中获得SinkWriter、SinkCommitter等实例。其关键方法createWriter(Context)L82创建SinkWriter。restoreWriter(Context, states)L84默认直接调用createWriter需要从状态恢复时重写。createCommitter()L104与createAggregatedCommitter()L124分别返回SinkCommitter与SinkAggregatedCommitter均为Optional可只实现其一或两者都实现。配套序列化器getWriterStateSerializer()、getCommitInfoSerializer()、getAggregatedCommitInfoSerializer()用于跨进程传递 Writer 状态、提交信息与聚合提交信息。Sink 侧的重要特性是分布式事务处理。SeaTunnel 定义了两个不同的 CommitterSinkCommitter处理不同子任务subTask的事务。SinkAggregatedCommitter在单个节点上统一处理所有节点的事务结果可避免二阶段提交第二阶段失败导致的状态不一致问题。SinkAggregatedCommitter的combine()方法用于聚合SinkWriter.prepareCommit返回的事务信息生成聚合后的事务信息。6.2 SinkWriter直接写入数据SinkWriter.java 直接与输出端交互把 SeaTunnel 从数据源获得的数据交给 Writer 写入。关键方法write(element)L47负责把数据交给SinkWriter。可以选择直接写入也可以缓冲一定量后再写入。目前仅支持SeaTunnelRow数据类型。prepareCommit()L65在提交前执行。可直接在此写数据也可以在 2PC 中实现第一阶段phase one第二阶段在SinkCommitter或SinkAggregatedCommitter中实现。该方法返回的提交信息commit info会交给SinkCommitter和SinkAggregatedCommitter用于下一阶段的处理。snapshotState(checkpointId)L71返回 Writer 的状态。abortPrepare()L81用于中止prepareCommit的副作用。若prepareCommit失败则没有 CommitInfoT无法通过SinkCommitter回滚此时可用该方法回滚目前仅在 Spark 引擎使用。6.3 实现 SinkCommitter 还是 SinkAggregatedCommitter指南给出了明确的推荐结论当前版本优先推荐实现SinkAggregatedCommitter它能在 Flink/Spark 中提供强一致性保证。同时commit 必须做到幂等idempotent以保证引擎重试能正常工作。七、工厂与参数体系TableSourceFactory / TableSinkFactory为了让框架自动创建Source Connector、Sink Connector 与 Transform ConnectorConnector 需要返回创建它们所需的参数以及每个参数的校验规则。为此 SeaTunnel 定义了TableSourceFactory与TableSinkFactory官方建议将它们放在SeaTunnelSource/SeaTunnelSink实现类的同目录下方便检索。两个 Factory 共同继承自 Factory.java其中声明了两个核心方法factoryIdentifier()L32表示当前 Factory 的名称。该值必须与getPluginName返回的值保持一致这样未来使用 Factory 创建 Source/Sink 时可以无缝切换。optionRule()L42返回参数规则用于声明 Connector 支持哪些参数、哪些必填、哪些可选、哪些互斥exclusive、哪些需要捆绑必填bundledRequired。它既用于可视化Web-UI创建 Connector 逻辑也用于根据用户配置生成完整的参数对象从而使 Connector 开发者无需在 Config 中逐一判断参数是否存在。注意createSink和createSource是创建 Source/Sink 的方法目前不需要自行实现框架已通过 Factory 上下文自动完成实例装配。7.1 参考实现ConsoleSinkFactory仓库中最简洁的 Factory 参考实现是 ConsoleSinkFactory.javaAutoService(Factory.class) public class ConsoleSinkFactory implements TableSinkFactory { public static final OptionBoolean LOG_PRINT_DATA Options.key(log.print.data) .booleanType() .defaultValue(true) .withDescription( Flag to determine whether data should be printed in the logs.); public static final OptionInteger LOG_PRINT_DELAY Options.key(log.print.delay.ms) .intType() .defaultValue(0) .withDescription( Delay in milliseconds between printing each data item to the logs.); Override public String factoryIdentifier() { return Console; } Override public OptionRule optionRule() { return OptionRule.builder().build(); } Override public TableSink createSink(TableSinkFactoryContext context) { ReadonlyConfig options context.getOptions(); return () - new ConsoleSink( context.getCatalogTable().getTableSchema().toPhysicalRowDataType(), options); } }该示例展示了几个关键实践AutoService(Factory.class)注解L31必不可少这是 Java SPI 自动注册机制Factory是TableSourceFactory与TableSinkFactory的父接口。忘记加注解插件将无法被框架发现。factoryIdentifier()返回ConsoleL49与ConsoleSink.getPluginName()、plugin-mapping.properties 中的seatunnel.sink.Console connector-console三者保持一致构成完整的插件解析链路。Option 定义直接内联在 Factory 中log.print.data默认true与log.print.delay.ms默认0两个参数通过Options.key(...).booleanType()/intType().defaultValue(...).withDescription(...)声明式构建。另外指南还提到两个可参考的现有实现org.apache.seatunnel.connectors.seatunnel.elasticsearch.source.ElasticsearchSourceFactory位于 connector-elasticsearch——许多 Source 都支持配置 Schema因此使用了一个公共的 Option如果需要 schema可参考org.apache.seatunnel.api.table.catalog.CatalogTableUtil.SCHEMA即 CatalogTableUtil.java。八、Options 与 OptionMark声明式的参数体系实现TableSourceFactory和TableSinkFactory时会同步创建对应的Option。每个Option对应一个配置项不同配置具有不同类型常见类型booleanType()、intType()、stringType()等可直接通过Options.key(...)链式方法创建。如果参数类型是对象可以使用 POJO 来表示对象类型的参数并在每个字段上使用 OptionMark.java 注解标记这是子 Option。OptionMark有两个参数nameL35声明字段对应的参数名。为空时默认将 Java 小驼峰转换为下划线风格例如myUserPassword→my_user_password。大多数情况下默认留空即可。descriptionL38当前参数的描述可选建议与文档保持一致。具体示例可参考org.apache.seatunnel.connectors.seatunnel.assertion.sink.AssertSinkFactory位于 connector-assert。九、启动类与多引擎适配除旧版启动类外SeaTunnel 新建了两个启动模块seatunnel-core/seatunnel-flink-starter与seatunnel-core/seatunnel-spark-starter本仓库中分别细分为 seatunnel-flink-13-starter、seatunnel-flink-15-starter 以及 seatunnel-spark-2-starter、seatunnel-spark-3-starter。你可以在这里找到「如何把配置文件解析成可执行的 Flink/Spark 进程」的完整实现。而seatunnel-api注意区别于旧模块seatunnel-apis则存放 SeaTunnel API 定义的新接口——通过实现这些接口开发者就能完成支持多引擎的 SeaTunnel Connector。翻译层位于 seatunnel-translation它通过适配不同引擎的接口实现 SeaTunnel API 与引擎 API 之间的转换让 Connector 支持在多个不同引擎上运行。如果你感兴趣可以阅读该模块代码并帮助改进。十、插件落地的最终检查清单结合指南的「Result」章节与仓库实际约束新插件合入前请对照以下清单代码位置所有 Connector 实现必须位于seatunnel-connectors-v2目录下可参考该模块内现有连接器如 connector-fake、connector-kafka作为范例。命名与包结构模块名为connector-{名称}包结构为org.apache.seatunnel.connectors.seatunnel.{名称}.source与...sinkSource 实现必须返回SeaTunnelRow类型数据。插件映射在 plugin-mapping.properties 中登记seatunnel.source.XXX/seatunnel.sink.XXX且XXX必须与getPluginName()/factoryIdentifier()返回值完全一致。发布集成在 seatunnel-dist/pom.xml 添加依赖确保 Jar 进入二进制发行包。Factory 注册实现类上必须标注AutoService(Factory.class)实现optionRule()声明参数校验规则对象类型参数用OptionMark标注。事务一致性Sink 优先实现SinkAggregatedCommitter以获得强一致性commit 必须幂等以支持引擎重试。测试与示例可参考seatunnel-examples中的本地调试示例为插件补充可运行示例并在 seatunnel-e2e/seatunnel-connector-v2-e2e 下补充 E2E 测试用例。通过这套流程开发出的 Connector将同时具备「多引擎可移植性」Source/Sink API 与引擎解耦 翻译层适配、「声明式参数校验」Factory Option/OptionMark与「一致性保证」两阶段提交体系三大特性这也是 SeaTunnel Connector-V2 相较于旧版连接器最核心的能力升级。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel Connector-V2 插件贡献指南从 FAQ 导航层到源码级实现SeaTunnel Connector V2 插件贡献指南从 FAQ 导航层到源码级实现 Connector V2 是 SeaTunnel 面向多引擎解耦设计数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Transform-V2 插件贡献指南从贡献入口到源码级实现SeaTunnel Transform V2 插件贡献指南从贡献入口到源码级实现 导读 本文面向想要为 SeaTunnel 贡献 Transform V2 插数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Transform-V2 插件贡献指南从入门路径到核心契约、源码实现与 E2E 验证SeaTunnel Transform V2 插件贡献指南从入门路径到核心契约、源码实现与 E2E 验证 本篇基于 SeaTunnel 仓库中的 Transf数据集成ETL大数据批处理流处理变更数据捕获上一篇AltStore存储优化终极指南快速清理缓存与冗余数据的5个技巧下一篇10分钟生成专业短视频MoneyPrinterTurbo如何颠覆传统视频创作创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表