Kafka数据流向是生产者、Broker、消费者协同的闭环链路:生产者按key分区、批量压缩发送并依赖acks确认;Broker以Partition为单位存储日志、维护ISR动态同步;消费者主动拉取、组内分配Partition并手动提交offset保障可靠性。

Kafka 的数据流向不是单线程的“推”或“拉”,而是一套由生产者、Broker、消费者协同完成的闭环链路。理解它,关键在于看清每个环节的职责和衔接逻辑——数据从生成到最终被业务消费,每一步都依赖明确的协议与状态管理。
生产者:消息发出前的关键控制
生产者是数据进入 Kafka 的第一道关卡。它不直接写磁盘,而是把消息发给对应 Partition 的 Leader Broker,并根据配置决定“等不等确认”:
- acks=all:必须等所有 ISR 副本同步完成才返回成功,这是保障不丢消息的底线配置;
- key 决定分区:比如用订单 ID 作 key,相同 ID 的消息总落在同一 Partition,保证局部有序;
- 批量发送 + 压缩:消息先缓存在 RecordAccumulator 中,达到 batch.size 或 linger.ms 后统一发出,减少网络开销。
Broker:存储与副本同步的核心枢纽
Broker 接收消息后,不是简单存盘,而是按 Partition 组织成日志段(segment),并维护索引文件加速查询:
- Leader 写入本地 .log 文件后,触发 Follower 拉取复制,是否完成取决于 acks 设置;
- ISR 列表动态维护:Follower 落后太多会被踢出,确保“同步中”的副本始终可用;
- 消息在 Partition 内严格按 offset 递增,但不同 Partition 之间无序——全局有序需靠业务层设计。
消费者:按需拉取与进度管理
消费者不被动接收,而是主动向 Broker 发起 fetch 请求,指定要读哪个 Partition 的哪段 offset:
- 消费组(Consumer Group)内自动分配 Partition,一个 Partition 只由组内一个 Consumer 拉取;
- offset 提交方式影响可靠性:自动提交可能丢数据,手动提交需在业务处理成功后再调用 commitSync();
- 如果消费者崩溃,Group Coordinator 会触发 Rebalance,重新分配 Partition,新 Consumer 从上次提交的 offset 继续读。
端到端可靠性:三段式保障缺一不可
一条消息不丢失,需要生产、存储、消费三个环节全部守住底线:
- 生产端开启幂等性(enable.idempotence=true)+ acks=all + 重试机制;
- Broker 端保证至少一个 ISR 副本存活,且日志刷盘策略(flush.messages/flush.ms)合理;
- 消费端处理完业务逻辑再提交 offset,避免“消费了但没处理完就提交”的空转风险。


















