第 12 章 Kafka:消息队列与异步
Kafka 在系统设计里的核心价值,是把同步强耦合链路改造成可缓冲、可回放、可扩展的异步数据流。它并不只是“消息中转站”,而是围绕顺序写、分区并行和副本复制构建出来的一整套高吞吐日志系统。
Kafka 基础模型
Kafka 的基本对象包括 Broker、Topic、Partition、Producer、Consumer 和 Consumer Group。理解它们之间的关系,是看懂顺序性、吞吐与高可用的前提。
| 组件 | 作用 |
|---|---|
| Broker | 存储和转发消息的服务节点 |
| Topic | 业务主题,逻辑上的消息分类 |
| Partition | Topic 的物理分片,承载并行与扩展 |
| Producer | 生产消息 |
| Consumer | 消费消息 |
| Consumer Group | 多消费者协同消费的逻辑组 |
Kafka 的一个关键设计是“分区内有序、全局无序”。只要消息落在同一个 Partition,就会按 offset 单调递增排列;但跨 Partition 不承诺全局顺序。因此,业务若要求“同一订单内的状态变更严格有序”,通常做法不是把 Topic 只设成一个分区,而是按订单 ID 作为 key 路由到固定分区。
副本模型则决定了 Kafka 的高可用边界:
- Leader:对外承担读写。
- Follower:从 Leader 复制数据。
- ISR:与 Leader 保持同步的副本集合。
Leader 故障时,只从 ISR 中选新 Leader,才能尽量避免数据倒退。如果允许落后副本直接上位,就会提升可用性但增加丢数据风险。
写入、消费与消费组
写入路径。 Producer 发送消息时,通常会经历序列化、分区选择、批量打包、网络发送、Broker 追加写日志、Follower 同步、返回确认这一流程。影响写入体验的几个关键参数包括:
acks:控制等待多少副本确认。batch.size与linger.ms:控制批量发送效果。- 重试与幂等配置:决定失败重发后的可靠性和顺序风险。
消费路径。 Consumer 不是被推送消息,而是主动拉取。拉取模型让消费者能按自己的节奏处理,也方便批量读取与偏移量管理。消费完成后是否提交 offset,直接决定可靠性语义。
Consumer Group。 Consumer Group 用于实现负载均衡。一个分区在同一消费组内只能被一个消费者消费,所以:
- 分区数少于消费者数时,多余消费者会空闲。
- 要提升组内并行度,分区数必须足够。
消费组最大的治理点在于 Rebalance。以下情况会触发重新分配:
- 消费者加入或退出。
- Topic 分区数变化。
- 心跳超时或消费线程长时间卡住。
Rebalance 的代价并不小,它会导致一段时间消费暂停,还可能造成重复消费。因此要特别关注超时参数、消费逻辑时长和实例重启策略。
顺序性、可靠性与幂等
顺序性。 Kafka 的顺序保证是分层次的:
- 同一 Partition 内天然有序。
- 多 Partition 并行时,全局无序。
- Producer 开启重试且允许过多飞行中的请求时,可能破坏同分区重试顺序。
所以对顺序敏感的业务,一般要同时做三件事:
- 用业务 key 固定分区。
- 控制生产端重试配置,避免乱序重发。
- 让消费端按分区串行处理关键链路。
可靠性语义。
| 语义 | 特点 | 常见做法 |
|---|---|---|
| At-most-once | 可能丢,不重复 | 自动提交 offset,弱确认 |
| At-least-once | 不丢,但可能重复 | acks=all + 手动提交 offset |
| Exactly-once | 端到端成本最高 | Producer 幂等/事务 + 消费端幂等 |
常见的“不丢消息”配置是:
- Producer 用
acks=all - Broker 设置合理副本数和
min.insync.replicas - 关闭
unclean.leader.election - Consumer 业务处理成功后再提交 offset
但要明确,Kafka 的可靠性是链路整体属性,不是单个参数就能保证。比如 acks=all 如果此时 ISR 已经只剩 Leader,本质上依然是退化状态。
幂等。 Kafka 的 Producer 幂等可以减少重试导致的重复写入,但它主要解决“Producer 到 Broker”这一段。真正的端到端幂等,仍然要靠业务侧保证,例如:
- 用订单号做数据库唯一键。
- 用状态机限制非法重复更新。
- 用去重表或 Redis 记录已处理消息 ID。
因此,“Kafka 开了幂等就万事大吉”是一个常见误解。
高吞吐设计原理
Kafka 的高吞吐并不是因为它“把数据放进内存队列”,而是因为它把磁盘和网络用到了极致。
顺序写磁盘。 Kafka 把消息追加到日志尾部,避免了随机写的磁盘寻道开销。顺序写对磁盘极其友好,这是它在持久化前提下仍能做出高吞吐的基础。
Page Cache。 Kafka 尽量利用操作系统页缓存,而不是自己在 JVM 堆中维护大块缓存。好处是:
- 热数据读写往往直接命中内存。
- JVM 堆不必过大,GC 压力更小。
- 文件缓存由 OS 统一管理,更适合顺序日志场景。
这也是为什么 Kafka 机器的内存不能被 JVM 堆吃光,否则 Page Cache 不足会让性能明显下滑。
零拷贝。 Kafka 在发送消息给消费者时会尽量利用 sendfile 等能力,让数据从 Page Cache 直接进入网卡,减少用户态与内核态之间的多次拷贝和上下文切换。
批量与压缩。 Producer 会把多条消息打成 batch,再统一发送;Broker 和网络层又能对 batch 做压缩。这意味着单条消息的协议开销和系统调用次数被均摊,从而进一步提升吞吐。
分区并行。 Partition 把单线程日志扩展成了“多条顺序日志并行处理”的模型。它是 Kafka 横向扩容的核心手段,但也是顺序性与运维复杂度的来源。
积压、延迟与回压处理
消费积压不是一个单一问题,它可能来自生产端突增、消费逻辑变慢、下游依赖超时、分区数不足,或者 Rebalance 频繁引发的停顿。
排查顺序。
- 先看 Lag,确认是哪些 Topic、哪些分区积压。
- 再看消费端处理时长,确认是不是业务逻辑慢。
- 检查消费者实例数与分区数的关系,避免“加了实例却没有并行收益”。
- 检查是否频繁 Rebalance、GC、网络抖动或下游超时。
常见应对手段。
- 临时扩消费者实例,但不超过分区数。
- 增加分区数以提升并行度,但要评估顺序和迁移影响。
- 把重逻辑异步化,缩短消费线程阻塞时间。
- 对非核心消息允许降级、丢弃或跳过。
- 调整
max.poll.interval.ms、批量参数和消费线程模型。
回压思路。 Kafka 自身能缓冲流量,但它不能无限吞掉下游故障。真正稳定的系统需要在 Producer、Consumer、下游服务之间形成回压闭环,例如:
- Producer 端限速或降级。
- Consumer 端控制并发和批量。
- 下游数据库或搜索服务承压时主动减速。
否则 Kafka 只是把问题从“实时失败”变成“延迟爆炸”。
常见线上问题与排查
Rebalance 频繁,导致重复消费。 当消费逻辑过慢、心跳超时或实例频繁重启时,消费组会不断重平衡。排查时重点看:
session.timeout.msmax.poll.interval.ms- 消费业务是否阻塞
- 是否缺少静态成员配置
分区数量不足,扩容无效。 很多“加机器不生效”的根因,不是消费者不够,而是 Topic 分区不够。一个分区只能被组内一个消费者处理,多加出来的实例根本拿不到任务。
acks=1 或副本配置不当导致丢消息。 只等 Leader 确认时,如果 Leader 写完但 Follower 还没同步就宕机,消息就可能丢失。此类问题通常要结合 acks、min.insync.replicas 和 ISR 监控一起看。
Page Cache 不足,吞吐骤降。 如果把 JVM 堆设得过大,OS 没有足够内存做页缓存,Kafka 会频繁读盘,吞吐和延迟都会恶化。
retention 配置不合理,磁盘打满。 Kafka 的本质是日志系统,不设置合理保留时间和容量上限,就等于默认“消息永久堆积”。日志类 Topic 特别容易在这件事上出事故。
重点监控项。
| 类别 | 指标 |
|---|---|
| 吞吐 | MessagesInPerSec、BytesInPerSec |
| 延迟 | 请求平均延迟、批量发送耗时 |
| 副本 | UnderReplicatedPartitions、ISR 收缩次数 |
| 消费 | Consumer Lag、Rebalance 次数 |
| 资源 | 磁盘使用率、Page Cache 压力、网络带宽 |
本章小结
Kafka 的工程价值可以总结为三层:
- 用 Topic 和 Partition 把同步调用改造成可扩展的异步日志流。
- 用副本、ISR、offset 与消费组支撑可靠性和高可用。
- 用顺序写、Page Cache、零拷贝、批量压缩支撑高吞吐。
真正落地时,最难的往往不是“会不会发消息”,而是能否在顺序、可靠性、吞吐、积压治理之间做清晰取舍。
本章面试题与追问
-
Kafka 的核心组件有哪些,它和普通消息队列最大的差异是什么? 追问:为什么说 Kafka 更像分布式日志系统,而不只是消息转发器?
-
为什么 Kafka 只能保证分区内有序,不能天然保证全局有序? 追问:如果订单状态要求严格有序,你会怎么设计分区策略?
-
Consumer Group 是如何实现负载均衡的? 追问:为什么消费者数量超过分区数量后,再扩容也没有收益?
-
Rebalance 是什么,为什么它会影响线上稳定性? 追问:如果线上频繁重复消费,你会优先检查哪几个超时和消费参数?
-
Kafka 如何保证消息不丢失? 追问:为什么
acks=all也不能脱离 ISR 状态单独谈可靠性? -
Producer 幂等和业务幂等有什么区别? 追问:为什么 Kafka 自带幂等仍然不能替代数据库唯一键或状态机设计?
-
Kafka 为什么这么快? 追问:顺序写、Page Cache、零拷贝、批量压缩分别解决了什么瓶颈?
-
什么是 HW 和 LEO? 追问:为什么消费者通常只能读到 HW 之前的数据?
-
Kafka 出现消费积压时,你会如何分层排查? 追问:如果 Lag 很高,但消费者 CPU 并不高,可能说明什么问题?
-
Topic 的 retention 应该如何设置? 追问:如果日志类 Topic 把磁盘打满,除了临时清理,你会如何从治理角度避免再次发生?