Kafka
本文介绍如何将 FlowMQ 数据写入 Kafka(含 FlowMQ 数据流或其他 Kafka 兼容集群)。控制台 Sink 选项名为 Kafka Franz。
前提条件
- 已了解管道结构与创建流程,参考 数据集成;Sink 连接器总览见 Sink
- Kafka / FlowMQ 服务可访问
- 已规划目标 Topic,且与本管道 Source 读取的 Topic 区分开,避免消息环路
配置步骤
- 在“数据管道”中创建管道,完成基础信息配置。
- 在 Source 步骤选择
Kafka(Topic:flowmq.mqtt.kafka,消费组:flowmq-sink-kafka),Processors 按需配置。 - 在 Sink 步骤选择
Kafka Franz,配置seed_brokers、topic等参数。 - 点击“测试”验证连接。
- 继续完成确认步骤。
配置示例
以下示例将 Kafka Topic flowmq.mqtt.kafka 中的样例消息写入目标系统。样例 JSON 与 Source 约定见 Sink · 示例场景。
目标 Topic 准备
建管道前,先在「数据流」中创建 Source Topic flowmq.mqtt.kafka 和目标 Topic flowmq.sink.demo。Topic 不存在时,连通性测试可能通过,但无法完成实际转发。
写入后可在「数据流」中查看 flowmq.sink.demo。按下文填写 key 时,样例消息的 Kafka key 为 1SDA526VD_POP_P1。
可参考以下参数示例:
| 参数 | 建议值 |
|---|---|
| seed_brokers | 127.0.0.1:9092 |
| topic | flowmq.sink.demo |
| key | ${! json("source") }(可选) |
| client_id | sink-kafka-demo |
| max_in_flight | 10 |
| timeout | 10s |
| compression | none |
必要表单参数
参数名后带 * 表示控制台必填。「支持表达式」为「是」时,可使用 ${! ... } 按消息动态取值,详见 Processors。
| 参数 | 支持表达式 | 说明 |
|---|---|---|
seed_brokers* | 否 | Broker 地址列表 |
topic* | 是 | 写入目标 Topic |
key | 是 | 消息 key;不填则无 key |
partitioner | 否 | 分区策略,默认 murmur2_hash |
partition | 是 | 仅当 partitioner 为 manual 时生效,须为整数 |
client_id* | 否 | 客户端标识 |
max_in_flight* | 否 | 并行发送的批次上限 |
timeout* | 否 | 发送超时时间 |
max_message_bytes* | 否 | 单条消息最大字节数 |
max_buffered_records* | 否 | 客户端缓冲记录上限 |
compression* | 否 | 压缩类型(如 none、gzip、snappy、lz4、zstd) |
sasl | 否 | 启用鉴权时需要配置 |
参数建议
topic宜按业务域命名;也可按消息动态指定,例如${! json("source") }- 写入 FlowMQ 数据流时,
seed_brokers填本集群 Kafka 地址即可 max_in_flight可先沿用默认值,压测后再上调- 需要按设备/测点分区时,填写
key(如${! json("source") }),并保持默认murmur2_hash
注意事项
- 若目标 Topic 的数据会再次流入本管道读取的 Kafka Topic,Sink 请勿写入该 Topic,以免形成消息环路
include_prefixes/include_patterns若出现空行,连通性「测试」按钮可能不出现;可删掉空行后再测partition仅在partitioner为manual时使用,表达式结果必须是有效整数