消息中间件技术
消息队列是分布式系统中最重要的中间件之一。
一、核心主线
它在生产者和消费者之间增加一个消息存储与转发层:
1 | 生产者 |
生产者只负责把消息投递到队列,不需要等待消费者立即处理;消费者则可以根据自身能力选择消费消息的时机和速度。
本节重点介绍两部分内容:
- 消息队列的三大用途:异步化、流量削峰、服务解耦;
- Kafka 的核心概念、消息读写原理和高可用机制。
[!note]
本节对 Kafka 的介绍以依赖 ZooKeeper 的传统 Kafka 架构为背景。
二、消息队列的通信模式
1. 三个核心角色
| 角色 | 职责 |
|---|---|
| Producer | 创建并发送消息 |
| Message Queue | 保存消息,连接生产者与消费者 |
| Consumer | 从队列中读取并处理消息 |
可以将消息队列理解为电子邮箱:
- 发件人是生产者;
- 邮箱系统是消息队列;
- 收件人是消费者。
生产者只要把消息成功放入队列,就完成了当前阶段的通信,不需要等待消费者立即处理。
2. 通信特点
消息队列把生产和消费分成了两个相对独立的过程:
1 | 生产消息的速度、时间 |
因此:
- 生产者不必阻塞等待消费者;
- 消费者可以根据自身负载控制消费速度;
- 消费者短暂不可用时,消息可以暂存在队列中;
- 生产者与消费者不必同时在线。
三、消息队列的核心用途
1. 异步化
1.1 同步调用的问题
假设一个请求要依次经过 A、B、C、D 四个服务:
| 服务 | 处理耗时 |
|---|---|
| A:核心服务 | 10 ms |
| B:非核心服务 | 200 ms |
| C:非核心服务 | 300 ms |
| D:非核心服务 | 100 ms |
采用串行同步调用时:
1 | 用户 → A → B → C → D → 返回 |
非核心逻辑拉长了整个请求的响应时间。
1.2 使用消息队列异步执行
引入消息队列后:
1 | 用户请求 |
用户只需要等待 A 的核心逻辑,响应时间可由 610 ms 降低到约 10 ms。
1.3 异步化的价值
- 缩短用户请求链路;
- 降低接口响应时间;
- 非核心任务不再阻塞核心流程;
- 下游服务短暂故障时,不一定直接影响当前用户请求。
适合异步执行的任务通常具有以下特征:
- 不要求立即完成;
- 不影响当前请求的核心结果;
- 可以稍后处理;
- 处理失败后允许重试或补偿。
2. 流量削峰
2.1 突发流量的问题
如果后端服务每秒只能处理 100 个请求,却突然收到 10000 QPS:
1 | 10000 QPS → 服务 → 数据库 |
服务和数据库可能被瞬间压垮。
2.2 消息队列作为缓冲区
引入消息队列后:
1 | 突发流量 10000 QPS |
消息队列将短时间内的大量请求转化为较长时间内的平稳消费。
2.3 削峰的本质
流量削峰并没有减少总工作量,而是改变了工作执行的时间分布:
1 | 瞬时高峰 |
消费者因此获得处理请求的主动权,可以根据自身吞吐能力拉取消息,避免下游系统被突发流量击垮。
3. 服务解耦
3.1 直接 RPC 的耦合问题
假设点赞服务处理点赞请求后,需要通知:
- 热点服务;
- 策略服务;
- 未来可能增加的其他服务。
采用直接 RPC 时:
1 | 点赞服务 |
这会导致:
- 点赞服务需要知道所有下游服务;
- 新增消费者需要修改点赞服务;
- 下线消费者也需要修改点赞服务;
- 下游故障可能影响点赞主流程;
- 多团队之间形成紧密依赖。
3.2 基于事件的解耦
使用消息队列后,点赞服务只负责发布点赞事件:
1 | 点赞服务 |
任何需要点赞行为的服务都可以自行订阅消息,而无须改造点赞服务。
3.3 解耦的价值
- 生产者不需要知道消费者是谁;
- 新增或删除消费者不影响生产者;
- 不同团队可以独立开发和部署;
- 降低服务之间的直接依赖;
- 减少业务变更的影响范围。
四、Kafka 整体架构
Kafka 是一个分布式、高性能、可扩展的消息队列系统,可用于消息中间件和日志收集等场景。
传统 Kafka 架构可以概括为:
1 | Producer |
五、Kafka 的核心概念
1. Producer 与 Consumer
- Producer:生产并发送消息;
- Consumer:从 Kafka 中读取和处理消息。
2. Topic
Topic 表示消息的逻辑分类或消息类型。
例如:
1 | user-like-events:用户点赞事件 |
生产者向指定 Topic 发送消息,消费者订阅感兴趣的 Topic。
3. Partition
一个 Topic 的消息会分散存储在多个 Partition 中。
1 | Topic |
Partition 类似存储系统中的数据分片,其作用包括:
- 拆分海量消息;
- 将读写压力分散到多个 Broker;
- 提升 Topic 的吞吐能力;
- 为消费者并行消费提供基础。
顺序性范围:
- 同一个 Partition 内的消息基于写入顺序保存;
- 不同 Partition 之间不保证全局有序。
4. Broker
Broker 是 Kafka 的服务节点,负责:
- 接收生产者发送的消息;
- 将消息写入 Partition;
- 保存消息数据;
- 响应消费者的拉取请求;
- 管理 Partition 副本。
多个 Broker 共同组成 Kafka Cluster,每个 Broker 使用全局唯一的 Broker ID 标识。
5. Consumer Group
多个 Consumer 实例可以组成一个 Consumer Group。
Topic 面向 Consumer Group 进行消费:
- 同一个 Topic 可以被多个 Consumer Group 独立消费;
- 一个 Consumer Group 可以消费多个 Topic;
- 在一个 Consumer Group 内,一个 Partition 同一时刻只能分配给一个 Consumer;
- 一个 Consumer 可以同时消费多个 Partition。
六、ZooKeeper 在传统 Kafka 中的作用
1. Broker 注册与服务发现
Broker 启动后,将以下信息注册到 ZooKeeper:
- Broker ID;
- IP 地址;
- 端口号。
Broker 宕机或断开连接后,ZooKeeper 会删除相应节点信息,使集群感知 Broker 状态变化。
2. Topic 元信息管理
ZooKeeper 维护:
- Topic 拥有哪些 Partition;
- Partition 分布在哪些 Broker;
- Partition 副本与 Broker 的对应关系。
3. 生产者负载均衡
生产者可以根据 Broker 地址和 Topic 分区元信息,决定消息应写入哪个 Partition。
4. 消费者负载均衡
ZooKeeper 记录 Consumer 与 Partition 的分配关系。
当某个 Consumer 宕机时,Consumer Group 会发生 Rebalance,将原 Consumer 负责的 Partition 重新分配给其他 Consumer。
5. 消费进度管理
消费者通过 Offset 表示自己已经消费到 Partition 的哪个位置。
早期 Kafka 将消费 Offset 保存到 ZooKeeper;从书中提到的 Kafka 0.9 开始,Offset 改为保存在 Kafka Broker 的内部存储中,以改善写入性能。
七、生产者写入消息
1. Partition 选择规则
生产者向 Topic 发送消息前,需要决定消息进入哪个 Partition。
选择规则:
1 | 是否显式指定 Partition? |
2. Key 的意义
使用相同 Key 的消息通常会被路由到同一个 Partition。
这样可以利用 Partition 内部有序的特性,保证同一业务对象的消息顺序。
例如:
1 | Key = user_id |
同一个用户的操作事件会尽量写入同一个 Partition。
3. 顺序写磁盘
生产者找到目标 Partition 后,将消息发送给该 Partition 所在的 Broker。
Broker 收到消息后,按照追加方式将消息顺序写入磁盘。
顺序写减少了磁盘随机寻址,是 Kafka 获得高吞吐能力的重要原因之一。
八、消费者消费消息
1. Partition 与 Consumer 的关系
在同一个 Consumer Group 中:
1 | 一个 Partition → 最多由一个 Consumer 消费 |
例如,一个 Topic 有 10 个 Partition:
- Consumer Group 中有 5 个 Consumer:每个 Consumer 可负责多个 Partition;
- Consumer Group 中有 10 个 Consumer:最多可实现 10 路并行;
- Consumer Group 中有 20 个 Consumer:其中 10 个 Consumer 会处于空闲状态。
因此:
Consumer Group 的最大有效并行度通常受 Partition 数量限制。
2. Pull 消费模式
消费者使用拉取模式从 Broker 获取消息:
1 | Consumer → Broker:我现在可以处理,请给我消息 |
由消费者主动控制拉取速度,可以:
- 根据自身处理能力消费;
- 在繁忙时降低拉取速度;
- 避免 Broker 持续推送导致消费者过载;
- 为流量削峰提供实现基础。
3. Offset
Offset 表示 Consumer 在某个 Partition 中的消费位置。
记录 Offset 后,Consumer 重启或 Consumer Group 发生 Rebalance 时,可以从之前的位置继续消费。
九、Kafka 的多副本高可用
1. Leader 与 Follower
Kafka 允许一个 Partition 保存多个副本:
1 | Partition |
- Producer 向 Leader 写消息;
- Consumer 从 Leader 读取消息;
- Follower 从 Leader 复制数据;
- Leader 故障后,Follower 可以被提升为新 Leader。
Follower 的主要作用是数据备份和故障接管,而不是正常情况下直接处理客户端读写。
2. 副本的跨 Broker 分布
同一个 Partition 的多个副本应尽量部署在不同 Broker 上。
错误方式:
1 | Broker 1 |
Broker 1 宕机后,所有副本同时不可用。
正确思路:
1 | Broker 1:Leader |
这样单个 Broker 故障后仍然存在可用副本。
十、同步复制、异步复制与 ISR
1. 同步复制
所有 Follower 都复制完成后,才认为消息写入成功。
优点
- 副本数据一致性高;
- 消息丢失风险较低。
缺点
- 任意一个 Follower 速度慢都会拖慢写入;
- 某个 Follower 宕机可能降低 Partition 可用性;
- 写入延迟较高。
2. 异步复制
Leader 收到消息后立即认为写入成功,不等待 Follower。
优点
- 写入性能高;
- 延迟低;
- 不容易被慢 Follower 影响。
缺点
- Leader 宕机时,尚未复制的消息可能丢失;
- 副本数据一致性较弱。
3. ISR 机制
Kafka 使用 ISR(In-Sync Replica,同步副本集合)在性能、一致性和可用性之间折中。
每个 Partition 的 Leader 维护一个 ISR 列表,列表中包含与 Leader 保持较好同步状态的 Follower。
Follower 出现以下情况时,可能被移出 ISR:
- 长时间没有向 Leader 发起复制;
- 数据进度落后 Leader 太多;
- 节点或网络发生异常。
书中所述的提交逻辑是:
1 | Leader 收到消息 |
ISR 的价值在于:
- 不必等待所有副本;
- 慢副本或故障副本可以被动态移出;
- 降低单个 Follower 对写入性能和可用性的影响;
- 同时保持一定程度的数据一致性。
十一、Controller 与故障恢复
1. Controller 的职责
Kafka 会在 Broker 中选举一个 Controller。
Controller 主要负责:
- 监控 Broker 状态变化;
- 处理 Partition Leader 选举;
- 处理副本重新分配;
- 将 Leader 和 Follower 变化通知相关 Broker。
在传统架构中,Broker 借助 ZooKeeper 分布式锁竞争 Controller:
1 | 哪个 Broker 先获得锁 |
2. Broker 故障后的恢复流程
当保存 Partition Leader 的 Broker 宕机:
- 故障 Broker 与 ZooKeeper 断开连接;
- ZooKeeper 删除该 Broker 的注册节点;
- ZooKeeper 通知 Controller;
- Controller 查询受影响的 Partition;
- Controller 从 Partition 的 ISR 中选择一个 Follower;
- Controller 通知相关 Broker 最新选举结果;
- 被选中的 Follower 提升为 Leader;
- 其他 Follower 改为向新 Leader 复制数据。
1 | Leader Broker 宕机 |
十二、Kafka 的关键设计权衡
| 设计 | 解决的问题 | 带来的价值 |
|---|---|---|
| Topic | 消息分类 | 生产者和消费者按业务类型通信 |
| Partition | 数据分片 | 提升存储容量和并发能力 |
| Consumer Group | 并行消费 | 扩展消费者处理能力 |
| Pull 模式 | 消费速度控制 | 防止消费者被推送流量压垮 |
| Offset | 消费进度 | 支持重启恢复和 Rebalance |
| Leader / Follower | 数据副本 | 提升 Partition 高可用性 |
| ISR | 动态同步副本 | 平衡一致性、性能与可用性 |
| Controller | 集群协调 | 自动选主和故障恢复 |
十三、核心架构思想
- 消息队列通过中间存储,使生产者和消费者不必同时工作。
- 异步化把非核心任务移出同步请求链路,从而降低响应延迟。
- 流量削峰使用队列缓冲突发请求,让消费者按自身能力处理。
- 服务解耦让生产者只发布事件,不需要感知具体消费者。
- Kafka 使用 Topic 对消息分类,使用 Partition 对消息分片。
- Partition 是 Kafka 扩展存储容量、写入吞吐和消费并行度的基础。
- Kafka 只保证单个 Partition 内部的顺序,不保证跨 Partition 全局有序。
- 同一个 Consumer Group 内,一个 Partition 只能由一个 Consumer 消费。
- 消费者采用 Pull 模式,可以自主控制消费速度。
- Partition 多副本通过 Leader 和 Follower 实现数据备份。
- 副本应分布在不同 Broker,避免单机故障导致全部副本丢失。
- ISR 是同步复制与异步复制之间的折中方案。
- Controller 负责 Partition Leader 选举和故障恢复。
- 消息系统设计需要在吞吐量、顺序性、一致性、可用性和消费延迟之间进行权衡。
十四、一句话总结
消息队列通过异步化、流量削峰和服务解耦提升分布式系统的性能与稳定性;Kafka 则通过 Topic 分类、Partition 分片、Consumer Group 并行消费、Offset 进度管理、Leader/Follower 多副本、ISR 同步集合以及 Controller 自动选主,实现高吞吐、可扩展和高可用的消息存储与消费体系。