流计算能力内置了多种连接器(Connector),用于在 流计算SQL 中对接上下游数据系统。本文档面向使用流计算 SQL 开发的客户,详细说明当前支持的连接器的 DDL 建表语法及 WITH 参数。
每个连接器均支持作为 Source(数据源,读取消息) 与 Sink(数据汇,写入消息) 使用。您可以在流计算控制台的"连接器"页面点击"作为 Source 创建"或"作为 Sink 创建"快速生成建表模板,再参照本文档补全参数。
连接器范围
连接器 | 源 | 目标 |
Kafka | 支持 | 支持 |
MQTT | 支持 | 支持 |
RocketMQ | 支持 | 支持 |
说明:连接器会逐渐开放,如有紧急需求,请通过 工单 提出。
通用说明
建表语法结构
所有连接器统一采用标准 CREATE TABLE 语法定义:
CREATE TABLE table_name (
-- 字段定义
col1 STRING,
col2 BIGINT,
...
) WITH (
'connector' = '<连接器类型>',
-- 连接器专属参数
...
);参数标记约定
下文参数表格中,"是否必填"列标记为 是 的参数在建表时必须提供;标记为 否 的参数可省略,省略时采用"默认值"列所示的取值。
数据格式(format)
Source 与 Sink 表均需通过 format 或 value.format 指定消息体的序列化格式。常用取值:json、csv、avro、raw、debezium-json、canal-json。不同 format 有各自的附加参数,详见文末「附录:Format 常用参数」。
Kafka 连接器
Kafka 连接器用于对接云消息队列 Kafka 版的 Topic,支持读取和写入消息,并可利用消息的时间戳、分区、offset 等元数据。
作为 Source
CREATE TABLE kafka_source (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING,
-- 元数据字段(可选)
`ts` TIMESTAMP(3) METADATA FROM 'timestamp',
`partition_id` BIGINT METADATA FROM 'partition' VIRTUAL,
`offset` BIGINT METADATA FROM 'offset' VIRTUAL
) WITH (
'connector' = 'kafka',
'topic' = 'your_topic',
'properties.bootstrap.servers' = 'alikafka-xxx:9092',
'properties.group.id' = 'your_consumer_group',
'scan.startup.mode' = 'group-offsets',
'format' = 'json'
);Source 参数说明
参数 | 是否必填 | 默认值 | 说明 |
| 是 | 无 | 固定填 |
| 是 | 无 | 读取的 Topic 名称。多个 Topic 用分号分隔,如 |
| 否 | 无 | 以正则匹配一批 Topic,与 |
| 是 | 无 | Kafka 接入点地址,多个用逗号分隔。可在实例详情页获取。 |
| 是 | 无 | 消费组 ID,用于记录消费位点。 |
| 否 |
| 启动消费位点: |
| 否 | 无 | 当模式为 |
| 否 | 无 | 当模式为 |
| 是 | 无 | 消息体解析格式,如 |
| 否 | 无 | 消息 Key 的解析格式,需配合 |
| 否 | 无 | 指定哪些字段来自消息 Key。 |
| 否 | 无 | 透传给 Kafka 客户端的原生参数,如 |
| 否 | 无 | 动态分区发现间隔,如 |
作为 Sink
CREATE TABLE kafka_sink (
`user_id` BIGINT,
`item_id` BIGINT,
`behavior` STRING
) WITH (
'connector' = 'kafka',
'topic' = 'your_sink_topic',
'properties.bootstrap.servers' = 'alikafka-xxx:9092',
'format' = 'json',
'sink.partitioner' = 'default'
);Sink 参数说明
参数 | 是否必填 | 默认值 | 说明 |
| 是 | 无 | 固定填 |
| 是 | 无 | 写入的目标 Topic。 |
| 是 | 无 | Kafka 接入点地址。 |
| 是 | 无 | 消息体序列化格式。 |
| 否 | 无 | Key 的序列化格式,配合 |
| 否 | 无 | 指定作为消息 Key 写入的字段。 |
| 否 |
| 分区策略: |
| 否 |
| 投递语义: |
| 否 | 无 |
|
| 否 | 无 | 透传给 Kafka Producer 的原生参数。 |
MQTT 连接器
在流计算中,MQTT 消息通过 RocketMQ 连接器(connector = 'rocketmq')读写,并使用一组 mqtt 专属子参数,而非独立的 mqtt 连接器。适用于物联网(IoT)设备上下行消息场景。
作为 Source
CREATE TABLE mqtt_source (
`user_id` STRING,
`user_name` STRING,
`action` STRING,
`event_time` TIMESTAMP(3),
`mqtt_client_id` STRING METADATA VIRTUAL,
`mqtt_qos_level` STRING METADATA VIRTUAL,
`mqtt_real_topic` STRING METADATA VIRTUAL
) WITH (
'connector' = 'rocketmq',
'rocketmq.source.mqtt.first.topic' = 'flink-source-1',
'rocketmq.client.mqtt.instance.id' = 'mqtt-cn-xxxxxxxxx',
'rocketmq.source.mqtt.group.id' = 'GID_test',
'rocketmq.source.group' = 'GID-flink',
'rocketmq.client.endpoints' = 'mqtt-cn-xxx-single-internal-vpc.mqtt.aliyuncs.com:5673',
'rocketmq.client.accessKey' = '${secret_ak}',
'rocketmq.client.secretKey' = '${secret_sk}',
'rocketmq.source.filter.tag' = '*',
'rocketmq.source.startup.scan.mode' = 'earliest',
'rocketmq.source.pull.batch.size' = '256',
'rocketmq.source.pull.interval.ms' = '10',
'rocketmq.source.pull.suspend.timeout' = '30000',
'rocketmq.client.message.field.delimiter' = ';'
);作为 Sink
CREATE TABLE mqtt_sink (
`user_id` STRING,
`user_name` STRING,
`action` STRING,
`keys` STRING METADATA,
`tags` STRING METADATA
) WITH (
'connector' = 'rocketmq',
'rocketmq.sink.mqtt.topic' = 'flink-sink-1/C',
'rocketmq.client.mqtt.instance.id' = 'mqtt-cn-xxxxxxxxx',
'rocketmq.sink.group' = 'PID-flink',
'rocketmq.client.endpoints' = 'mqtt-cn-xxx-single-internal-vpc.mqtt.aliyuncs.com:5673',
'rocketmq.client.accessKey' = '${secret_ak}',
'rocketmq.client.secretKey' = '${secret_sk}',
'rocketmq.sink.delivery.guarantee' = 'AT_LEAST_ONCE',
'rocketmq.client.message.field.delimiter' = ';'
);MQTT 专属参数
参数 | 适用端 | 是否必填 | 默认值 | 说明 |
rocketmq.client.mqtt.instance.id | Source/Sink | 是 | 无 | MQTT 实例 ID,如 mqtt-cn-xxxxxxxxx。 |
rocketmq.source.mqtt.first.topic | Source | 是 | 无 | 订阅的 MQTT 一级 Topic(父 Topic),对应 LMQ 映射的 RocketMQ Topic。 |
rocketmq.source.mqtt.group.id | Source | 是 | 无 | MQTT 消费组 ID。 |
rocketmq.sink.mqtt.topic | Sink | 是 | 无 | 发送的 MQTT 目标 Topic,可含二级 Topic 与 QoS 后缀,如 flink-sink-1/C(/C 表示 CleanSession 等发布语义)。 |
MQTT 可用元数据列
元数据 | 类型 | 读/写 | 说明 |
mqtt_client_id | STRING | 读 | 发送该消息的 MQTT 客户端 ID。 |
mqtt_qos_level | STRING | 读 | 消息 QoS 等级(0/1/2)。 |
mqtt_real_topic | STRING | 读 | 消息实际的完整 MQTT Topic(含二级 Topic)。 |
RocketMQ 通用元数据列(topic、msg_id、store_timestamp、born_timestamp、queue_id、queue_offset、keys、tags)在 MQTT 场景同样可用。
RocketMQ 连接器
RocketMQ 连接器用于对接云消息队列 RocketMQ 版,支持普通消息的读写,并可利用 Tag、Key、Queue 等属性进行过滤与路由。所有参数均以 rocketmq. 为前缀,分为三类:rocketmq.client.*(连接与消息编解码,Source/Sink 通用)、rocketmq.source.*(消费端专属)、rocketmq.sink.*(生产端专属)。
作为 Source
CREATE TABLE rocketmq_source (
`user_id` STRING,
`user_name` STRING,
`action` STRING,
`event_time` TIMESTAMP(3),
`topic` STRING METADATA VIRTUAL,
`msg_id` STRING METADATA VIRTUAL,
`store_timestamp` BIGINT METADATA VIRTUAL,
`born_timestamp` BIGINT METADATA VIRTUAL,
`queue_id` INT METADATA VIRTUAL,
`queue_offset` BIGINT METADATA VIRTUAL,
`keys` STRING METADATA VIRTUAL,
`tags` STRING METADATA VIRTUAL
) WITH (
'connector' = 'rocketmq',
'rocketmq.source.topic' = 'source',
'rocketmq.source.group' = 'GID-flink',
'rocketmq.client.endpoints' = 'rmq-cn-xxx-vpc.cn-hangzhou.rmq.aliyuncs.com:8080',
-- 4.x 独享实例需带 namespace(实例 ID):
-- 'rocketmq.client.namespace' = 'MQ_INST_xxxxxxxx',
'rocketmq.client.accessKey' = '${secret_ak}',
'rocketmq.client.secretKey' = '${secret_sk}',
'rocketmq.source.filter.tag' = '*',
'rocketmq.source.startup.scan.mode' = 'latest',
'rocketmq.source.pull.batch.size' = '32',
'rocketmq.source.pull.interval.ms' = '100',
'rocketmq.source.pull.suspend.timeout' = '30000',
'rocketmq.client.message.field.delimiter' = ';',
'rocketmq.client.message.line.delimiter' = '\n',
'rocketmq.client.message.encoding' = 'UTF-8',
'rocketmq.client.message.length.check' = 'NONE',
'rocketmq.client.timeZone' = 'Asia/Shanghai',
'rocketmq.client.partition.discovery.interval.ms' = '10000'
);连接与鉴权参数(rocketmq.client.*,Source/Sink 通用)
参数 | 是否必填 | 默认值 | 说明 |
rocketmq.client.endpoints | 是 | 无 | 接入点地址(含端口)。RocketMQ 5.x 形如 rmq-cn-xxx-vpc.cn-hangzhou.rmq.aliyuncs.com:8080;4.x 独享实例形如 http://MQ_INST_xxx.cn-hangzhou.mq-internal.aliyuncs.com:8080。 |
rocketmq.client.namespace | 否 | 无 | 实例命名空间 / 实例 ID(如 MQ_INST_xxx)。使用 4.x 独享实例、Topic/Group 需按实例隔离时必填;5.x 独立域名接入点通常无需填写。 |
rocketmq.client.accessKey | 否 | 无 | 阿里云 AccessKey ID,开启 ACL 鉴权的实例必填。建议用密钥管理引用。 |
rocketmq.client.secretKey | 否 | 无 | 阿里云 AccessKey Secret,开启 ACL 鉴权的实例必填。严禁明文,建议密钥管理引用。 |
rocketmq.client.message.field.delimiter | 否 | , | 字段分隔符:将消息体按位置拆分为各列(Source)/ 将各列拼接为消息体(Sink)。样例用 ;。 |
rocketmq.client.message.line.delimiter | 否 | \n | 行分隔符:一条消息含多行记录时用于切分。 |
rocketmq.client.message.encoding | 否 | UTF-8 | 消息体字符编码。 |
rocketmq.client.message.length.check | 否 | NONE | 列数长度校验:NONE 不校验;开启后字段数与表结构不符会报错/丢弃。 |
rocketmq.client.timeZone | 否 | 系统时区 | 时间类型字段解析所用时区,如 Asia/Shanghai。 |
rocketmq.client.partition.discovery.interval.ms | 否 | 无(关闭) | 队列(分区)动态发现间隔(毫秒),用于自动感知 Topic 扩缩容,如 10000。 |
消费端参数(rocketmq.source.*)
参数 | 是否必填 | 默认值 | 说明 |
rocketmq.source.topic | 是 | 无 | 读取的 Topic 名称。 |
rocketmq.source.group | 是 | 无 | 消费组(Group ID),需在控制台预先创建。 |
rocketmq.source.filter.tag | 否 | * | Tag 过滤表达式,* 表示不过滤;多个 Tag 用 ` |
rocketmq.source.startup.scan.mode | 否 | latest | 启动消费位点:earliest(最早)/ latest(最新)/ timestamp(指定时间)/ group_offsets(消费组已提交位点)。 |
rocketmq.source.pull.batch.size | 否 | 32 | 每次拉取的最大消息条数。压测大流量时可调大(样例 MQTT 场景用 256)。 |
rocketmq.source.pull.interval.ms | 否 | 0 | 相邻两次拉取的间隔(毫秒),控制拉取频率与背压。 |
rocketmq.source.pull.suspend.timeout | 否 | 30000 | 长轮询在服务端的挂起超时(毫秒):无消息时请求最长挂起时长。 |
作为 Sink
CREATE TABLE rocketmq_sink (
`user_id` STRING,
`user_name` STRING,
`action` STRING,
`process_time` TIMESTAMP(3),
`field_a` STRING,
`field_b` STRING,
`keys` STRING METADATA,
`tags` STRING METADATA
) WITH (
'connector' = 'rocketmq',
'rocketmq.sink.topic' = 'sink',
'rocketmq.sink.group' = 'PID-flink',
'rocketmq.client.endpoints' = 'rmq-cn-xxx-vpc.cn-hangzhou.rmq.aliyuncs.com:8080',
'rocketmq.client.accessKey' = '${secret_ak}',
'rocketmq.client.secretKey' = '${secret_sk}',
'rocketmq.sink.delivery.guarantee' = 'AT_LEAST_ONCE',
'rocketmq.client.message.field.delimiter' = ';',
'rocketmq.sink.key.columns' = 'field_a,user_name',
'rocketmq.sink.key.columns.included' = 'true',
'rocketmq.sink.tag.column' = 'field_b',
'rocketmq.sink.tag.column.included' = 'false'
);生产端参数(rocketmq.sink.*)
参数 | 是否必填 | 默认值 | 说明 |
rocketmq.sink.topic | 是 | 无 | 写入的目标 Topic。 |
rocketmq.sink.group | 是 | 无 | 生产组(Producer Group ID)。 |
rocketmq.sink.delivery.guarantee | 否 | AT_LEAST_ONCE | 投递语义:NONE / AT_LEAST_ONCE / EXACTLY_ONCE(需实例与连接器版本支持事务)。 |
rocketmq.sink.key.columns | 否 | 无 | 指定作为消息 Key 的字段,多个用逗号分隔,如 field_a,user_name。用于消息查询与幂等去重。 |
rocketmq.sink.key.columns.included | 否 | true | 作为 Key 的字段是否同时写入消息体。true 保留在 body,false 仅作 Key 不进 body。 |
rocketmq.sink.tag.column | 否 | 无 | 指定作为消息 Tag 的字段,供下游按 Tag 过滤。 |
rocketmq.sink.tag.column.included | 否 | true | 作为 Tag 的字段是否同时写入消息体。false 表示仅作 Tag 不进 body。 |
说明:keys / tags 也可通过可写元数据列(METADATA,不带 VIRTUAL)在 SELECT 中直接赋值,与 rocketmq.sink.key.columns / rocketmq.sink.tag.column 是两种等效路径,按需二选一。
RocketMQ 可用元数据列
元数据 | 类型 | 读/写 | 说明 |
topic | STRING | 读 | 消息所属 Topic。 |
msg_id | STRING | 读 | 消息全局唯一 ID。 |
store_timestamp | BIGINT | 读 | 服务端存储时间戳(毫秒)。 |
born_timestamp | BIGINT | 读 | 消息产生时间戳(毫秒)。 |
queue_id | INT | 读 | 队列 ID。 |
queue_offset | BIGINT | 读 | 队列内 offset。 |
keys | STRING | 读/写 | 消息 Key。 |
tags | STRING | 读/写 | 消息 Tag。 |
Format 常用参数
不同 format 需配合以下附加参数(在 WITH 中以 format名.参数 形式书写,如 json.ignore-parse-errors):
JSON
参数 | 默认值 | 说明 |
|
| 解析出错时是否跳过该条并置 null,避免作业失败。 |
|
| 缺字段时是否报错。 |
|
| 时间格式标准: |
CSV
参数 | 默认值 | 说明 |
|
| 字段分隔符。 |
|
| 解析出错时是否跳过。 |
| 无 | 表示 null 的字面量。 |
其他
raw(原始字节,适合单字段透传)、avro(需配合 Schema)、debezium-json / canal-json(用于订阅 CDC 变更数据)。
说明: 以上参数以流计算控制台当前支持的连接器版本为准。实际接入点地址、鉴权信息请在对应实例详情页获取;涉及密码、AccessKey 等敏感信息建议通过密钥管理引用,避免在 SQL 中明文书写。如需 exactly-once 等高级语义,请先确认目标实例与连接器版本是否支持。