消息中间件技术

消息队列是分布式系统中最重要的中间件之一。

一、核心主线

它在生产者和消费者之间增加一个消息存储与转发层:

1
2
3
4
5
生产者
↓ 发送消息
消息队列
↓ 消费消息
消费者

生产者只负责把消息投递到队列,不需要等待消费者立即处理;消费者则可以根据自身能力选择消费消息的时机和速度。

本节重点介绍两部分内容:

  1. 消息队列的三大用途:异步化、流量削峰、服务解耦
  2. Kafka 的核心概念、消息读写原理和高可用机制。

[!note]
本节对 Kafka 的介绍以依赖 ZooKeeper 的传统 Kafka 架构为背景。


二、消息队列的通信模式

1. 三个核心角色

角色 职责
Producer 创建并发送消息
Message Queue 保存消息,连接生产者与消费者
Consumer 从队列中读取并处理消息

可以将消息队列理解为电子邮箱:

  • 发件人是生产者;
  • 邮箱系统是消息队列;
  • 收件人是消费者。

生产者只要把消息成功放入队列,就完成了当前阶段的通信,不需要等待消费者立即处理。

2. 通信特点

消息队列把生产和消费分成了两个相对独立的过程:

1
2
3
生产消息的速度、时间

消费消息的速度、时间

因此:

  • 生产者不必阻塞等待消费者;
  • 消费者可以根据自身负载控制消费速度;
  • 消费者短暂不可用时,消息可以暂存在队列中;
  • 生产者与消费者不必同时在线。

三、消息队列的核心用途

1. 异步化

1.1 同步调用的问题

假设一个请求要依次经过 A、B、C、D 四个服务:

服务 处理耗时
A:核心服务 10 ms
B:非核心服务 200 ms
C:非核心服务 300 ms
D:非核心服务 100 ms

采用串行同步调用时:

1
2
3
用户 → A → B → C → D → 返回

总耗时 = 10 + 200 + 300 + 100 = 610 ms

非核心逻辑拉长了整个请求的响应时间。

1.2 使用消息队列异步执行

引入消息队列后:

1
2
3
4
5
6
7
8
9
10
11
12
用户请求

A 完成核心逻辑

发送消息到队列

立即向用户返回

消息队列
├─→ B 异步处理
├─→ C 异步处理
└─→ D 异步处理

用户只需要等待 A 的核心逻辑,响应时间可由 610 ms 降低到约 10 ms。

1.3 异步化的价值

  • 缩短用户请求链路;
  • 降低接口响应时间;
  • 非核心任务不再阻塞核心流程;
  • 下游服务短暂故障时,不一定直接影响当前用户请求。

适合异步执行的任务通常具有以下特征:

  • 不要求立即完成;
  • 不影响当前请求的核心结果;
  • 可以稍后处理;
  • 处理失败后允许重试或补偿。

2. 流量削峰

2.1 突发流量的问题

如果后端服务每秒只能处理 100 个请求,却突然收到 10000 QPS:

1
10000 QPS → 服务 → 数据库

服务和数据库可能被瞬间压垮。

2.2 消息队列作为缓冲区

引入消息队列后:

1
2
3
4
5
6
7
突发流量 10000 QPS

消息队列暂存

消费者按自身能力以 100 QPS 处理

数据库

消息队列将短时间内的大量请求转化为较长时间内的平稳消费。

2.3 削峰的本质

流量削峰并没有减少总工作量,而是改变了工作执行的时间分布:

1
2
3
4
5
瞬时高峰

队列积压

持续、稳定地消费

消费者因此获得处理请求的主动权,可以根据自身吞吐能力拉取消息,避免下游系统被突发流量击垮。


3. 服务解耦

3.1 直接 RPC 的耦合问题

假设点赞服务处理点赞请求后,需要通知:

  • 热点服务;
  • 策略服务;
  • 未来可能增加的其他服务。

采用直接 RPC 时:

1
2
3
点赞服务
├─RPC→ 热点服务
└─RPC→ 策略服务

这会导致:

  • 点赞服务需要知道所有下游服务;
  • 新增消费者需要修改点赞服务;
  • 下线消费者也需要修改点赞服务;
  • 下游故障可能影响点赞主流程;
  • 多团队之间形成紧密依赖。

3.2 基于事件的解耦

使用消息队列后,点赞服务只负责发布点赞事件:

1
2
3
4
5
6
点赞服务
↓ 发布点赞事件
消息队列
├─→ 热点服务
├─→ 策略服务
└─→ 其他服务

任何需要点赞行为的服务都可以自行订阅消息,而无须改造点赞服务。

3.3 解耦的价值

  • 生产者不需要知道消费者是谁;
  • 新增或删除消费者不影响生产者;
  • 不同团队可以独立开发和部署;
  • 降低服务之间的直接依赖;
  • 减少业务变更的影响范围。

四、Kafka 整体架构

Kafka 是一个分布式、高性能、可扩展的消息队列系统,可用于消息中间件和日志收集等场景。

传统 Kafka 架构可以概括为:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
Producer
↓ 生产消息
Kafka Cluster
├─ Broker 1
│ └─ Topic / Partition
├─ Broker 2
│ └─ Topic / Partition
└─ Broker N
└─ Topic / Partition
↓ 消费消息
Consumer Group
├─ Consumer 1
├─ Consumer 2
└─ Consumer N

ZooKeeper:维护集群元信息和协调信息

五、Kafka 的核心概念

1. Producer 与 Consumer

  • Producer:生产并发送消息;
  • Consumer:从 Kafka 中读取和处理消息。

2. Topic

Topic 表示消息的逻辑分类或消息类型。

例如:

1
2
3
user-like-events:用户点赞事件
order-created-events:订单创建事件
login-events:用户登录事件

生产者向指定 Topic 发送消息,消费者订阅感兴趣的 Topic。

3. Partition

一个 Topic 的消息会分散存储在多个 Partition 中。

1
2
3
4
Topic
├─ Partition 0
├─ Partition 1
└─ Partition 2

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
2
3
4
5
6
7
是否显式指定 Partition?
├─ 是 → 写入指定 Partition
└─ 否

是否设置消息 Key?
├─ 是 → 对 Key 做哈希,选择 Partition
└─ 否 → 轮询选择 Partition

2. Key 的意义

使用相同 Key 的消息通常会被路由到同一个 Partition。

这样可以利用 Partition 内部有序的特性,保证同一业务对象的消息顺序。

例如:

1
Key = user_id

同一个用户的操作事件会尽量写入同一个 Partition。

3. 顺序写磁盘

生产者找到目标 Partition 后,将消息发送给该 Partition 所在的 Broker。

Broker 收到消息后,按照追加方式将消息顺序写入磁盘。

顺序写减少了磁盘随机寻址,是 Kafka 获得高吞吐能力的重要原因之一。


八、消费者消费消息

1. Partition 与 Consumer 的关系

在同一个 Consumer Group 中:

1
2
一个 Partition → 最多由一个 Consumer 消费
一个 Consumer → 可以消费多个 Partition

例如,一个 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
2
3
4
Partition
├─ Leader
├─ Follower 1
└─ Follower 2
  • Producer 向 Leader 写消息;
  • Consumer 从 Leader 读取消息;
  • Follower 从 Leader 复制数据;
  • Leader 故障后,Follower 可以被提升为新 Leader。

Follower 的主要作用是数据备份和故障接管,而不是正常情况下直接处理客户端读写。

2. 副本的跨 Broker 分布

同一个 Partition 的多个副本应尽量部署在不同 Broker 上。

错误方式:

1
2
3
4
Broker 1
├─ Leader
├─ Follower 1
└─ Follower 2

Broker 1 宕机后,所有副本同时不可用。

正确思路:

1
2
3
Broker 1:Leader
Broker 2:Follower 1
Broker 3:Follower 2

这样单个 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
2
3
4
5
Leader 收到消息

ISR 中的 Follower 完成同步确认

Leader 认为消息写入成功

ISR 的价值在于:

  • 不必等待所有副本;
  • 慢副本或故障副本可以被动态移出;
  • 降低单个 Follower 对写入性能和可用性的影响;
  • 同时保持一定程度的数据一致性。

十一、Controller 与故障恢复

1. Controller 的职责

Kafka 会在 Broker 中选举一个 Controller。

Controller 主要负责:

  • 监控 Broker 状态变化;
  • 处理 Partition Leader 选举;
  • 处理副本重新分配;
  • 将 Leader 和 Follower 变化通知相关 Broker。

在传统架构中,Broker 借助 ZooKeeper 分布式锁竞争 Controller:

1
2
3
哪个 Broker 先获得锁

哪个 Broker 成为 Controller

2. Broker 故障后的恢复流程

当保存 Partition Leader 的 Broker 宕机:

  1. 故障 Broker 与 ZooKeeper 断开连接;
  2. ZooKeeper 删除该 Broker 的注册节点;
  3. ZooKeeper 通知 Controller;
  4. Controller 查询受影响的 Partition;
  5. Controller 从 Partition 的 ISR 中选择一个 Follower;
  6. Controller 通知相关 Broker 最新选举结果;
  7. 被选中的 Follower 提升为 Leader;
  8. 其他 Follower 改为向新 Leader 复制数据。
1
2
3
4
5
6
7
8
9
10
11
Leader Broker 宕机

ZooKeeper 感知

通知 Controller

从 ISR 选举新 Leader

更新集群元信息

恢复 Partition 读写

十二、Kafka 的关键设计权衡

设计 解决的问题 带来的价值
Topic 消息分类 生产者和消费者按业务类型通信
Partition 数据分片 提升存储容量和并发能力
Consumer Group 并行消费 扩展消费者处理能力
Pull 模式 消费速度控制 防止消费者被推送流量压垮
Offset 消费进度 支持重启恢复和 Rebalance
Leader / Follower 数据副本 提升 Partition 高可用性
ISR 动态同步副本 平衡一致性、性能与可用性
Controller 集群协调 自动选主和故障恢复

十三、核心架构思想

  1. 消息队列通过中间存储,使生产者和消费者不必同时工作。
  2. 异步化把非核心任务移出同步请求链路,从而降低响应延迟。
  3. 流量削峰使用队列缓冲突发请求,让消费者按自身能力处理。
  4. 服务解耦让生产者只发布事件,不需要感知具体消费者。
  5. Kafka 使用 Topic 对消息分类,使用 Partition 对消息分片。
  6. Partition 是 Kafka 扩展存储容量、写入吞吐和消费并行度的基础。
  7. Kafka 只保证单个 Partition 内部的顺序,不保证跨 Partition 全局有序。
  8. 同一个 Consumer Group 内,一个 Partition 只能由一个 Consumer 消费。
  9. 消费者采用 Pull 模式,可以自主控制消费速度。
  10. Partition 多副本通过 Leader 和 Follower 实现数据备份。
  11. 副本应分布在不同 Broker,避免单机故障导致全部副本丢失。
  12. ISR 是同步复制与异步复制之间的折中方案。
  13. Controller 负责 Partition Leader 选举和故障恢复。
  14. 消息系统设计需要在吞吐量、顺序性、一致性、可用性和消费延迟之间进行权衡。

十四、一句话总结

消息队列通过异步化、流量削峰和服务解耦提升分布式系统的性能与稳定性;Kafka 则通过 Topic 分类、Partition 分片、Consumer Group 并行消费、Offset 进度管理、Leader/Follower 多副本、ISR 同步集合以及 Controller 自动选主,实现高吞吐、可扩展和高可用的消息存储与消费体系。

© 2026 DadaVinCi's Blog

Elegant theme by Shiro · Made by Acris with ❤️