跳到主要内容
版本:3.0.0

RocketMQ

RocketMQ 源连接器

支持的 Apache RocketMQ 版本​

  • 4.9.0 或更新版本

支持的引擎​

Spark
Flink
SeaTunnel Zeta

主要特性​

描述​

从 Apache RocketMQ topic 读取消息。连接器既可以用一套 schema 读取一个或多个 topic,也可以通过 tables_configs 读取多张不同结构的表。

源参数​

参数名类型是否必填默认值描述
name.srv.addrString是-RocketMQ NameServer 地址,例如 localhost:9876。
topicsString否-topic 名称,多个 topic 使用逗号分隔,例如 "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_errors。
table_listList否-已废弃,请使用 tables_configs。
tagsString否-tag 名称,多个 tag 使用逗号分隔。只消费 RocketMQ tag 与配置值完全匹配的消息。
acl.enabledBoolean否false是否启用 RocketMQ ACL 鉴权。
access.keyString否-访问密钥。acl.enabled = true 时必填。
secret.keyString否-秘密密钥。acl.enabled = true 时必填。
batch.sizeint否100每次最多拉取的消息数。
consumer.groupString否SeaTunnel-Consumer-GroupRocketMQ 消费者组 ID。
commit.on.checkpointBoolean否true是否在 SeaTunnel checkpoint 完成后提交消费位点。
schemaconfig否-消息结构。详情请参考 Schema 特性。不配置时,连接器按文本读取消息体。
formatString否json消息格式。支持 json 和 text。
field.delimiterString否,format = text 时使用的字段分隔符。
start.modeString否CONSUME_FROM_GROUP_OFFSETS启动消费位置。支持:CONSUME_FROM_LAST_OFFSET、CONSUME_FROM_FIRST_OFFSET、CONSUME_FROM_GROUP_OFFSETS、CONSUME_FROM_TIMESTAMP、CONSUME_FROM_SPECIFIC_OFFSETS。
start.mode.offsetsMap否-start.mode = CONSUME_FROM_SPECIFIC_OFFSETS 时必填。key 格式为 topic-queueId,例如 test_topic-0。
start.mode.timestampLong否-start.mode = CONSUME_FROM_TIMESTAMP 时必填,单位是毫秒时间戳。
partition.discovery.interval.millislong否-1动态发现 topic 和分区的间隔,单位毫秒。-1 表示不启用动态发现。
ignore_parse_errorsBoolean否false是否跳过解析失败的 JSON 消息。
consumer.poll.timeout.millislong否5000拉取消息的超时时间,单位毫秒。
common-optionsconfig否-源连接器通用参数,详情请参考 源通用参数。

参数说明​

启动消费位置​

start.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 指定的位点开始读。

当 start.mode = CONSUME_FROM_TIMESTAMP 时,start.mode.timestamp 必须是非负的毫秒时间戳,并且不能晚于任务运行时的当前时间。

start.mode = "CONSUME_FROM_SPECIFIC_OFFSETS"
start.mode.offsets = {
test_topic-0 = 50
}
start.mode = "CONSUME_FROM_TIMESTAMP"
start.mode.timestamp = 1667179890315

消息格式​

当 format = json 时,请配置 schema,SeaTunnel 会按 schema 将 JSON 消息体解析成有类型的字段。配置 ignore_parse_errors = true 后,遇到无法解析的 JSON 消息会跳过,而不是让任务失败。

当 format = text 时,SeaTunnel 会按 field.delimiter 拆分消息体,并按照 schema 字段顺序映射数据。如果不配置 schema,消息体会作为单个文本值读取。

标签​

tags 使用逗号分隔,例如 tag_a,tag_b。连接器会把拉取到的消息 tag 和这些值逐个做精确匹配,因此这里不要使用 RocketMQ 的 tag 表达式语法,例如 tag_a || tag_b。在多表读取任务中,每个 tables_configs 条目都可以配置自己的 tags 过滤条件。

多表读取​

当不同 topic 的字段结构不一样时,使用 tables_configs。每一项都必须包含 topics,并且可以单独配置 schema、format、tags 和启动消费位置。如果没有配置 schema.table,输出表名默认使用 topic 名称。

topics、tables_configs 和已废弃的 table_list 互斥,只能配置其中一个。在 tables_configs 中,单个条目未配置的参数会沿用顶层默认值,因此每个条目只需要覆盖该 topic 特有的 schema、tag 或启动位置。

如果某个 tables_configs 条目使用 start.mode = CONSUME_FROM_TIMESTAMP,必须同时配置 start.mode.timestamp。如果使用 start.mode = CONSUME_FROM_SPECIFIC_OFFSETS,必须同时配置非空的 start.mode.offsets。

任务示例​

读取 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 读取文本消息​

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"
}
}

从指定 offset 读取​

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"
}
}

读取多个不同结构的 topic​

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 {}
}

从指定时间戳读取​

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"
}
}

变更日志​

Change Log
ChangeCommitVersion
[Fix][Zeta] Reuse fixed slots after master failover (#11458)https://github.com/apache/seatunnel/commit/a9cda80f43.0.0
[Improve][Common] Add HashUtils.bucketIndex for hash-to-bucket routing (#11987)https://github.com/apache/seatunnel/commit/be53a1d3d3.0.0
[Improve][Connector-V2] Migrate RocketMQ Source validation to declarative OptionRule (#11158)https://github.com/apache/seatunnel/commit/c8fb493583.0.0
Test/rocketmq restore e2e (#10778)https://github.com/apache/seatunnel/commit/6d2a104653.0.0
[Improve][Connector-V2] Complete OptionRule declarations for RocketMQ source and sink (#10701)https://github.com/apache/seatunnel/commit/7becf69b73.0.0
[Feature][Connector-V2] Support multi-table read for RocketMQ source (#10619)https://github.com/apache/seatunnel/commit/85615468b3.0.0