第 10 章 · 消息队列与Kafka 当系统之间需要传递数据与事件,消息队列登场。Kafka 是事实标准的分布式事件流平台——它本质上是一个"分布式日志"。本节讲透消息队列的价值、Kafka 的五大核心抽象(Topic/Partition/Broker/Consumer Group/Offset)、消息不丢不重的工程手段,以及 ZooKeeper 到 KRaft 的演进。 学习目标 说出消息队列解决的四大问题:削峰、解耦、异步、顺序 描述 Kafka 的 Topic/Partition/Broker/Offset 组织方式,解释"分区内有序、全局无序" 掌握 Consumer Group 的负载均衡规则与"消费者多于分区则闲置"的约束 理解 Offset
当系统之间需要传递数据与事件,消息队列登场。Kafka 是事实标准的分布式事件流平台——它本质上是一个"分布式日志"。本节讲透消息队列的价值、Kafka 的五大核心抽象(Topic/Partition/Broker/Consumer Group/Offset)、消息不丢不重的工程手段,以及 ZooKeeper 到 KRaft 的演进。
在引入 Kafka 之前,先回答"为什么需要消息队列"。四个核心价值:
| 价值 | 说明 | 场景 |
|---|---|---|
| 削峰填谷 | 突发流量先写入队列,下游按自身能力消费 | 秒杀、大促、突发的 Webhook 洪峰 |
| 解耦 | 生产者不关心消费者是谁、有几个、是否在线 | 订单系统发事件,多个下游订阅 |
| 异步 | 非核心链路后台处理,主链路快速返回 | 下单后发短信、更新推荐 |
| 顺序保障 | 单分区内保证消息有序 | 同一订单的状态流转必须有序 |
经典比喻:消息队列像邮局——寄信人(生产者)把信投进邮筒就不管了,收信人(消费者)什么时候取、怎么处理,互不打扰。
Kafka 自称"分布式事件流平台",本质是分布式的提交日志(Commit Log):消息追加写入、顺序读取、不删除(保留期内),像一本持续追加的账本。
五个核心抽象:
Topic(主题):消息的逻辑分类,如 order-events、user-login。生产者往 Topic 写,消费者从 Topic 读。
Partition(分区):Topic 被拆成多个分区,是并行与有序的基本单位。每条消息在分区内按写入顺序排列、分区内有序;但跨分区的消息之间没有全局顺序。因此"Kafka 保证消息有序"是常见误解——只能说"单分区内有序"。想要全局有序,只能使用单分区(牺牲并行度)。
Broker(代理节点):Kafka 集群中的一台服务器。一个 Topic 的分区会分布在多个 Broker 上,每个分区有 Leader 副本负责读写、Follower 副本负责备份(副本因子通常 3),Broker 宕机时 Leader 自动切换。
Consumer Group(消费组):一组共同消费某 Topic 的消费者。核心规则:
Offset(偏移量):分区内消息的位置编号,消费者用它记录"读到哪里了"。这是消息语义的开关:
Kafka 集群由多个 Broker 组成,元数据与 Leader 选举历史上由 ZooKeeper 承担:管理 Topic/分区元数据、Broker 注册、Controller 选举。Kafka 2.8 起引入 KRaft,逐步用 Kafka 内置的 Raft 协议替代 ZooKeeper,到 3.x 已可完全去 ZooKeeper 运行——好处是组件更少、部署更简单、元数据管理更可靠。
副本机制:每个分区有 1 个 Leader + N 个 Follower。所有读写都走 Leader,Follower 异步同步数据。ISR(In-Sync Replicas)是与 Leader 保持同步的副本集合;同步落后的副本被踢出 ISR,Leader 宕机时从 ISR 中选出新 Leader。
生产者的可靠写入配置要点:acks 决定确认级别——acks=0 不等确认(可能丢)、acks=1 Leader 确认(Leader 宕机可能丢)、acks=all 全部 ISR 确认(最可靠,配 min.insync.replicas 保证至少几个副本同步)。
"消息丢了"通常发生在三个环节,对策不同:
| 环节 | 丢消息的典型原因 | 对策 |
|---|---|---|
| 生产者 → Broker | acks=0/1,网络异常或 Leader 宕机 | acks=all + 重试(retries) |
| Broker 内部 | 副本数不足,Leader 宕机丢数据 | 副本因子 ≥ 3 + min.insync.replicas |
| Broker → 消费者 | Offset 先提交后消费 | 先消费后提交(可能重复但不会丢)+ 幂等消费 |
幂等是应对重复消费的标准答案:消费逻辑天然幂等(如 SET 操作),或记录已处理的消息 ID 去重。
| 系统 | 模型 | 特点 | 典型场景 |
|---|---|---|---|
| Kafka | 分布式日志、拉取模型 | 高吞吐、消息可回溯(重放)、分区有序 | 事件流、日志管道、大数据 |
| RabbitMQ | 队列/交换机、推送模型 | 灵活路由、功能丰富、低延迟 | 任务分发、业务解耦 |
| Redis 流/列表 | 轻量队列 | 简单易用、基于内存 | 小型异步任务 |
| Pulsar | 分层存储、多租户 | 存算分离、云原生 | 大规模多租户事件流 |
选型思路:数据量大、要重放、要吞吐 → Kafka;路由灵活、任务型 → RabbitMQ;简单场景不想引入重组件 → Redis。
到这里,第 10 章《可观测性与数据》收束:可观测性让我们看得见系统,数据库让数据存得下,消息队列让事件传得动——三大支柱构成生产系统的"神经系统"。下一章开始,你将把这些能力串起来,进入完整的 DevOps 工程实践与案例综合应用。