ARTICLE DETAIL

资讯详情

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

SeaTunnel RocketMQ 连接器全解析:Source 与 Sink 配置实战、事务写入与版本演进

SeaTunnel RocketMQ 连接器全解析:Source 与 Sink 配置实战、事务写入与版本演进 SeaTunnel RocketMQ 连接器全解析Source 与 Sink 配置实战、事务写入与版本演进【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelRocketMQ 是 SeaTunnel Connector-V2 体系中消息队列场景的高频连接器本文以其 变更日志 为时间骨架完整梳理 RocketMQ Source / Sink 的全部配置项、启动位置控制、多表读取、Tag 过滤、事务消息等核心能力并结合connector-rocketmq模块源码说明其底层实现原理与版本演进脉络。读完本文你将能够在 SeaTunnel 中独立配置从 RocketMQ 消费、向 RocketMQ 写入的端到端数据管道并理解每个关键行为背后对应的代码实现。连接器概览能力边界与适用引擎RocketMQ 连接器由 connector-rocketmq 模块实现其插件标识为Rocketmq定义于 RocketMqBaseOptions.java 的CONNECTOR_IDENTITY常量。该连接器以统一的代码同时承担 Source消费与 Sink生产两个角色支持以下运行环境与能力见 Source 文档 与 Sink 文档维度支持情况支持的 RocketMQ 版本Apache RocketMQ 4.9.0 及以上支持引擎SeaTunnel Zeta、Spark、FlinkSource 特性batch / stream 模式、exactly-once、并行度扩展SupportParallelism、多表读取、不支持列投影与用户自定义 splitSink 特性exactly-once事务消息、不支持 CDC 与定时 flushSource 端实现了SeaTunnelSource与SupportParallelism接口见 RocketMqSource.java其有界性由作业模式决定JobMode.BATCH对应BOUNDEDSTREAMING对应UNBOUNDED因此同一份 Source 配置既可跑批任务也可跑流任务。Source 端配置全解Source 端所有参数集中定义在 RocketMqSourceOptions.java与 Sink 共享的基础参数则位于父类 RocketMqBaseOptions.java。完整参数如下参数名类型必填默认值说明name.srv.addrString是-RocketMQ NameServer 地址例如localhost:9876topicsString否-逗号分隔的 Topic 列表如topic_a,topic_b与tables_configs、table_list三选一tables_configsList否-多表读取配置每项必须包含topics可含format、schema、tags、start.mode、start.mode.timestamp、start.mode.offsets、ignore_parse_errorstable_listList否-已弃用请改用tables_configstagsString否-逗号分隔的 Tag 列表只有 RocketMQ 消息 Tag 精确匹配才被消费acl.enabledBoolean否false是否启用 RocketMQ ACL 鉴权access.keyString否-Access Keyacl.enabled true时必填secret.keyString否-Secret Keyacl.enabled true时必填batch.sizeint否100每次拉取的最大消息条数consumer.groupString否SeaTunnel-Consumer-Group消费者组 IDcommit.on.checkpointBoolean否true是否在 SeaTunnel checkpoint 完成时提交 offsetschemaconfig否-消息 Schema省略时按文本读取消息体formatString否json消息格式支持json与textfield.delimiterString否,format text时的字段分隔符start.modeString否CONSUME_FROM_GROUP_OFFSETS起始消费位置详见下文start.mode.offsetsMap否-CONSUME_FROM_SPECIFIC_OFFSETS时必填Key 格式为topic-queueIdstart.mode.timestampLong否-CONSUME_FROM_TIMESTAMP时必填毫秒时间戳partition.discovery.interval.millislong否-1Topic/分区动态发现间隔毫秒-1关闭ignore_parse_errorsBoolean否false是否跳过无法解析的 JSON 消息而非让作业失败consumer.poll.timeout.millislong否5000拉取超时时间毫秒common-optionsconfig否-Source 公共参数见 Source Common Options启动位置start.mode的 5 种语义start.mode决定 Source 从何处开始消费其取值由 StartMode.java 枚举约束CONSUME_FROM_GROUP_OFFSETS默认从消费者组已提交的 offset 开始消费CONSUME_FROM_FIRST_OFFSET从最早的可用 offset 开始CONSUME_FROM_LAST_OFFSET从最新的可用 offset 开始CONSUME_FROM_TIMESTAMP从start.mode.timestamp非负毫秒时间戳对应的首个 offset 开始且该时间戳不能晚于作业运行时刻CONSUME_FROM_SPECIFIC_OFFSETS从start.mode.offsets指定的具体 offset 开始。对应的 HOCON 配置示例节选自 Source 文档start.mode CONSUME_FROM_SPECIFIC_OFFSETS start.mode.offsets { test_topic-0 50 }start.mode CONSUME_FROM_TIMESTAMP start.mode.timestamp 1667179890315从源码看CONSUME_FROM_TIMESTAMP的校验发生在配置解析阶段RocketMqSourceConfig.java 会将start.mode.timestamp与System.currentTimeMillis()比较大于当前时间则直接抛出IllegalArgumentException而CONSUME_FROM_SPECIFIC_OFFSETS则把形如test_topic-0的 Key 按最后一个-拆分为 topic 与 queueId构造MessageQueue后存入specificStartOffsets。消息格式json 与 textformat参数由 SchemaFormat.java 枚举定义仅支持json与text两种取值。当format json时需要配合schema声明消息体的字段类型SeaTunnel 使用JsonDeserializationSchema将 JSON 消息体解析为强类型字段若ignore_parse_errors true无法解析的 JSON 消息会被跳过而不是导致作业失败当format text时SeaTunnel 按field.delimiter切分消息体并按 schema 顺序映射到字段如果完全省略schema消息体会被当作单个文本值读取。该分支逻辑位于 RocketMqSourceConfig.java配置了schema时按format选择反序列化器未配置schema时则退化为按\002ASCII 2分隔的文本反序列化最后统一包装为带表 ID 的RocketMqTableIdDeserializationSchema。Tag 过滤的精确匹配语义tags使用逗号分隔列表如tag_a,tag_b连接器对拉取到的消息 Tag 做精确匹配因此不要在此使用 RocketMQ 的 Tag 表达式语法如tag_a || tag_b。在多表作业中每个tables_configs条目可设置各自的tags过滤。配置解析见 RocketMqSourceConfig.javaTag 列表经逗号切分、去空白、distinct()去重后保存在TopicTableConfig中。多表读取tables_configs当不同 Topic 的 schema 不同时应使用tables_configs进行多表读取。topics、tables_configs与已弃用的table_list三者互斥。规则如下每个条目必须包含topics可定义自己的schema、format、tags与启动位置条目中未设置的选项继承顶层默认值因此只需覆盖与主题相关的 schema、tags 或启动位置若条目未设置schema.table输出表名默认取 Topic 名条目使用CONSUME_FROM_TIMESTAMP时必须同时设置start.mode.timestamp使用CONSUME_FROM_SPECIFIC_OFFSETS时必须设置非空start.mode.offsets。源码层面RocketMqSourceConfig.java 优先读取tables_configs其次读取table_list逐条解析每个TopicTableConfig并聚合所有 Topic找不到多表配置时退回单表兼容模式。Source 配置示例读取 JSON 消息含 schema 定义env { parallelism 1 job.mode BATCH } source { Rocketmq { name.srv.addr rocketmq-e2e:9876 topics test_topic_json plugin_output rocketmq_table format json schema { fields { id bigint c_string string c_int int c_timestamp timestamp } } } } sink { Console { plugin_input rocketmq_table } }多表读取不同 Topic 不同 schema、不同 Tag、不同启动位置source { Rocketmq { name.srv.addr rocketmq-e2e:9876 start.mode CONSUME_FROM_LAST_OFFSET tables_configs [ { topics test_topic_multi_a start.mode CONSUME_FROM_FIRST_OFFSET format json schema { fields { id bigint c_string string } } }, { topics test_topic_multi_b start.mode CONSUME_FROM_FIRST_OFFSET tags tag_b format json schema { table rocketmq_multi_custom fields { id bigint description string } } } ] } }Sink 端配置全解Sink 端参数定义在 RocketMqSinkOptions.java完整参数如下参数名类型必填默认值说明topicString是-RocketMQ Topic 名name.srv.addrString是-RocketMQ NameServer 地址acl.enabledBoolean否false是否启用 ACL 鉴权access.keyString否-acl.enabled true时必填secret.keyString否-acl.enabled true时必填producer.groupString否SeaTunnel-Producer-Group生产者组 IDtagString否-随每条消息写入的 RocketMQ Tagpartition.key.fieldsList否-序列化为消息 Key 的字段名列表字段必须存在于上游 schemaformatString否json消息格式支持json与textfield.delimiterString否,format text时的字段分隔符producer.send.syncBoolean否false是否同步发送false为异步发送exactly.onceBoolean否false是否通过事务消息实现精确一次投递max.message.sizeint否4194304最大消息体大小字节源码常量DEFAULT_MAX_MESSAGE_SIZE 1024 * 1024 * 4send.message.timeoutint否3000发送超时毫秒源码常量DEFAULT_SEND_MESSAGE_TIMEOUT_MILLIS 3000common-optionsconfig否-Sink 公共参数见 Sink Common Optionspartition.key.fields控制消息 Key 与队列路由partition.key.fields决定 RocketMQ 消息 KeySeaTunnel 将配置字段的值序列化为 JSON 并写入Message.keys。非事务发送场景下同一 Key 会经过 RocketMQ 的 hash 队列选择器进入同一队列从而保证相同 Key 的行落入同一 queue未配置时由 RocketMQ 自行选择队列。例如输入 schema 包含c_int则如下配置用c_int构造消息 Keypartition.key.fields [c_int]exactly.once基于事务消息的精确一次写入Sink 通过 RocketMQ 事务消息支持精确一次写入默认关闭。开启exactly.once true前需确认 RocketMQ 集群与作业 checkpoint 配置已为事务写入做好准备。源码实现位于 RocketMqTransactionSender.java它构建TransactionMQProducer并注册TransactionListener——executeLocalTransaction与checkLocalTransaction均直接返回LocalTransactionState.COMMIT_MESSAGE通过sendMessageInTransaction提交消息这一设计与作业 checkpoint 配合将“预提交 确认提交”纳入 SeaTunnel 的精确一次保障链路。producer.send.sync同步与异步发送producer.send.sync true时生产者等待 RocketMQ 对每次发送返回确认默认false表示异步发送。Sink 组装逻辑见 RocketMqSink.java该处同时完成 Topic、Tag、ACL、生产者组、最大消息体、发送超时、分区 Key 字段、精确一次与同步标志的装配。Sink 配置示例写入带 Tag 的消息Source 为 FakeSource 造数source { FakeSource { row.num 10 schema { fields { c_string string c_int int } } } } sink { Rocketmq { name.srv.addr rocketmq-e2e:9876 topic test_topic_message_tag tag test_tag partition.key.fields [c_string] producer.send.sync true } }文本格式写入sink { Rocketmq { name.srv.addr rocketmq-e2e:9876 topic test_text_topic format text field.delimiter , producer.send.sync true } }端到端RocketMQ 读 RocketMQ 写STREAMING 模式下将test_topic_source的消息消费后写入test_topic_sink并开启精确一次env { parallelism 1 job.mode STREAMING checkpoint.interval 5000 } source { Rocketmq { name.srv.addr rocketmq-e2e:9876 topics test_topic_source plugin_output rocketmq_table format json start.mode CONSUME_FROM_FIRST_OFFSET consumer.group rocketmq_to_rocketmq_group schema { fields { id bigint c_string string } } } } sink { Rocketmq { plugin_input rocketmq_table name.srv.addr rocketmq-e2e:9876 topic test_topic_sink partition.key.fields [id] exactly.once true } }底层实现原理Split 枚举与动态发现RocketMQ Source 的并行读取基于 SeaTunnel 的 SourceSplitEnumerator 机制实现。核心类是 RocketMqSourceSplitEnumerator.java每个 RocketMQMessageQueuetopic queueId对应一个RocketMqSourceSplit枚举器负责把 split 分配给各并行 readerpartition.discovery.interval.millis控制 Topic/分区的动态发现间隔源码默认常量DEFAULT_DISCOVERY_INTERVAL_MILLIS 60 * 1000毫秒用户配置-1时关闭动态发现通过restoredSplits与isRestored标记区分“从 checkpoint 恢复的 split”与“新发现的 split”避免恢复的 split 起始 offset 被覆盖RocketMqSourceSplitEnumerator.java构造函数中设置rocketmq.client.logUseSlf4j true以避免 RocketMQ 客户端创建大量AsyncAppender-Dispatcher-Thread这一细节也与变更日志中线程相关修复一脉相承。从源码结构看连接器内部还划分了serializeSeaTunnelRowSerializer与DefaultSeaTunnelRowSerializer、sink普通生产者RocketMqNoTransactionSender与事务生产者RocketMqTransactionSender、source消费者线程RocketMqConsumerThread、Reader、Split 状态等等包职责边界清晰。版本演进从 2.3.2 到 2.3.12 的关键变更解读connector-rocketmq 变更日志 完整记录了该连接器自诞生以来的演进路线。逐条解读如下可与上文配置项一一对应2.3.2连接器诞生与基础能力[Feature][Connector-V2] Add rocketmq source and sink (#4007)RocketMQ 连接器首次引入同时提供 Source 与 Sink这是后续所有能力的地基[Hotfix][Connector-V2][RocketMQ] Fix rocketmq spark e2e test cases (#4583)修复 Spark 引擎下的 E2E 测试用例[Improve][pom] Formatting pom (#4761)工程层面的 pom 规范化整理。2.3.4稳定性与工程规范[Fix] [Connector] Rocketmq source startOffset greater than endOffset error (#6287)修复 Source 端起始 offset 大于结束 offset 时的错误——对应start.mode CONSUME_FROM_SPECIFIC_OFFSETS指定了越界 offset 的场景[Fix][connector-rocketmq]Fix a NPE problem when checkpoint.interval is set too small (#6624) (#6625)修复checkpoint.interval设置过小时出现的 NPE提醒流式作业中 checkpoint 频率与消费线程生命周期需要匹配[Test][E2E] Add thread leak check for connector (#5773)为连接器增加线程泄漏检查结合上文rocketmq.client.logUseSlf4j的设置可见该连接器对线程治理的重视[Improve][Common] Introduce new error define rule (#5793)、[Improve] Remove use SeaTunnelSink::getConsumedType method and mark it as deprecated (#5755)、[Improve][CheckStyle] Remove useless SuppressWarnings annotation of checkstyle (#5260)均为引擎层 API 语义收敛与代码规范改进。2.3.5checkpoint 联动修复[fix][connector-rocketmq]Fix a NPE problem when checkpoint.interval is set too small(#6624) (#6625)2.3.5 版本再次强化对过小 checkpoint 间隔的防御该问题与commit.on.checkpoint的提交时机直接相关。2.3.6offset 提交与多表读取[Feature][Kafka] Support multi-table source read (#5992)多表 Source 读取能力落地该特性同时惠及 Kafka 与 RocketMQ 类消息连接器对应tables_configs参数[Fix][connector-rocketmq] commit a correct offset to broker reduce ThreadInterruptedException log (#6668)向 broker 提交正确 offset并减少ThreadInterruptedException的日志噪音——修正了 checkpoint 提交与线程中断在消费侧的竞态表现。2.3.8异常兜底[Fix][Connector-V2] Fix some throwable error not be caught (#7657)修复部分Throwable未被捕获的问题避免偶发异常导致作业状态异常属于连接器健壮性补强。2.3.9Tag 与可观测性[Improve][Connector-V2] RocketMQ Sink add message tag config (#7996)Sink 端新增tag配置对应 RocketMqSinkOptions.java[Feature][Restapi] Allow metrics information to be associated to logical plan nodes (#7786)RestAPI 允许将指标信息关联到逻辑计划节点增强作业可观测性。2.3.10Source Tag 与容错消费[Improve][Connector-V2] RocketMQ Source add message tag config (#8825)Source 端新增tags配置与 2.3.9 的 Sink 端tag形成闭环消费侧可通过 RocketMqSourceOptions.java 的tags做精确匹配过滤[Improve][Connector-V2] Add optional flag for rocketmq connector to skip parse errors instead of failing (#8737)新增ignore_parse_errors可选标志让连接器跳过解析失败的消息而不是让整个作业失败对应 RocketMqSourceOptions.java。2.3.11状态类序列化规范[Feature][Checkpoint] Add check script for source/sink state class serialVersionUID missing (#9118)新增对 Source/Sink 状态类缺少serialVersionUID的检查脚本。RocketMQ 侧的RocketMqSourceState等状态类需要保证serialVersionUID稳定这是分布式作业 checkpoint 恢复的正确性前提。2.3.12枚举器性能与配置优化[Improve][API] Optimize the enumerator API semantics and reduce lock calls at the connector level (#9671)优化SourceSplitEnumeratorAPI 语义减少连接器层面的锁调用直接受益者正是 RocketMqSourceSplitEnumerator.java 中 split 分配路径的并发开销[improve] rocketmq options (#9251)对 RocketMQ 相关选项定义做整体改进沉淀为当前版本的参数体系。质量保障单元测试与 E2E 验证连接器配套了完整的测试体系可在仓库中直接查阅单元测试位于 connector-rocketmq 测试目录覆盖RocketMqFactoryTest工厂与插件名注册、RocketMqSourceConfigTest配置解析含 start.mode 各分支与RocketMqSourceSplitEnumeratorTestsplit 分配与恢复逻辑端到端测试位于 RocketMqIT.java配合RocketMqContainer以 Testcontainers 拉起真实 RocketMQ 4.x 环境验证读写全链路变更日志中 2.3.2 的 Spark E2E 修复与 2.3.4 的线程泄漏检查均以此为基础。总结SeaTunnel RocketMQ 连接器从 2.3.2 首次引入到 2.3.12 的持续打磨已经形成一套覆盖消费5 种启动位置、Tag 精确过滤、多表读取、容错解析、生产事务消息、同步/异步发送、分区 Key、消息 Tag与工程治理offset 提交、checkpoint 联动、状态序列化检查、线程治理的完整方案。实践中建议消费侧按需选用start.mode与ignore_parse_errors多 schema 场景优先使用tables_configs写入侧在需要精确一次时开启exactly.once并配合合理的 checkpoint 间隔需要顺序性时利用partition.key.fields控制队列路由。相关配置示例与源码实现均可在本文引用的仓库路径中进一步查阅。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表