介绍
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 |