Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

第 12 章 Kafka:消息队列与异步

Kafka 在系统设计里的核心价值,是把同步强耦合链路改造成可缓冲、可回放、可扩展的异步数据流。它并不只是“消息中转站”,而是围绕顺序写、分区并行和副本复制构建出来的一整套高吞吐日志系统。

Kafka 基础模型

Kafka 的基本对象包括 Broker、Topic、Partition、Producer、Consumer 和 Consumer Group。理解它们之间的关系,是看懂顺序性、吞吐与高可用的前提。

组件作用
Broker存储和转发消息的服务节点
Topic业务主题,逻辑上的消息分类
PartitionTopic 的物理分片,承载并行与扩展
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.sizelinger.ms:控制批量发送效果。
  • 重试与幂等配置:决定失败重发后的可靠性和顺序风险。

消费路径。 Consumer 不是被推送消息,而是主动拉取。拉取模型让消费者能按自己的节奏处理,也方便批量读取与偏移量管理。消费完成后是否提交 offset,直接决定可靠性语义。

Consumer Group。 Consumer Group 用于实现负载均衡。一个分区在同一消费组内只能被一个消费者消费,所以:

  • 分区数少于消费者数时,多余消费者会空闲。
  • 要提升组内并行度,分区数必须足够。

消费组最大的治理点在于 Rebalance。以下情况会触发重新分配:

  1. 消费者加入或退出。
  2. Topic 分区数变化。
  3. 心跳超时或消费线程长时间卡住。

Rebalance 的代价并不小,它会导致一段时间消费暂停,还可能造成重复消费。因此要特别关注超时参数、消费逻辑时长和实例重启策略。

顺序性、可靠性与幂等

顺序性。 Kafka 的顺序保证是分层次的:

  • 同一 Partition 内天然有序。
  • 多 Partition 并行时,全局无序。
  • Producer 开启重试且允许过多飞行中的请求时,可能破坏同分区重试顺序。

所以对顺序敏感的业务,一般要同时做三件事:

  1. 用业务 key 固定分区。
  2. 控制生产端重试配置,避免乱序重发。
  3. 让消费端按分区串行处理关键链路。

可靠性语义。

语义特点常见做法
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 频繁引发的停顿。

排查顺序。

  1. 先看 Lag,确认是哪些 Topic、哪些分区积压。
  2. 再看消费端处理时长,确认是不是业务逻辑慢。
  3. 检查消费者实例数与分区数的关系,避免“加了实例却没有并行收益”。
  4. 检查是否频繁 Rebalance、GC、网络抖动或下游超时。

常见应对手段。

  • 临时扩消费者实例,但不超过分区数。
  • 增加分区数以提升并行度,但要评估顺序和迁移影响。
  • 把重逻辑异步化,缩短消费线程阻塞时间。
  • 对非核心消息允许降级、丢弃或跳过。
  • 调整 max.poll.interval.ms、批量参数和消费线程模型。

回压思路。 Kafka 自身能缓冲流量,但它不能无限吞掉下游故障。真正稳定的系统需要在 Producer、Consumer、下游服务之间形成回压闭环,例如:

  • Producer 端限速或降级。
  • Consumer 端控制并发和批量。
  • 下游数据库或搜索服务承压时主动减速。

否则 Kafka 只是把问题从“实时失败”变成“延迟爆炸”。

常见线上问题与排查

Rebalance 频繁,导致重复消费。 当消费逻辑过慢、心跳超时或实例频繁重启时,消费组会不断重平衡。排查时重点看:

  • session.timeout.ms
  • max.poll.interval.ms
  • 消费业务是否阻塞
  • 是否缺少静态成员配置

分区数量不足,扩容无效。 很多“加机器不生效”的根因,不是消费者不够,而是 Topic 分区不够。一个分区只能被组内一个消费者处理,多加出来的实例根本拿不到任务。

acks=1 或副本配置不当导致丢消息。 只等 Leader 确认时,如果 Leader 写完但 Follower 还没同步就宕机,消息就可能丢失。此类问题通常要结合 acksmin.insync.replicas 和 ISR 监控一起看。

Page Cache 不足,吞吐骤降。 如果把 JVM 堆设得过大,OS 没有足够内存做页缓存,Kafka 会频繁读盘,吞吐和延迟都会恶化。

retention 配置不合理,磁盘打满。 Kafka 的本质是日志系统,不设置合理保留时间和容量上限,就等于默认“消息永久堆积”。日志类 Topic 特别容易在这件事上出事故。

重点监控项。

类别指标
吞吐MessagesInPerSecBytesInPerSec
延迟请求平均延迟、批量发送耗时
副本UnderReplicatedPartitions、ISR 收缩次数
消费Consumer Lag、Rebalance 次数
资源磁盘使用率、Page Cache 压力、网络带宽

本章小结

Kafka 的工程价值可以总结为三层:

  • 用 Topic 和 Partition 把同步调用改造成可扩展的异步日志流。
  • 用副本、ISR、offset 与消费组支撑可靠性和高可用。
  • 用顺序写、Page Cache、零拷贝、批量压缩支撑高吞吐。

真正落地时,最难的往往不是“会不会发消息”,而是能否在顺序、可靠性、吞吐、积压治理之间做清晰取舍。

本章面试题与追问

  1. Kafka 的核心组件有哪些,它和普通消息队列最大的差异是什么? 追问:为什么说 Kafka 更像分布式日志系统,而不只是消息转发器?

  2. 为什么 Kafka 只能保证分区内有序,不能天然保证全局有序? 追问:如果订单状态要求严格有序,你会怎么设计分区策略?

  3. Consumer Group 是如何实现负载均衡的? 追问:为什么消费者数量超过分区数量后,再扩容也没有收益?

  4. Rebalance 是什么,为什么它会影响线上稳定性? 追问:如果线上频繁重复消费,你会优先检查哪几个超时和消费参数?

  5. Kafka 如何保证消息不丢失? 追问:为什么 acks=all 也不能脱离 ISR 状态单独谈可靠性?

  6. Producer 幂等和业务幂等有什么区别? 追问:为什么 Kafka 自带幂等仍然不能替代数据库唯一键或状态机设计?

  7. Kafka 为什么这么快? 追问:顺序写、Page Cache、零拷贝、批量压缩分别解决了什么瓶颈?

  8. 什么是 HW 和 LEO? 追问:为什么消费者通常只能读到 HW 之前的数据?

  9. Kafka 出现消费积压时,你会如何分层排查? 追问:如果 Lag 很高,但消费者 CPU 并不高,可能说明什么问题?

  10. Topic 的 retention 应该如何设置? 追问:如果日志类 Topic 把磁盘打满,除了临时清理,你会如何从治理角度避免再次发生?