帮你快速理解、总结文档立即下载

消息队列 RocketMQ

最近更新时间:2026-09-01 15:12:30
我的收藏

介绍

RocketMQ 连接器用于在 Flink 作业与 Apache RocketMQ 消息队列之间进行数据读写,同时支持 RocketMQ 4.x 和 RocketMQ 5.x 两个服务端版本。用户既可以把 RocketMQ Topic 作为 Flink 作业的数据源(Source),也可以把它作为数据目的(Sink)。
RocketMQ 支持同一个 Topic 多分区读写,数据可以从多个分区读入,也可以写入到多个分区,以提供更高的吞吐量,减少数据倾斜和热点。

版本说明

Flink 版本
RocketMQ 4.x
RocketMQ 5.x
1.16
支持
支持
1.18
支持
支持
1.20
支持
支持

使用范围

RocketMQ 支持用作数据源表(Source),也可以作为 Tuple 数据流的目的表(Sink)。

DDL 定义

用作数据源(Source)

RocketMQ 5.x - 基础示例

CREATE TABLE rocketmq_source (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING
) WITH (
'connector' = 'rocketmq5', -- 固定为 rocketmq5
'rocketmq.client.endpoints' = 'rmq-xxx.rocketmq.gz.qcloud.tencenttdmq.com:8080',
-- 必选,替换为您的 RocketMQ 5.x 接入点
'rocketmq.source.topic' = 'YourTopic', -- 必选,替换为您要消费的 Topic
'rocketmq.source.group' = 'YourConsumerGroup', -- 必选,替换为您的消费者分组
'rocketmq.source.startup.scan.mode' = 'earliest',
'rocketmq.client.accessKey' = 'YourAccessKey', -- 可选,替换为您的 AccessKey
'rocketmq.client.secretKey' = 'YourSecretKey' -- 可选,替换为您的 SecretKey
);

RocketMQ 4.x - 基础示例

CREATE TABLE rocketmq_source (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING
) WITH (
'connector' = 'rocketmq4', -- 固定为 rocketmq4
'rocketmq.client.endpoints' = 'http://rocketmq-xxx.rocketmq.gz.qcloud.tencenttdmq.com:9876',
-- 必选,替换为您的 RocketMQ 4.x 接入点
'rocketmq.source.topic' = 'YourTopic', -- 必选,替换为您要消费的 Topic
'rocketmq.source.group' = 'YourConsumerGroup', -- 必选,替换为您的消费者分组
'rocketmq.client.accessKey' = 'YourAccessKey', -- 可选,替换为您的 AccessKey
'rocketmq.client.secretKey' = 'YourSecretKey' -- 可选,替换为您的 SecretKey
);

用作数据目的(Sink)

RocketMQ 5.x - 基础示例

CREATE TABLE rocketmq_sink (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING
) WITH (
'connector' = 'rocketmq5', -- 固定为 rocketmq5
'rocketmq.sink.topic' = 'YourTopic', -- 必选,替换为您要写入的 Topic
'rocketmq.client.endpoints' = 'rmq-xxx.rocketmq.gz.qcloud.tencenttdmq.com:8080',
-- 必选,替换为您的 RocketMQ 5.x 接入点
'rocketmq.client.accessKey' = 'YourAccessKey', -- 可选,替换为您的 AccessKey
'rocketmq.client.secretKey' = 'YourSecretKey' -- 可选,替换为您的 SecretKey
);

RocketMQ 4.x - 基础示例

CREATE TABLE rocketmq_sink (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING
) WITH (
'connector' = 'rocketmq4', -- 固定为 rocketmq4
'rocketmq.sink.topic' = 'YourTopic', -- 必选,替换为您要写入的 Topic
'rocketmq.client.endpoints' = 'http://rocketmq-xxx.rocketmq.gz.qcloud.tencenttdmq.com:9876',
-- 必选,替换为您的 RocketMQ 4.x 接入点
'rocketmq.sink.group' = 'YourProducerGroup', -- 必选,替换为您的生产者分组
'rocketmq.client.accessKey' = 'YourAccessKey', -- 可选,替换为您的 AccessKey
'rocketmq.client.secretKey' = 'YourSecretKey' -- 可选,替换为您的 SecretKey
);

WITH 参数

通用客户端参数(rocketmq.client.*

参数
必填
默认值
描述
rocketmq.client.endpoints
RocketMQ 服务地址。RocketMQ 4.x 填 NameServer 地址(如 127.0.0.1:9876),RocketMQ 5.x 填 Proxy 接入点
rocketmq.client.namespace
RocketMQ 实例命名空间(多租户场景使用)
rocketmq.client.channel
CLOUD
RocketMQ 接入通道。可选 LOCAL / CLOUD / ALIYUN
rocketmq.client.partition.discovery.interval.ms
10000
从 NameServer / Proxy 拉取路由信息的周期(毫秒)
rocketmq.client.timeZone
解析时间字符串所用的时区(如 Asia/Shanghai)
rocketmq.client.message.encoding
UTF-8
消息体的字符编码
rocketmq.client.message.field.delimiter
\\u0001
字段分隔符(默认 ASCII 0x01)
rocketmq.client.message.line.delimiter
\\n
消息内行分隔符
rocketmq.client.message.length.check
NONE
字段数不匹配时的处理策略,可选 NONE / SKIP / EXCEPTION / PAD
rocketmq.client.accessKey
ACL 访问凭据(AccessKey)
rocketmq.client.secretKey
ACL 访问凭据(SecretKey)

源表(Source)参数(rocketmq.source.*

参数
必填
默认值
描述
rocketmq.source.topic
订阅的 Topic 名称,可使用 ; 分隔多个 Topic
rocketmq.source.group
消费者分组名称(用于 Group 维度的消费位点管理)
rocketmq.source.filter.tag
*
消息 Tag 过滤条件,RocketMQ 仅支持单 Tag 过滤
rocketmq.source.filter.sql
消息过滤表达式(与 Tag 二选一)
rocketmq.source.startup.scan.mode
latest
启动消费位点模式:earliest / latest / timestamp / specific_offset / group_offsets
rocketmq.source.startup.offset.specific
-1
启动消费位点的具体偏移量。设置此参数即启用 specific_offset 模式,不要同时设置 startup.scan.mode
rocketmq.source.startup.offset.timestamp
-1
启动消费位点对应的时间戳(毫秒)。设置此参数即启用 timestamp 模式,不要同时设置 startup.scan.mode
rocketmq.source.startup.offset.date
启动消费位点对应的日期字符串(yyyy-MM-dd HH:mm:ss),需配合 rocketmq.client.timeZone。设置此参数即启用 timestamp 模式,不要同时设置 startup.scan.mode
rocketmq.source.stop.offset.timestamp
有限流场景下停止消费的时间戳;提供时 Source 进入有界模式
rocketmq.source.allocate.strategy
hash
队列分配策略,可选 hash/ broadcast/ average
rocketmq.source.pull.batch.size
32
每次拉取消息的最大条数
rocketmq.source.offset.commit.checkpoint
true
是否在 Flink checkpoint 时 commit 位点
说明:
启动模式互斥说明:startup.scan.mode 与 startup.offset.specific / startup.offset.timestamp / startup.offset.date 互斥;同时设置会导致作业启动失败。

结果表(Sink)参数(rocketmq.sink.*

参数
必填
默认值
描述
rocketmq.sink.topic
写入的 Topic 名称
rocketmq.sink.group
PID-flink-producer
生产者分组
rocketmq.sink.delivery.guarantee
AT_LEAST_ONCE
暂时不支持设置
rocketmq.sink.send.retry.times
3
同步发送的重试次数
rocketmq.sink.send.timeout
5000
单次 send 调用超时(毫秒)
rocketmq.sink.tag
消息 Tag

元数据列

源表可读元数据

源表支持的元数据列(在 DDL 中需声明为 METADATA VIRTUAL):
列名
数据类型
描述
topic
STRING NOT NULL
消息所属 Topic
ingestion-time
TIMESTAMP_LTZ(3) NOT NULL
消息进入 Flink 引擎的时间
event-time
TIMESTAMP_LTZ(3) NOT NULL
消息生产时间(Born Timestamp)
queue-id
INT NOT NULL
消息所在队列 ID
queue-offset
BIGINT NOT NULL
消息的消费位点
msg-id
STRING NOT NULL
消息唯一 ID
keys
STRING NOT NULL
消息的 Keys
tags
STRING NOT NULL
消息的 Tags
properties
MAP<STRING, STRING> NOT NULL
消息自定义属性

结果表可写元数据

结果表支持的元数据列:
列名
数据类型
描述
keys
STRING
消息 Keys,写入 Message 的 Keys 字段
tags
STRING
消息 Tags,写入 Message 的 Tags 字段

代码示例

SQL:从 RocketMQ 读 + 写入 RocketMQ

-- Source:定义 RocketMQ 源表
CREATE TABLE rocketmq_source (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING
) WITH (
'connector' = 'rocketmq5', -- 固定为 rocketmq5
'rocketmq.client.endpoints' = 'rmq-xxx.rocketmq.gz.qcloud.tencenttdmq.com:8080',
-- 必选,替换为您的 RocketMQ 5.x 接入点
'rocketmq.source.topic' = 'YourTopic', -- 必选,替换为您要消费的 Topic
'rocketmq.source.group' = 'YourConsumerGroup', -- 必选,替换为您的消费者分组
'rocketmq.source.startup.scan.mode' = 'earliest', -- 可选,启动消费位点模式:earliest / latest
'rocketmq.client.accessKey' = 'YourAccessKey', -- 可选,替换为您的 AccessKey
'rocketmq.client.secretKey' = 'YourSecretKey' -- 可选,替换为您的 SecretKey
);

-- Sink:定义 RocketMQ 结果表
CREATE TABLE rocketmq_sink (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING
) WITH (
'connector' = 'rocketmq5', -- 固定为 rocketmq5
'rocketmq.sink.topic' = 'YourTopic', -- 必选,替换为您要写入的 Topic
'rocketmq.client.endpoints' = 'rmq-xxx.rocketmq.gz.qcloud.tencenttdmq.com:8080',
-- 必选,替换为您的 RocketMQ 5.x 接入点
'rocketmq.client.accessKey' = 'YourAccessKey', -- 可选,替换为您的 AccessKey
'rocketmq.client.secretKey' = 'YourSecretKey' -- 可选,替换为您的 SecretKey
);

INSERT INTO rocketmq_sink SELECT * FROM rocketmq_source;

二进制格式示例

当 RocketMQ 消息体是 Protobuf / Thrift / Avro / 自定义二进制协议等无法按字符切分的内容时,DDL 中只声明一个 VARBINARY 字段即可。连接器会把整条消息的字节数组原样写入该字段,不做切分;后续可在 Flink 作业里用 UDF 自行解析。
CREATE TABLE rocketmq_source (
msg_body VARBINARY
) WITH (
'connector' = 'rocketmq5', -- 必选,固定为 rocketmq5
'rocketmq.client.endpoints' = 'YourEndpoint', -- 必选,替换为您的 RocketMQ 5.x 接入点
'rocketmq.source.topic' = 'YourTopic', -- 必选,替换为您要消费的 Topic
'rocketmq.source.group' = 'YourConsumerGroup', -- 必选,替换为您的消费者分组
'rocketmq.client.accessKey' = 'YourAccessKey', -- 可选,替换为您的 AccessKey
'rocketmq.client.secretKey' = 'YourSecretKey' -- 可选,替换为您的 SecretKey
);
二进制模式下需要注意:
整张表只能声明一个数据列且类型必须为 VARBINARY;
字段分隔符(rocketmq.client.message.field.delimiter)和字符编码(rocketmq.client.message.encoding)在二进制模式下不生效。
若需要在二进制数据中同时保留消息属性(如 keys / tags / properties),可以在 DDL 中追加元数据列(参考「元数据列」章节),二者在解析时互不干扰。

类型映射

连接器对常见 Flink 类型与消息体字段的映射关系:
Flink 字段类型
RocketMQ 字段类型
BOOLEAN
STRING
VARBINARY
VARCHAR
TINYINT / SMALLINT / INT / BIGINT
INTEGER / BIGINT
FLOAT / DOUBLE
DOUBLE
DECIMAL
DECIMAL