
SeaTunnel RocketMQ 源连接器使用指南启动位点、多表读取与 Source 端实现解析【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本篇技术指南以 SeaTunnel 的 RocketMQ 源连接器Rocketmqsource为主线覆盖其全部配置选项、五种启动位点模式、JSON/Text 消息解析、Tag 过滤与tables_configs多 Topic 读取等核心能力并结合连接器源码解析 Split 分配、消费者线程拉取与检查点提交位点的底层机制读完后可直接在 SeaTunnel 作业中正确配置 RocketMQ 数据源并理解其容错行为。连接器概览SeaTunnel 的 RocketMQ 源连接器从 Apache RocketMQ 4.9.0 或更高版本的主题中读取消息。连接器支持两种读取形态单表模式通过topics读取一个或多个共享同一 Schema 的主题多表模式通过tables_configs读取多个 Schema 各不相同的主题每个条目可独立定义schema、format、tags与启动位点。该连接器支持 Spark、Flink 与 SeaTunnel Zeta 三种引擎。对照 连接器特性说明其当前能力矩阵为特性支持情况批处理batch支持流处理stream支持exactly-once 语义支持列裁剪column projection不支持并行度parallelism支持用户自定义 split不支持多表读取multiple table read支持从源码结构看topics、tables_configs与已废弃的table_list三个选项被声明为互斥mutually exclusive只能配置其中一项。工厂类 RocketMqSourceFactory 通过OptionRule中的.exclusive(...)规则在作业提交前强制执行这一约束违反时会抛出选项校验异常。配置选项全解以下选项表完整继承自 官方文档并结合 RocketMqSourceOptions 与 RocketMqBaseOptions 中的定义核对过默认值。名称类型必填默认值说明name.srv.addrString是-RocketMQ NameServer 地址例如localhost:9876topicsString否-以逗号分隔的主题列表例如topic_a,topic_b。topics、tables_configs、table_list三者只配其一tables_configsList否-多表读取配置。每项必须包含topics可选format、schema、tags、start.mode、start.mode.timestamp、start.mode.offsets、ignore_parse_errorstable_listList否-已废弃请使用tables_configs替代tagsString否-以逗号分隔的 Tag 列表。仅消费 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-GroupRocketMQ 消费组 IDcommit.on.checkpointBoolean否trueSeaTunnel 检查点完成后是否向 Broker 提交位点schemaconfig否-消息 Schema参见 Schema 特性。省略时消息体按文本读取formatString否json消息格式支持json与textfield.delimiterString否,format text时使用的字段分隔符start.modeString否CONSUME_FROM_GROUP_OFFSETS启动位点取值见下文“启动位点”一节start.mode.offsetsMap否-start.mode CONSUME_FROM_SPECIFIC_OFFSETS时必填键格式为topic-queueId例如test_topic-0start.mode.timestampLong否-start.mode 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除上述选项外RocketMqSourceFactory#optionRule()还内置了若干条件校验start.mode CONSUME_FROM_TIMESTAMP时start.mode.timestamp必填且必须 0start.mode CONSUME_FROM_SPECIFIC_OFFSETS时start.mode.offsets必填且不能为空 Mapacl.enabled true时access.key与secret.key均为必填。这些规则由单元测试 RocketMqFactoryTest 逐项验证负数时间戳、空 offsets Map、多表条目缺少topics、多表条目在时间戳模式下未配置start.mode.timestamp等场景均会抛出OptionValidationException。启动位点start.modestart.mode控制源连接器从何处开始读取五种取值语义如下取值语义CONSUME_FROM_GROUP_OFFSETS默认从消费组已提交的位点开始CONSUME_FROM_FIRST_OFFSET从最早可用位点开始CONSUME_FROM_LAST_OFFSET从最新可用位点开始CONSUME_FROM_TIMESTAMP从start.mode.timestamp对应时间点该时刻或之后的第一个位点开始CONSUME_FROM_SPECIFIC_OFFSETS从start.mode.offsets指定的位点开始使用CONSUME_FROM_TIMESTAMP时start.mode.timestamp必须是非负毫秒时间戳且不能晚于作业运行时的当前时间——这一约束在源码 RocketMqSourceConfig#buildConsumerMetadata 中通过System.currentTimeMillis()做了运行时检查超限会抛出IllegalArgumentException。start.mode CONSUME_FROM_SPECIFIC_OFFSETS start.mode.offsets { test_topic-0 50 }start.mode CONSUME_FROM_TIMESTAMP start.mode.timestamp 1667179890315从源码实现看位点解析集中在 RocketMqSourceSplitEnumerator#setPartitionStartOffset 中几个值得注意的行为GROUP_OFFSETS 的回退逻辑连接器通过 Admin 客户端查询消费组在各队列上的已提交位点若查询结果为空例如该消费组从未提交过位点会自动回退为CONSUME_FROM_FIRST_OFFSET从最早位点开始消费而不是直接报错SPECIFIC_OFFSETS 的键解析start.mode.offsets的键按最后一个-拆分为主题名与队列 ID例如test_topic_source-0解析为主题test_topic_source的队列 0多表模式下的按主题生效位点setPartitionStartOffset会先通过getEffectiveStartMode(topic)查找该主题的独立start.mode未单独配置时才回落到全局默认值这正是多表模式中“每项只覆盖差异项”的底层实现。消息格式与解析format / schema / ignore_parse_errorsformat json需要配合schema定义字段SeaTunnel 使用JsonDeserializationSchema把 JSON 消息体解析为类型化字段。ignore_parse_errors true时无法解析的 JSON 消息会被跳过而不是让作业失败format text消息体按field.delimiter拆分并按 Schema 字段顺序一一映射省略schema消息体被当作单个文本值读取。从源码 RocketMqSourceConfig#buildDeserialization 看此时会构造一个使用\002作为分隔符的TextDeserializationSchema由于该分隔符不会出现在普通文本消息体中等效于将整体消息体作为一个文本字段输出。所有解析结果最终都会被一层RocketMqTableIdDeserializationSchema包装为每条行数据附加其来源主题的表标识——这是多表读取模式下下游按表路由数据的关键机制见下文。Tag 过滤tags使用逗号分隔的普通列表例如tag_a,tag_b。连接器在拉取到消息后将消息的实际 Tag 与配置值做精确匹配比较因此不要在这里使用 RocketMQ 的 Tag 表达式语法如tag_a || tag_b。在tables_configs多表作业中每个条目可以设置自己的tags过滤器。源码印证parseTags会把逗号分隔的字符串拆分为去重后的列表存入每个主题的TopicTableConfig随后 RocketMqSourceReader 在pollNext中逐条执行tags.contains(record.getTags())判断只有匹配的消息才会进入反序列化并输出且被过滤掉的消息不会推进该队列的有效输出位点。E2E 测试中的rocketmq-source_text_error_tag_to_console.conf场景即验证了错误 Tag 消息被正确过滤的行为。多表读取tables_configs当不同主题对应不同 Schema 时使用tables_configs。每个条目必须包含topics并可定义自己的schema、format、tags与启动位点条目中未设置的选项继承顶层默认值因此每个条目只需覆盖该主题特有的差异。关键约束与默认行为topics、tables_configs与废弃的table_list三者互斥若schema.table未设置输出表名默认为主题名源码中通过TablePath.DEFAULT判断后以首个主题名回填表标识某条目使用start.mode CONSUME_FROM_TIMESTAMP时必须同时设置start.mode.timestamp使用start.mode CONSUME_FROM_SPECIFIC_OFFSETS时必须设置非空的start.mode.offsets。这些规则由 RocketMqSourceFactory 内部的TableConfigsValidator在提交前逐条校验校验失败会给出带条目下标的明确报错。任务示例以下五个示例完整继承自 官方文档示例中的rocketmq-e2e:9876为 E2E 测试环境中的 NameServer 地址实际使用时替换为你自己的 NameServer 即可。示例一读取 JSON 消息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 } }示例二读取带 Tag 过滤的 Text 消息env { parallelism 1 job.mode BATCH } source { Rocketmq { name.srv.addr rocketmq-e2e:9876 topics test_topic_text plugin_output rocketmq_table format text field.delimiter , tags tag_a,tag_b schema { fields { id bigint content string } } } } sink { Console { plugin_input rocketmq_table } }示例三从指定位点开始读取env { parallelism 1 job.mode BATCH } source { Rocketmq { name.srv.addr rocketmq-e2e:9876 topics test_topic_source plugin_output rocketmq_table format json start.mode CONSUME_FROM_SPECIFIC_OFFSETS start.mode.offsets { test_topic_source-0 50 } schema { fields { id bigint } } } } sink { Console { plugin_input rocketmq_table } }示例四多主题不同 Schema 读取env { parallelism 1 job.mode BATCH } 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 { Console {} }注意该示例中第二个条目通过schema.table rocketmq_multi_custom覆盖了默认以主题名作为输出表名的行为。示例五从时间戳开始读取env { parallelism 1 job.mode BATCH } source { Rocketmq { name.srv.addr rocketmq-e2e:9876 topics test_topic_source plugin_output rocketmq_table format json start.mode CONSUME_FROM_TIMESTAMP start.mode.timestamp 1667179890315 schema { fields { id bigint } } } } sink { Console { plugin_input rocketmq_table } }源码结构Split 分配与容错机制从源码结构看整个源端遵循 SeaTunnel V2 Source 的 Enumerator/Reader/Split 三段式架构核心文件位于 connector-rocketmq 的 source 包Split 发现与分配RocketMqSourceSplitEnumerator通过 Admin 工具类 RocketMqAdminUtil 查询各主题的MessageQueue及 min/max offset每个队列生成一个RocketMqSourceSplit分配时通过getSplitOwner以队列 ID 做哈希取模把队列分散到不同并行度实例上动态发现open()中会启动一个名为RocketMq-messageQueue-dynamic-discovery的守护线程周期重新发现队列。从源码看当partition.discovery.interval.millis配置为非正值时枚举器会将其替换为内置的 60 秒默认间隔再启动定时任务与文档“-1 禁用动态发现”的描述存在差异实际使用时建议以 源码行为 为准恢复场景检查点恢复后引擎会通过addSplitsBack把携带正确startOffset的 Split 归还枚举器源码用restoredSplits集合标记这些队列并在setPartitionStartOffset中跳过位点重置避免把持久化位点覆盖为 Broker 当前值。消息拉取与输出RocketMqSourceReader每个被分配的MessageQueue对应一个独立的RocketMqConsumerThread拉取线程每次poll以consumer.poll.timeout.millis为超时单批最多拉取batch.size条拉取结果先按 topic/broker/queueId 校验队列归属再做 Tag 精确匹配过滤最后交给按主题路由的反序列化 Schema批处理模式下当消费到queueOffset endOffset即停止并调用signalNoMoreElement通知引擎数据流结束。位点提交Reader 在snapshotState(checkpointId)中把各 Split 的当前startOffset快照到checkpointOffsetscommit.on.checkpoint true默认时notifyCheckpointComplete收到引擎回调后通过 RocketMQ 客户端的OffsetStore.updateOffset persist把位点写回 Broker。这正是文档中commit.on.checkpoint选项“检查点完成后提交位点”的实现也是连接器提供 exactly-once 语义与具备事务能力的 Sink 配合的位点基础。版本演进与 E2E 验证该源连接器自 SeaTunnel 2.3.2 引入初始提交同时包含 source 与 sink完整变更记录见 connector-rocketmq 变更日志。几个与本篇主题相关的节点2.3.10源连接器增加tags消息标签过滤#8825、增加ignore_parse_errors跳过解析失败消息#87372.3.6修复向 Broker 提交错误位点的问题#66682.3.12RocketMQ 选项体系优化#9251。端到端测试 RocketMqIT 基于 Testcontainers 拉起真实 RocketMQ 容器覆盖本文各场景五种启动位点分别对应rocketmq/目录下的rocketmq_source_earliest_to_console.conf、rocketmq_source_latest_to_console.conf、rocketmq_source_group_offset_to_console.conf、rocketmq_source_timestamp_to_console.conf、rocketmq_source_specific_offsets_to_console.conf五个作业配置此外还有 Tag 过滤、文本格式、检查点恢复rocketmq_source_restore.conf以及多表读取multiTableIT/rocketmq_multi_source_to_assert.conf场景配置均位于 connector-rocketmq-e2e 测试资源目录。常见问题与使用建议消费组无历史位点时从哪开始默认CONSUME_FROM_GROUP_OFFSETS模式下若查询到的消费组位点为空源码会回退为从最早位点开始消费时间戳晚于当前时间作业启动即抛出异常这是 RocketMqSourceConfig 的显式运行时校验属于预期行为Tag 过滤不生效检查是否误用了||等 RocketMQ 表达式语法本连接器只做逗号分隔后的精确匹配ACL 集群acl.enabled true时必须同时给出access.key与secret.key否则提交阶段的选项校验会直接失败位点提交若下游 Sink 无法支持事务性写入或你希望完全依赖 Broker 侧位点做故障恢复请保持commit.on.checkpoint true关闭后 SeaTunnel 不再向 Broker 持久化位点。相关文档RocketMQ 源连接器官方文档Source 公共选项Schema 特性说明连接器 V2 特性说明【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考