流计算连接器 SQL 使用说明

更新时间:
复制 MD 格式

流计算能力内置了多种连接器(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 表均需通过 formatvalue.format 指定消息体的序列化格式。常用取值:jsoncsvavrorawdebezium-jsoncanal-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 参数说明

参数

是否必填

默认值

说明

connector

固定填 kafka

topic

读取的 Topic 名称。多个 Topic 用分号分隔,如 topic-1;topic-2

topic-pattern

以正则匹配一批 Topic,与 topic 二选一。

properties.bootstrap.servers

Kafka 接入点地址,多个用逗号分隔。可在实例详情页获取。

properties.group.id

消费组 ID,用于记录消费位点。

scan.startup.mode

group-offsets

启动消费位点:earliest-offset(最早)/latest-offset(最新)/group-offsets(消费组已提交位点)/timestamp(指定时间)/specific-offsets(指定 offset)。

scan.startup.specific-offsets

当模式为 specific-offsets 时指定,格式 partition:0,offset:42;partition:1,offset:300

scan.startup.timestamp-millis

当模式为 timestamp 时指定,起始时间的毫秒时间戳。

format / value.format

消息体解析格式,如 jsoncsvavro

key.format

消息 Key 的解析格式,需配合 key.fields 使用。

key.fields

指定哪些字段来自消息 Key。

properties.*

透传给 Kafka 客户端的原生参数,如 properties.security.protocol

scan.topic-partition-discovery.interval

动态分区发现间隔,如 10s,用于自动感知 Topic 扩容。

作为 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 参数说明

参数

是否必填

默认值

说明

connector

固定填 kafka

topic

写入的目标 Topic。

properties.bootstrap.servers

Kafka 接入点地址。

format / value.format

消息体序列化格式。

key.format

Key 的序列化格式,配合 key.fields

key.fields

指定作为消息 Key 写入的字段。

sink.partitioner

default

分区策略:default(Kafka 默认)/fixed(每个并行度对应固定分区,减少小文件)/round-robin(轮询)/自定义分区器类名。

sink.delivery-guarantee

at-least-once

投递语义:none/at-least-once/exactly-once(需配合事务)。

sink.transactional-id-prefix

exactly-once 语义下的事务 ID 前缀,必须唯一。

properties.*

透传给 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

参数

默认值

说明

json.ignore-parse-errors

false

解析出错时是否跳过该条并置 null,避免作业失败。

json.fail-on-missing-field

false

缺字段时是否报错。

json.timestamp-format.standard

SQL

时间格式标准:SQLISO-8601

CSV

参数

默认值

说明

csv.field-delimiter

,

字段分隔符。

csv.ignore-parse-errors

false

解析出错时是否跳过。

csv.null-literal

表示 null 的字面量。

其他

raw(原始字节,适合单字段透传)、avro(需配合 Schema)、debezium-json / canal-json(用于订阅 CDC 变更数据)。

说明: 以上参数以流计算控制台当前支持的连接器版本为准。实际接入点地址、鉴权信息请在对应实例详情页获取;涉及密码、AccessKey 等敏感信息建议通过密钥管理引用,避免在 SQL 中明文书写。如需 exactly-once 等高级语义,请先确认目标实例与连接器版本是否支持。