跳到主要内容
版本:3.0.0

RocketMQ

RocketMQ Sink 连接器

支持的 Apache RocketMQ 版本​

  • 4.9.0 或更新版本

支持的引擎​

Spark
Flink
SeaTunnel Zeta

主要特性​

描述​

将 SeaTunnel 数据行写入 Apache RocketMQ topic。该 Sink 支持 JSON 和文本消息体、消息 tag、同步发送、按字段选择分区,以及在 exactly.once = true 时使用事务消息保证精确一次写入。

Sink 参数​

参数名类型是否必填默认值描述
topicString是-RocketMQ topic 名称。
name.srv.addrString是-RocketMQ NameServer 地址,例如 localhost:9876。
acl.enabledBoolean否false是否启用 RocketMQ ACL 鉴权。
access.keyString否-访问密钥。acl.enabled = true 时必填。
secret.keyString否-秘密密钥。acl.enabled = true 时必填。
producer.groupString否SeaTunnel-Producer-GroupRocketMQ 生产者组 ID。
tagString否-写入每条消息时使用的 RocketMQ tag。
partition.key.fieldsList否-会被序列化为 RocketMQ 消息 key 的字段名。配置的字段必须存在于上游 schema 中。
formatString否json消息格式。支持 json 和 text。
field.delimiterString否,format = text 时使用的字段分隔符。
producer.send.syncBoolean否false是否同步发送消息。为 false 时异步发送。
exactly.onceBoolean否false是否使用事务消息实现精确一次写入。
max.message.sizeint否4194304最大消息体大小,单位字节。
send.message.timeoutint否3000发送消息超时时间,单位毫秒。
common-optionsconfig否-Sink 连接器通用参数,详情请参考 Sink 通用参数。

参数说明​

partition.key.fields​

partition.key.fields 控制 RocketMQ 消息 key。SeaTunnel 会把这些字段的值序列化成 JSON,并写入 Message.keys。在非事务发送时,同一个 key 还会交给 RocketMQ 的哈希队列选择器,因此相同 key 的数据会进入同一个队列。如果不配置该参数,则由 RocketMQ 自行选择队列。

例如,上游字段中有 c_int 时,可以这样把 c_int 用作消息 key:

partition.key.fields = ["c_int"]

exactly.once​

Sink 支持通过 RocketMQ 事务消息实现精确一次写入。该能力默认关闭。确认 RocketMQ 集群和作业 checkpoint 配置满足事务写入要求后,可设置 exactly.once = true。

当 format = text 时,SeaTunnel 会按上游 schema 的字段顺序序列化,并使用 field.delimiter 拼接字段。当 format = json 时,每行数据会写成一个 JSON 对象。

producer.send.sync​

producer.send.sync = true 表示生产者会等待 RocketMQ 确认每次发送请求。默认值为 false 时,消息会异步发送。

任务示例​

写入 JSON 消息​

env {
parallelism = 1
job.mode = "BATCH"
}

source {
FakeSource {
row.num = 10
schema = {
fields {
c_string = string
c_int = int
c_timestamp = timestamp
}
}
}
}

sink {
Rocketmq {
name.srv.addr = "rocketmq-e2e:9876"
topic = "test_topic"
partition.key.fields = ["c_int"]
producer.send.sync = true
}
}

写入文本消息​

env {
parallelism = 1
job.mode = "BATCH"
}

source {
FakeSource {
row.num = 10
schema = {
fields {
id = bigint
content = string
}
}
}
}

sink {
Rocketmq {
name.srv.addr = "rocketmq-e2e:9876"
topic = "test_text_topic"
format = text
field.delimiter = ","
producer.send.sync = true
}
}

写入带 tag 的消息​

env {
parallelism = 1
job.mode = "BATCH"
}

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

RocketMQ 读写 RocketMQ​

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

变更日志​

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