CC 咖啡猫的工作空间 Coding Space

消息队列与事件驱动

原始素材:4-服务器与后端/4.13-message-queues.md

当系统耦合严重、流量突增时,如何保证核心链路稳定? 消息队列是现代分布式系统的"缓冲器"和"解耦器"。本文通过真实案例(餐厅叫号、快递分拣、秒杀系统)深入理解消息队列的设计哲学和工程实践。

1. 为什么要"消息队列"?

1.1 从一个真实案例说起:淘宝订单系统的演进

2012年,淘宝订单系统遭遇了一次严重故障。双11零点,流量瞬间涌入,订单服务直接调用库存服务、支付服务、物流服务...整个链路像多米诺骨牌一样接连倒下。

当时的架构(紧耦合):

 用户下单 → 订单服务 → 同步调用库存服务 → 同步调用支付服务 → 同步调用物流服务
                    ↓                    ↓                    ↓
                 响应 200ms           响应 500ms           响应 300ms

⚠️ 紧耦合的致命问题

  • 总响应时间 = 200 + 500 + 300 = 1000ms(用户等1秒)
  • 库存服务挂了 → 订单服务也挂(线程池耗尽)
  • 支付服务慢了 → 整个链路被拖慢
  • 无法水平扩展 → 只能垂直加机器(贵且有限)

改进后的架构(引入消息队列):

 用户下单 → 订单服务 → 发送"订单创建"消息 → 立即返回(50ms)
                              ↓
                        消息队列(Kafka)
                              ↓
        ┌─────────────┬─────────────┬─────────────┐
        ▼             ▼             ▼             ▼
   库存服务      支付服务      物流服务      通知服务
   (异步扣减)  (异步处理)  (异步创建)  (异步发送)

✨ 改进后的效果

  • 用户响应时间 = 50ms(体验提升20倍)
  • 库存服务挂了 → 消息暂存队列,恢复后继续处理
  • 支付服务慢了 → 不影响订单创建
  • 可以水平扩展 → 增加消费者实例即可

1.2 消息队列的生活化比喻

餐厅叫号系统

想象你去一家网红餐厅:

  • 没有叫号系统 : 顾客必须站在窗口等,窗口有限,后面的人排长队,餐厅压力大
  • 有叫号系统 : 点完餐给你一个号,你可以先坐下,叫到号了去取餐

消息队列就是软件系统的"叫号系统" :

  • 生产者(点餐的人) → 把消息(订单)放到队列
  • 队列(叫号机) → 暂存消息
  • 消费者(厨师) → 按自己的节奏处理消息

核心原理:入站流量 (蓝色)超过处理能力 (绿色直线)时,多余的请求会被存入消息队列 (橙色区域)。 一旦流量高峰过去,系统会继续全速处理队列中的积压,直到队列清空。这就是"削峰填谷"。

2. 什么是消息队列?(定义 + 核心三要素)

2.1 什么是"消息队列"?

消息队列(Message Queue, MQ) 是一个存储消息的容器,生产者把消息放进去,消费者从里面取消息处理。它实现了"异步通信"——发送方不需要等待接收方处理完成。

同步 vs 异步 :

  • 同步 : 像打电话,对方必须接听才能交流
  • 异步 : 像发微信,发了就行,对方有空再看

这就像你给朋友打电话(同步) vs 发微信(异步)。

2.2 消息队列的核心三要素

要素一:生产者(Producer)

职责 : 创建并发送消息到队列。

生活化比喻 : 生产者就像"寄件人",把信件(消息)送到邮局(队列)。

关键设计要点

  • 发送方式 : 同步发送(可靠但阻塞) vs 异步发送(高性能但需处理回调)
  • 消息确认 : 等待 Broker 确认(At Least Once) vs 发送即忘(At Most Once)
  • 失败处理 : 重试策略、本地日志备份、死信队列

要素二:消费者(Consumer)

职责 : 从队列获取消息并处理。

生活化比喻 : 消费者就像"收件人",从邮箱(队列)取出信件(消息)并处理。

关键设计要点

  • 消费模式 : 推模式(Push,Broker主动推送) vs 拉模式(Pull,消费者主动拉取)
  • 消费确认 : 自动 ACK(高效但可能丢消息) vs 手动 ACK(可靠但需处理超时)
  • 并发控制 : 单线程顺序消费 vs 多线程并行消费
  • 失败处理 : 重试策略、死信队列、补偿机制

要素三:Broker(消息代理)

职责 : 接收、存储、转发消息。

生活化比喻 : Broker 就像"邮局"或"快递中转站",负责接收、分拣、派送信件。

关键设计要点

  • 存储模型 : 内存存储(低延迟) vs 磁盘存储(高可靠)
  • 复制策略 : 主从复制、多副本同步
  • 高可用机制 : 集群部署、自动故障转移
  • 扩展性 : 分区(Partition)、分片(Sharding)

3. 核心问题一:如何解耦系统,避免"牵一发而动全身"?

3.1 紧耦合的悲剧:一个服务挂了,全盘皆输

场景还原 : 某电商平台的早期架构

订单服务直接调用下游服务:
┌─────────────┐
│  订单服务   │
└──────┬──────┘
       │
       ├───────────┬───────────┬───────────┐
       ▼           ▼           ▼           ▼
┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐
│库存服务  │ │支付服务  │ │物流服务  │ │短信服务  │
│  200ms   │ │  500ms   │ │  300ms   │ │  100ms   │
└──────────┘ └──────────┘ └──────────┘ └──────────┘
痛点 具体表现 后果
级联故障 库存服务挂掉,订单服务同步调用超时 订单服务线程池耗尽,无法处理新请求
响应延迟 必须等待所有下游服务响应 用户等待1秒以上,体验极差
扩展困难 新增积分服务,需要修改订单服务代码 发布周期变长,风险增加
资源浪费 订单服务必须等待短信服务 数据库连接被长时间占用

3.2 解耦方案:引入消息队列作为"中间层"

解耦后的架构:

 订单服务只负责发消息,不关心谁消费:

┌─────────────┐
│  订单服务   │ ──发送"订单创建"消息──┐
└─────────────┘                       │
                                      ▼
                            ┌───────────────────┐
                            │   消息队列         │
                            │  (Kafka/RabbitMQ) │
                            │   - 可靠存储       │
                            │   - 多副本         │
                            │   - 顺序保证       │
                            └─────────┬─────────┘
                                      │
              ┌───────────────────────┼───────────────────────┐
              │                       │                       │
              ▼                       ▼                       ▼
       ┌──────────────┐      ┌──────────────┐      ┌──────────────┐
       │  库存服务     │      │  支付服务     │      │  物流服务     │
       │  订阅订单事件 │      │  订阅订单事件 │      │  订阅订单事件 │
       └──────────────┘      └──────────────┘      └──────────────┘

❌ 紧耦合的致命问题

⚠️依赖性强: 通知服务宕机,订单创建失败

⚠️响应慢: 总耗时 = 300ms + 500ms + 400ms = 1200ms

⚠️扩展难: 增加新服务需要修改订单代码

✨ 解耦的好处

维度 解耦前 解耦后
故障隔离 库存挂 = 订单挂 库存挂,消息暂存队列,恢复后消费
响应时间 1000ms(同步等待) 50ms(发完消息即返回)
扩展性 新增服务需改订单代码 新增服务只需订阅主题
系统复杂度 订单服务强依赖下游 订单服务只依赖消息队列

3.3 解耦的本质:从"直接调用"到"事件驱动"

思维模式的转变:

 传统思维(命令式):
"订单服务命令库存服务:给我扣库存!"
  ↓ 直接调用
  ↓ 耦合度高,被调用方必须在线
  ↓ 调用方需要知道被调用方的接口

事件驱动思维(声明式):
"订单服务声明:订单已创建,谁关心谁来处理。"
  ↓ 发送事件到消息队列
  ↓ 解耦,消费者可以离线
  ↓ 生产者不需要知道消费者的存在

4. 核心问题二:如何削峰填谷,应对流量突增?

4.1 秒杀场景:10万QPS如何平稳处理?

场景还原 : 某电商平台双11秒杀活动,预计峰值10万QPS,但数据库只能承受1000 QPS。

直接冲击的后果:

 用户请求 ──→ 应用服务器 ──→ 数据库
  10万/s       10万/s          1000/s(极限)
                              ↓
                         连接池耗尽
                         响应超时
                         数据库崩溃
                              ↓
                         雪崩效应(所有依赖数据库的服务都挂)

QPS(Queries Per Second) : 每秒查询数,衡量系统吞吐量的指标。

10万QPS 意味着每秒有10万个请求,就像10万人同时冲进商店。

4.2 削峰填谷方案:消息队列作为"蓄水池"

架构设计:

 ┌───────────────────────────────────────────────────────────────────────┐
│                        秒杀系统架构                           │
├───────────────────────────────────────────────────────────────────────┤
│                                                               │
│  第一层:网关层(硬限流)                                        │
│  ┌───────────────────────────────────────────────────────────────┐  │
│  │  - 令牌桶限流:10万/s → 1万/s(丢弃90%请求)          │  │
│  │  - CDN 缓存静态资源(商品详情页)                       │  │
│  │  - 验证码/排队页面(削峰第一层)                        │  │
│  └───────────────────────────────────────────────────────────────┘  │
│                            │                                 │
│                            ▼                                 │
│  第二层:服务层(软限流)                                        │
│  ┌───────────────────────────────────────────────────────────────┐  │
│  │  - Nginx限流:1万/s → 5000/s                         │  │
│  │  - Redis预扣库存(原子操作):                       │  │
│  │    * 使用 Lua 脚本保证原子性                          │  │
│  │    * 库存不足直接返回"已售罄"                         │  │
│  │  - 生成订单令牌(排队凭证)                             │  │
│  └───────────────────────────────────────────────────────────────┘  │
│                            │                                 │
│                            ▼                                 │
│  第三层:消息队列层(核心削峰)                                   │
│  ┌───────────────────────────────────────────────────────────────┐  │
│  │  Kafka/RocketMQ:                                     │  │
│  │  - 批量写入:5000/s → 1000/s(数据库承受能力)         │  │
│  │  - 消息持久化:落盘保证不丢消息                         │  │
│  │  - 多分区并行消费:提升吞吐量                           │  │
│  │  - 消费位点管理:支持故障恢复                           │  │
│  │                                                       │  │
│  │  关键指标监控:                                         │  │
│  │  - 生产速率(Produce Rate)                             │  │
│  │  - 消费速率(Consume Rate)                             │  │
│  │  - 消息堆积(Lag)                                      │  │
│  └───────────────────────────────────────────────────────────────┘  │
│                            │                                 │
│                            ▼                                 │
│  第四层:消费层(异步处理)                                        │
│  ┌───────────────────────────────────────────────────────────────┐  │
│  │  订单处理消费者(多实例):                              │  │
│  │  - 从 Kafka 拉取消息(1000/s,匹配数据库能力)           │  │
│  │  - 数据库事务:创建订单 + 扣减库存                        │  │
│  │  - 更新订单状态为"已创建"                               │  │
│  │  - 发送订单创建成功通知(邮件/短信/推送)                  │  │
│  │  - 确认消息消费(ACK)                                   │  │
│  │                                                         │  │
│  │  消费者扩容策略:                                        │  │
│  │  - 当 Lag > 10000 时,自动增加消费者实例                  │  │
│  │  - 当 Lag < 1000 时,减少消费者实例(节省成本)           │  │
│  └───────────────────────────────────────────────────────────────┘  │
│                                                               │
└───────────────────────────────────────────────────────────────────────┘

4.3 削峰填谷的核心指标与监控

削峰填谷的工作原理:当入站流量(生产者速率)超过处理能力(消费者速率)时,多余的请求会被存入消息队列(缓冲区)。一旦流量高峰过去,系统会继续全速处理队列中的积压消息,直到队列清空。

关键系统指标监控表

指标 含义 告警阈值
处理能力 (Throughput) 200 req/s 后端消费者的最大处理速度
队列容量 (Queue Capacity) 2000 msgs 消息队列能暂存的最大消息数 超过80%需要扩容
当前入站流量 (Inbound Rate) 100 req/s 实时生产者发送速率 > 处理能力时需要削峰
队列积压 (Queue Lag/Backlog) 0 msgs 当前队列中堆积的消息数 > 队列容量50%需要告警
实际处理速率 (Consumption Rate) 0 req/s 消费者实际处理速度 < 生产速率说明消费滞后
丢弃请求 (Rate Limited Requests) 0 reqs 网关限流丢弃的请求数 > 0 说明流量超过预期

三个流量场景对比

场景 入站流量 处理能力 队列变化 系统状态
正常期 100 req/s 200 req/s 队列空(消费 > 生产) ✅ 系统轻负荷,空闲资源可用
高峰期 500 req/s 200 req/s 队列积压 3000 msgs(生产 > 消费) ⚠️ 进入削峰模式,消息堆积
恢复期 100 req/s 200 req/s 队列快速清空 ✅ 队列积压被逐步消费

削峰效果的数学模型

原始流量曲线(无削峰):        平滑后流量曲线(有削峰):

流量 req/s                      流量 req/s
   │ ╱╲ ← 高峰期                    │ ████████████████ ← 恒定消费
   │╱  ╲                            │
1000│    ╲                        200│
   │     ╲                          │
   │      └─                        │
   └──────────────────时间           └──────────────────时间
     0s   1s    5s                   0s               100s

峰值:1000 req/s(持续 1 秒)    峰值:200 req/s(持续 100 秒)
问题:数据库崩溃                  方案:平稳处理,无故障

关键公式

1. 队列堆积数量 = (入站流量 - 处理能力) × 持续时间
   示例:(500 - 200) req/s × 10s = 3000 msgs

2. 清空时间 = 队列堆积数量 / 处理能力
   示例:3000 msgs / 200 req/s = 15 秒

3. 总处理时间 = 消息入队时间 + 队列等待时间 + 处理时间
   示例:假设消息在队列等待 10 秒,则用户体验到的延迟为 10+ 秒

实时监控仪表板

╔═════════════════════════════════════════════════════════════════╗
║           消息队列削峰监控面板 (Real-time Dashboard)             ║
╠═════════════════════════════════════════════════════════════════╣
║                                                                 ║
║  入站流量 (Inbound Rate)        处理能力 (Throughput)          ║
║  ████████ 145 req/s            ███████████ 200 req/s          ║
║                                                                 ║
║  队列堆积 (Queue Lag)            队列健康度 (Queue Health)      ║
║  ██████ 1234 msgs              ✅ 61.7% (健康范围)             ║
║                                                                 ║
║  消费延迟 (Consumption Lag)      告警状态 (Alerts)             ║
║  ⏱ 6.2 秒                       🟢 无告警                      ║
║                                                                 ║
║  建议操作:正常运行,无需干预  ✓                               ║
║                                                                 ║
╚═════════════════════════════════════════════════════════════════╝

削峰场景下的消费者扩容策略

监控指标 → 决策 → 执行

队列堆积 1000+ msgs    ⟹  Lag > 处理能力 50%    ⟹  增加 2 个消费者
   或消费延迟 > 30s

队列堆积 < 100 msgs    ⟹  Lag 接近 0           ⟹  等待(可能减少)
   且消费速率稳定

消费者异常              ⟹  自动触发告警         ⟹  优雅降级 / 转移任务

削峰填谷的本质

4.3 削峰填谷的数学原理

流量平滑效果:

 原始流量(尖峰):                平滑后流量:

10万/s │    ╱╲                  1000/s │████████████████
       │   ╱  ╲                        │
       │  ╱    ╲                       │
 1000/s│╱        ╲                 0/s │
       └───────────────               └────────────────
       0s   1s   2s                   0s              20s

原始:10万/s 峰值,持续1秒
平滑:1000/s 恒定速率,持续100秒

关键公式:

 队列长度 = 生产者速率 × 持续时间 - 消费者速率 × 持续时间
        = 100,000 × 1 - 1,000 × 1
        = 99,000 条消息(峰值时队列堆积)

消费完所有消息所需时间 = 队列长度 / 消费者速率
                      = 99,000 / 1,000
                      = 99 秒

5. 核心问题三:如何保证消息不丢失、不重复、有序?

5.1 消息可靠性:三道防线

消息可能在三个环节丢失:生产者发送时、Broker存储时、消费者处理时。通过三道防线确保消息不丢失。

消息传递链中的三个丢失风险点

生产者 ──(防线1)──> Broker ──(防线2)──> 消费者
   |                 |                  |
   ↓                 ↓                  ↓
 未确认            内存丢失         未处理完成

三道防线概览表

防线 位置 风险 解决方案 成本
防线 1 生产者 → Broker 网络中断或Broker拒收 生产者等待 ACK,超时重试 低(应该总是做)
防线 2 Broker 存储 重启或宕机丢失 消息落盘 + 多副本同步 中(性能换可靠性)
防线 3 消费者处理 处理中崩溃未确认 手动 ACK,处理完才确认 低(应该总是做)

防线 1:生产者确认 (Producer ACK)

发送流程:

生产者                              Broker
   │                                │
   ├─► 发送消息 ─────────────────→ │
   │   (待确认状态)                 │ 接收消息
   │                                │ 入队存储
   │                      ◄─ ACK ──┤
   │                      (已确认)
   ├─ 更新本地状态 ──────────────────┤
   │  消息已发送

生产者确认配置选项

配置 说明 优点 缺点
Wait All 等待 Broker 完全确认(同步副本都写入) ✅ 最安全,消息绝不丢 ❌ 最慢,延迟高
Wait ACK 等待 Broker ACK 即可(消息入队即返回) ✅ 平衡可靠性和性能 ⚠️ 极端场景可能丢
Fire and Forget 发送即忘,不等待任何确认 ✅ 最快 ❌ 可能丢消息

最佳实践:使用 Wait ACK 模式,再加上重试机制。

// 生产者配置示例
KafkaProducer<String, String> producer = new KafkaProducer<>(
    Map.of(
        ProducerConfig.ACKS_CONFIG, "all",  // 等待所有副本确认
        ProducerConfig.RETRIES_CONFIG, 3,    // 失败重试 3 次
        ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 1  // 确保顺序
    )
);

// 发送消息
Future<RecordMetadata> future = producer.send(
    new ProducerRecord<>("orders", orderId, orderData),
    (metadata, exception) -> {
        if (exception != null) {
            // 失败处理:记录日志、重试或发送死信
            log.error("消息发送失败,订单ID: " + orderId);
        } else {
            log.info("消息已确认,分区: " + metadata.partition());
        }
    }
);

防线 2:Broker 持久化

Broker 收到消息后,需要立即持久化,而非仅存内存。

存储方式对比

方面 内存存储 磁盘存储
速度 ⚡ 极快(微秒) 🐢 较慢(毫秒)
容量 📦 受限(GB级) 💾 充足(TB级)
可靠性 ❌ 重启丢失 ✅ 持久保存
典型场景 Redis 缓存 关键消息
故障影响 灾难性 完全可恢复

Broker 持久化机制

消息到达 Broker
   ↓
1️⃣ 写入内存缓冲区
   ↓
2️⃣ 发送 ACK 给生产者(此时若 Broker 宕机,消息会丢)
   ↓
3️⃣ 批量刷盘(Write to Disk)— 定期将缓冲区数据落磁盘
   ↓
4️⃣ 消息持久化完成 ✅ 此后任何故障都不会丢

故障场景恢复

场景 1:Broker 宕机前,消息已落盘
 → ✅ 重启后从磁盘恢复,消息完整

场景 2:Broker 宕机前,消息仅在内存
 → ❌ 丢失(除非有副本在其他节点)

场景 3:使用多副本(Replication)
 → ✅ 消息同步到 3 个节点
 → 即使 1 个节点宕机,其他 2 个节点完整保留

Kafka Broker 持久化配置示例

# Broker 配置文件(server.properties)

# 1. 日志持久化到磁盘路径
log.dirs=/var/kafka/logs

# 2. 刷盘策略:10000 条消息后立即刷盘
log.flush.interval.messages=10000

# 3. 或按时间:5 秒后必须刷盘(二选一)
log.flush.interval.ms=5000

# 4. 副本配置:消息写入 3 个副本才算成功
min.insync.replicas=3

# 5. 日志保留策略:保留 7 天的数据
log.retention.hours=168

防线 3:消费者确认 (Consumer ACK)

消费者处理完消息后,必须向 Broker 确认,Broker 才会标记该消息为已消费。

消费流程对比

环节 自动 ACK 手动 ACK
1. 拉取消息 从 Broker 获取消息 从 Broker 获取消息
2. 处理消息 执行业务逻辑 执行业务逻辑
3. 故障处理 ❌ 消息已标记为已消费 ⚠️ 消息仍未确认
4. 确认 自动确认(处理完前) 手动确认(处理完后)
故障后重消费 ❌ 无法重消 ✅ 从未确认的地方重新开始

推荐方案:手动 ACK(最安全)

// Kafka 消费者配置:手动确认
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(
    Map.of(
        ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false,  // 禁用自动确认
        ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"
    )
);

consumer.subscribe(Arrays.asList("orders"));

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));
    
    for (ConsumerRecord<String, String> record : records) {
        try {
            // 1. 处理消息
            processOrder(record.value());
            
            // 2. 处理成功,手动提交偏移量
            consumer.commitSync();
            log.info("订单已处理并确认");
            
        } catch (Exception e) {
            // 3. 处理失败,不确认
            log.error("订单处理失败,消息将重新投递", e);
            // Broker 会在超时后重新投递该消息给其他消费者
        }
    }
}

三种 ACK 模式工作流程

┌──────────────────────────────────────────────────────────────┐
│ 自动 ACK (Auto Commit)                                        │
├──────────────────────────────────────────────────────────────┤
│ 拉取 → 处理中 → 自动提交偏移量 → 处理完成/处理失败            │
│                   ↑                                            │
│                   (处理完前就提交了!)                       │
│ 风险:处理中宕机,消息已标记已消费,丢失 ❌                   │
└──────────────────────────────────────────────────────────────┘

┌──────────────────────────────────────────────────────────────┐
│ 手动 ACK (Manual Commit) — 推荐                               │
├──────────────────────────────────────────────────────────────┤
│ 拉取 → 处理 → 处理成功 → 手动提交偏移量                      │
│              ↓(处理失败不提交)                              │
│          未提交,等待重投递                                   │
│ 优点:完全控制,不丢消息 ✅                                   │
└──────────────────────────────────────────────────────────────┘

┌──────────────────────────────────────────────────────────────┐
│ 定时 ACK (Periodic Commit)                                    │
├──────────────────────────────────────────────────────────────┤
│ 拉取 → 处理 → 处理中 → 5秒后自动提交 → 继续处理              │
│             (可能重复消费)                                  │
│ 风险:5秒内宕机,消息未提交但已标记,可能重复 ⚠️             │
└──────────────────────────────────────────────────────────────┘

消费者异常处理策略

正常流:
  拉取 → 处理 → 成功 → ACK

异常流:
  拉取 → 处理 → 失败 → 不 ACK
              ↓
         重试逻辑(可选)
              ↓
         失败超过 N 次 → 发送死信队列 (DLQ)
              ↓
         告警通知,人工处理

三道防线完整验证流程

✅ 防线 1 正常工作的标志:
   - 生产者日志显示"消息已发送确认"
   - 无" Send Timeout "异常

✅ 防线 2 正常工作的标志:
   - Broker 磁盘有持久化日志文件
   - Broker 重启后消息未丢失

✅ 防线 3 正常工作的标志:
   - 消费者日志显示"消息已手动确认"
   - 消费者宕机后重启,消息从断点处重新消费

三道防线,缺一不可:生产者确认 + Broker 持久化 + 消费者确认

5.2 如何处理消息重复消费?

消息重复可能在以下场景发生:

  1. 生产者重试 : 生产者发送消息后未收到ACK,重试发送同一条消息
  2. 消费者ACK超时 : 消费者处理完成但ACK超时,Broker重新投递
  3. 网络抖动 : 消费者ACK未到达Broker,Broker认为未消费
  4. 消费者重启 : 消费者重启后重新消费同一批消息

幂等性 : 同一操作执行多次和执行一次的效果相同。

生活中的幂等性 :

  • 幂等 : 按电梯按钮(按10次和按1次,电梯都会来)
  • 非幂等 : 转账(转10元,执行两次会转20元)

技术解决方案:为每条消息生成唯一 ID,处理前检查是否已处理过。


银行转账场景对比:无幂等性 vs 有幂等性

┌─────────────────────────────────────────────────────────────────┐
│ 场景:消息重复投递导致重复扣款                                    │
├─────────────────────────────────────────────────────────────────┤
│                                                                 │
│ 初始状态:发送方余额 ¥1000,接收方余额 ¥500                     │
│ 转账指令:转账 ¥100 给接收方                                     │
│                                                                 │
└─────────────────────────────────────────────────────────────────┘

无幂等保护(❌ 高风险)

环节 操作 状态 结果
第 1 次消费 收到转账消息:转账 ID = msg_001,金额 ¥100 未检查重复 ✅ 发送方 -¥100,接收方 +¥100
消费中宕机 处理完成但 ACK 超时 Broker 认为未消费 ⚠️ 消息重新投递
第 2 次消费 再次收到同一消息:msg_001,金额 ¥100 无幂等性检查 又扣一遍 — 发送方 -¥200,接收方 +¥200
最终余额 发送方:¥800(应为 ¥900) 接收方:¥600(应为 ¥500)

问题:重复消费导致多次扣款,账户金额错误,数据不一致!


有幂等保护(✅ 推荐)

环节 操作 处理日志检查 结果
第 1 次消费 收到转账消息:msg_001,金额 ¥100 查询日志:未处理过 ✅ 执行转账:发送方 -¥100,接收方 +¥100
记录日志 写入处理日志:msg_001 → 已处理 日志持久化成功 📝 标记:msg_001 in processed_messages
消费中宕机 处理完成但 ACK 超时 日志已落盘 ⚠️ 消息重新投递
第 2 次消费 再次收到同一消息:msg_001,金额 ¥100 查询日志:msg_001 已处理过 跳过处理,直接返回成功
最终余额 发送方:¥900(✅ 正确) 接收方:¥600(✅ 正确)

优势:即使消息重复投递,业务结果也完全正确!


幂等性实现方案

方案 1:数据库唯一键约束(最简单)

-- 创建处理日志表,msg_id 作为唯一键
CREATE TABLE processed_messages (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    msg_id VARCHAR(50) UNIQUE NOT NULL,    -- 消息唯一 ID
    business_type VARCHAR(50),              -- 业务类型(如"transfer")
    business_id VARCHAR(50),                -- 业务 ID(如订单号)
    processed_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);

-- 处理消息时的 SQL
INSERT IGNORE INTO processed_messages (msg_id, business_type, business_id)
VALUES ('msg_001', 'transfer', 'order_123');

-- 如果 msg_id 重复,INSERT 会被忽略,不会重复处理

问题:只能记录"已处理",无法记录处理结果。


方案 2:业务状态机(推荐生产环保)

// 转账服务
public class TransferService {
    
    // 幂等转账处理
    public TransferResult handleTransfer(TransferMessage msg) {
        String msgId = msg.getId();           // "msg_001"
        String transferId = msg.getTransferId(); // "order_123"
        
        // 1️⃣ 检查幂等性:查询该转账是否已处理
        TransferRecord record = db.query(
            "SELECT * FROM transfers WHERE transfer_id = ?", 
            transferId
        );
        
        if (record != null) {
            // 2️⃣ 已处理过,返回之前的结果(幂等)
            log.info("消息 {} 重复,已处理过,返回缓存结果", msgId);
            return TransferResult.from(record);
        }
        
        // 3️⃣ 未处理过,执行业务逻辑
        try {
            // 数据库事务:原子操作
            db.transaction(() -> {
                // 扣款
                db.update("UPDATE accounts SET balance = balance - ? WHERE id = ?",
                    msg.getAmount(), msg.getSenderId());
                
                // 加款
                db.update("UPDATE accounts SET balance = balance + ? WHERE id = ?",
                    msg.getAmount(), msg.getReceiverId());
                
                // 4️⃣ 记录转账记录(关键:使用转账 ID 作为主键)
                db.insert("INSERT INTO transfers (transfer_id, sender_id, receiver_id, amount, status) VALUES (?, ?, ?, ?, ?)",
                    transferId, msg.getSenderId(), msg.getReceiverId(), msg.getAmount(), "completed");
            });
            
            log.info("转账成功:{}", msgId);
            return TransferResult.success(transferId);
            
        } catch (Exception e) {
            log.error("转账失败:{}", msgId, e);
            return TransferResult.failure(transferId, e.getMessage());
        }
    }
}

关键点

  • 使用业务 ID(transfer_id)作为幂等键
  • 每次处理前查询该业务是否已存在
  • 业务数据本身有唯一键约束(主键),重复 INSERT 会报错

方案 3:Redis 缓存幂等性(高并发)

public class TransferService {
    
    private RedisCache cache;
    private static final String IDEMPOTENT_KEY_PREFIX = "transfer:";
    private static final long IDEMPOTENT_TTL = 86400; // 24 小时
    
    public TransferResult handleTransfer(TransferMessage msg) {
        String msgId = msg.getId();
        String transferId = msg.getTransferId();
        String idempotentKey = IDEMPOTENT_KEY_PREFIX + transferId;
        
        // 1️⃣ 查询 Redis 缓存(极快,O(1) 时间)
        String cachedResult = cache.get(idempotentKey);
        if (cachedResult != null) {
            log.info("消息 {} 重复,直接返回缓存结果", msgId);
            return TransferResult.fromJson(cachedResult);
        }
        
        // 2️⃣ 缓存未命中,执行业务逻辑
        try {
            TransferResult result = executeTransfer(msg);
            
            // 3️⃣ 将结果存入 Redis(幂等保证)
            cache.setex(idempotentKey, IDEMPOTENT_TTL, result.toJson());
            
            return result;
            
        } catch (Exception e) {
            // 4️⃣ 即使失败也要缓存失败结果(避免重试时重复扣款)
            TransferResult failResult = TransferResult.failure(transferId, e.getMessage());
            cache.setex(idempotentKey, IDEMPOTENT_TTL, failResult.toJson());
            throw e;
        }
    }
    
    private TransferResult executeTransfer(TransferMessage msg) {
        // 实际的业务逻辑...
    }
}

优势

  • 查询速度极快(Redis O(1) vs 数据库 O(n))
  • 适合高并发场景
  • 自动过期(TTL),不需要手动清理

三种方案对比

方案 实现复杂度 性能 可靠性 适用场景
唯一键约束 ⭐ 简单 ⭐⭐ 中等 ⭐⭐ 一般 低并发、传统应用
业务状态机 ⭐⭐⭐ 复杂 ⭐⭐ 中等 ⭐⭐⭐⭐⭐ 最高 生产环境首选
Redis 缓存 ⭐⭐ 中等 ⭐⭐⭐⭐⭐ 极快 ⭐⭐⭐⭐ 很高 高并发、实时系统

幂等性核心原则:为每条消息生成唯一 ID,处理前检查是否已处理,避免重复操作。最佳实践是业务状态机 + 唯一键约束的组合,兼具可靠性和性能。

6. 实战:如何选择消息队列?

6.1 四大主流消息队列对比

特性 RabbitMQ Kafka RocketMQ Redis Stream
定位 传统消息队列 分布式日志流 电商级消息队列 轻量级队列
吞吐量 ~1万/秒 ~100万/秒 ~10万/秒 ~5万/秒
延迟 微秒级 毫秒级 毫秒级 毫秒级
可靠性 高(持久化) 高(多副本) 高(同步刷盘) 中(AOF)
消息回溯 不支持 支持 支持 支持
事务消息 支持(弱) 不支持 支持(强) 不支持
延迟消息 支持 不支持 支持 不支持
适用场景 传统企业应用 日志、大数据 电商、金融 小规模应用

决策树:

 选择消息队列:
│
├─ 需要事务消息(分布式事务)?
│  ├─ 是 → RocketMQ(首选)或 RabbitMQ
│  └─ 否 → 继续
│
├─ 需要处理海量日志/实时流?
│  ├─ 是 → Kafka(首选)
│  └─ 否 → 继续
│
├─ QPS > 1万/秒?
│  ├─ 是 → RocketMQ 或 Kafka
│  └─ 否 → 继续
│
├─ 需要复杂路由(如 headers 匹配)?
│  ├─ 是 → RabbitMQ
│  └─ 否 → 继续
│
├─ 已有 Redis 基础设施?
│  ├─ 是 → Redis Stream(快速开始)
│  └─ 否 → RabbitMQ(功能全面,学习曲线适中)

7. 总结:消息队列设计心法

7.1 核心原则回顾

原则 含义 实践要点
解耦 服务间不直接依赖 通过消息队列通信,消费者故障不影响生产者
削峰 平滑流量波动 消息队列作为蓄水池,消费者按恒定速率处理
可靠 消息不丢失 生产者确认 + Broker持久化 + 消费者确认
幂等 重复消费无影响 业务层面保证幂等性(唯一键、状态机)
有序 消息顺序保证 单分区有序或消费者端排序

7.2 设计检查清单

在引入消息队列前,问自己以下问题:

  • 是否真的需要消息队列?(简单异步可以用线程池)
  • 消息丢失是否可以接受?(决定可靠性级别)
  • 消息重复是否会影响业务?(决定幂等性投入)
  • 消息顺序是否重要?(决定分区策略)
  • 消费者处理能力如何?(决定队列大小和告警阈值)
  • 如何处理消费失败?(决定重试和死信策略)

8. 名词速查表

名词 全称 解释
MQ Message Queue 消息队列 。用于异步通信的中间件,实现生产者和消费者的解耦。
Producer - 生产者 。发送消息的一方。
Consumer - 消费者 。接收并处理消息的一方。
Broker - 消息代理 。存储和转发消息的服务端程序。
Topic - 主题 。消息的逻辑分类(如 "orders")。
Queue - 队列 。存储消息的物理容器。
Partition - 分区 。Kafka的概念,一个Topic可以分成多个Partition,提升并发。
ACK Acknowledgment 确认 。消费者处理完消息后,向Broker确认。
Pub/Sub Publish/Subscribe 发布订阅 。一种消息模式,一条消息可被多个消费者接收。
P2P Point-to-Point 点对点 。一种消息模式,一条消息只能被一个消费者接收。
DLQ Dead Letter Queue 死信队列 。存放无法消费的消息。
Idempotence - 幂等性 。多次执行结果相同。
Throughput - 吞吐量 。单位时间内处理的消息数量。
Latency - 延迟 。消息从发送到被接收的时间差。
Persistence - 持久化 。消息写入磁盘,而非仅存内存。
Replication - 副本 。为了高可用,消息被复制到多个节点。
Transaction Message - 事务消息 。保证本地事务和消息发送的一致性。
Backpressure - 背压 。消费者处理不过来时,通知生产者降速。
Offset - 偏移量 。消费者在分区中的消费位置。
Rebalance - 重平衡 。消费者组成员变化时,重新分配分区。