---
title: Kafka Source 集成
---

# Kafka Source 集成

本文介绍如何在管道的 Source 步骤中使用 Kafka 作为数据来源。

## 适用场景

- 从 Kafka Topic 持续消费业务事件
- 将存量 Kafka 数据接入 FlowMQ 处理链路

## 前提条件

- 已了解管道结构与创建流程，参考 [数据集成](data-integration.md)；Source 连接器总览见 [Source](data-integration-input.md)
- Kafka 服务可访问
- 已准备可消费的 Topic 与消费组策略

## 配置步骤

1. 在“数据管道”中创建管道，完成基础信息配置。
2. 在 Source 步骤选择 `Kafka`。
3. 配置 `seed_brokers`、`topics`、`consumer_group` 等参数。
4. 点击“测试”验证连接。
5. 继续完成 Processors、Sink 与确认步骤。

## 配置示例

可参考以下参数示例：

| 参数 | 建议值 |
|---|---|
| seed_brokers | `127.0.0.1:9092` |
| topics | `flowmq.mqtt.kafka` |
| consumer_group | `flowmq-source-mqtt-kafka` |

## 必要表单参数

以表单中带 `*` 的字段为必填：

| 参数 | 说明 |
|------|------|
| `seed_brokers*` | Broker 地址列表 |
| `topics*` | 要消费的 Topic 列表 |
| `regexp_topics*` | 是否将 `topics` 按正则匹配多 Topic |
| `auto_replay_nacks*` | 下游拒绝（nack）时是否自动重放 |
| `fetch_max_bytes*` | 单次 fetch 最大字节数 |
| `fetch_max_partition_bytes*` | 单分区单次 fetch 最大字节数 |
| `fetch_max_wait*` | Broker 等待凑齐最小字节的最长时间 |

常用非必填项：

| 参数 | 说明 |
|------|------|
| `consumer_group` | 消费组；填写后由组内协调分区与位移 |

## 参数建议

- `seed_brokers`：填写至少一个可达 broker 地址
- `topics`：按业务域拆分 Topic，避免单 Topic 过载
- `consumer_group`：按消费语义规划，避免重复消费；建议使用可读名称（如 `flowmq-source-mqtt-kafka`）

## 注意事项

- 关注分区数与任务数的匹配关系
- 明确 offset 管理策略与重放策略
- 接入后的字段清洗与映射见 [Processors](data-integration-processors-expressions.md)
