---
title: Kafka
---

# Kafka

本文介绍如何将 FlowMQ 数据写入 Kafka（含 FlowMQ 数据流或其他 Kafka 兼容集群）。控制台 Sink 选项名为 `Kafka Franz`。

## 前提条件

- 已了解管道结构与创建流程，参考 [数据集成](data-integration.md)；Sink 连接器总览见 [Sink](data-integration-output.md)
- Kafka / FlowMQ 服务可访问
- 已规划目标 Topic，且与本管道 Source 读取的 Topic 区分开，避免消息环路

## 配置步骤

1. 在“数据管道”中创建管道，完成基础信息配置。
2. 在 Source 步骤选择 `Kafka`（Topic：`flowmq.mqtt.kafka`，消费组：`flowmq-sink-kafka`），Processors 按需配置。
3. 在 Sink 步骤选择 `Kafka Franz`，配置 `seed_brokers`、`topic` 等参数。
4. 点击“测试”验证连接。
5. 继续完成确认步骤。

## 配置示例

以下示例将 Kafka Topic `flowmq.mqtt.kafka` 中的样例消息写入目标系统。样例 JSON 与 Source 约定见 [Sink · 示例场景](data-integration-output.md#示例场景)。

### 目标 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](data-integration-processors-expressions.md)。

| 参数 | 支持表达式 | 说明 |
|------|------------|------|
| `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` 时使用，表达式结果必须是有效整数
