Kafka Source 集成
本文介绍如何在管道的 Source 步骤中使用 Kafka 作为数据来源。
适用场景
- 从 Kafka Topic 持续消费业务事件
- 将存量 Kafka 数据接入 FlowMQ 处理链路
前提条件
配置步骤
- 在“数据管道”中创建管道,完成基础信息配置。
- 在 Source 步骤选择
Kafka。 - 配置
seed_brokers、topics、consumer_group等参数。 - 点击“测试”验证连接。
- 继续完成 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