Skip to content

消息队列

这里主要整理消息队列这条线,把下面几件事讲清楚:

  1. 为什么系统会需要消息队列
  2. 消息队列到底解决什么问题
  3. 消息为什么会丢、为什么会重复、为什么顺序难保证
  4. RabbitMQKafka 分别更适合什么场景

消息队列 不是某一个具体产品,而是一类中间件能力:系统把事件、任务或业务消息写进去,再由下游按约定方式消费处理。


1. 为什么需要消息队列

消息队列最常见的价值有三类:

  1. 解耦
  2. 异步
  3. 削峰

1.1 解耦

如果订单创建后还要同步调用库存、积分、短信、物流等多个系统,主流程会越来越重。

引入消息队列后,更常见的做法是:

  1. 主流程先完成核心事务
  2. 再把后续动作写成消息
  3. 下游系统各自订阅并处理

这样系统之间就不会强耦合在一条同步调用链上。

1.2 异步

很多动作不需要阻塞用户请求,例如:

  1. 发短信
  2. 发邮件
  3. 写审计日志
  4. 生成报表

把这些动作放进消息队列后,接口响应时间通常会更稳定。

1.3 削峰

削峰 可以理解成 把瞬时涌入的大量请求先缓存到队列里,再按消费者的处理能力慢慢消化

它解决的是高峰流量直接打爆下游系统的问题。

1.4 Java 线程池和消息队列有什么区别

很多人刚接触异步时,容易把 线程池消息队列 混在一起。它们确实都能把一部分工作从主线程、主请求里拆出去,但解决的问题不是一层:

  1. 线程池主要解决进程内的并发执行与资源复用
  2. 消息队列主要解决跨系统异步、解耦、削峰和可靠投递

直接看这张对比表:

对比维度Java 线程池消息队列
所在层级应用进程内应用进程外
主要目标复用线程、控制并发、调度任务解耦、异步、削峰、可靠传递
数据是否跨进程保存通常不会会,消息会先落到 Broker
调用关系更像本服务把任务交给本服务内部线程更像一个服务把事件交给另一个系统或消费者集群
服务重启后的影响内存队列里的任务可能丢失只要 Broker 做了持久化,消息通常还能继续消费
适合的任务范围发起并行调用、批处理、后台计算订单事件、通知链路、异步解耦、流量削峰
典型代价线程竞争、队列堆积、资源配置不当重复消费、顺序控制、消息积压、运维成本

1.4.1 为什么说线程池解决的是“进程内异步”

例如一个接口里需要:

  1. 并行调用 3 个下游服务
  2. 后台生成报表
  3. 异步执行一批图片处理任务

这种场景里,任务其实还是属于当前服务自己,只是希望:

  1. 不要阻塞主线程
  2. 不要每次都新建线程
  3. 能控制并发数量

这种场景通常就会用线程池。

线程池重点管理的是:

  1. 线程资源
  2. 任务排队
  3. 执行并发度

它本身并不天然解决:

  1. 跨系统解耦
  2. 消息可靠落盘
  3. 服务重启后的任务恢复

1.4.2 为什么说消息队列解决的是“跨系统异步”

例如订单创建后,后面还要做:

  1. 扣库存
  2. 发优惠券
  3. 发送短信
  4. 通知物流

如果全部靠当前服务的线程池异步执行,会有几个问题:

  1. 这些逻辑还是强依赖在当前服务里
  2. 当前服务挂了,内存里的待执行任务可能一起丢
  3. 下游处理速度跟不上时,当前服务自己要扛住堆积压力
  4. 很难自然扩展成多个独立消费者集群

消息队列更像这样:

  1. 当前服务把业务事件写入 Broker
  2. 不同下游系统各自订阅并消费
  3. 消费者可以独立扩缩容
  4. 消息失败后还能结合重试、死信队列、补偿机制继续处理

所以消息队列真正解决的,不只是“异步执行”,而是把系统间的依赖从同步调用改成基于消息的协作。

1.4.3 怎么选:什么时候用线程池,什么时候用消息队列

🌟 可以直接用一个很实用的判断方式:

  1. 任务只属于当前服务内部,而且结果不需要跨系统传递,优先考虑线程池
  2. 任务需要跨服务传播、解耦、削峰,或者希望和主流程彻底拆开,优先考虑消息队列

例如:

  1. 接口里并行查多个下游接口,适合线程池
  2. 订单创建后通知积分、营销、物流系统,适合消息队列
  3. 单机里的批量文件处理,通常先考虑线程池
  4. 秒杀高峰下把请求缓存下来慢慢消费,更适合消息队列

1.4.4 能不能一起用

可以,而且工程上很常见。

例如:

  1. 服务 A 把订单消息写入消息队列
  2. 服务 B 作为消费者拉取消息
  3. 服务 B 收到消息后,再把本地任务分发到线程池并发执行

这时候两者分工是:

  1. 消息队列负责系统之间的异步传递与削峰
  2. 线程池负责单个服务内部的任务并发执行

所以它们不是互相替代,而是经常上下配合。


2. 消息队列里最核心的几个角色

一条典型链路里通常会出现这些角色:

  1. Producer:生产者,负责发送消息
  2. Broker:消息中间件本身,负责存储和转发消息
  3. Consumer:消费者,负责读取并处理消息

这里的 Broker 直接看成:专门负责接收消息、保存消息、把消息分发给消费者的一层基础设施。


3. 消息队列真正难的地方是什么

很多人一开始会觉得 MQ 很简单,就是“把消息放进去,再拿出来”。但工程上真正麻烦的,通常是这几件事:

  1. 消息丢失
  2. 重复消费
  3. 顺序问题
  4. 消息积压
  5. 消费失败后的重试和补偿

3.1 为什么会丢消息

一条消息从发送到消费,至少会经过:

  1. 生产者发送
  2. Broker 接收
  3. Broker 持久化或复制
  4. 消费者处理
  5. 消费者确认

其中任何一个阶段出故障,都可能出现消息丢失风险。

3.1.1 如何尽量保证消息不丢失

先记一句最重要的话:消息队列很难承诺“绝对不丢”,工程上真正追求的是“尽量不丢、丢了能发现、发现后能补”。

如果把“防丢”拆开看,通常要同时守住 3 段链路:

  1. 生产者到 Broker
  2. Broker 自己的存储与复制
  3. Broker 到消费者,再到业务处理完成

看生产者到 Broker 这一段链路。

如果生产者把消息发出去了,但并不知道 Broker 是否真的收到了,这里就会有丢失风险。

所以常见做法通常包括:

  1. 发送确认机制
    例如 RabbitMQPublisher ConfirmKafkaacks=all
  2. 发送失败重试
    网络抖动或瞬时故障时,不要一失败就直接放弃
  3. 本地消息表 / Outbox
    把业务事件可靠写入本地数据库,再异步投递到消息中间件

再看 Broker 自己这一段。

如果 Broker 把消息只放在内存里,或者只有单副本,一旦进程崩溃、机器宕机、主节点损坏,就可能直接丢消息。

所以常见做法通常包括:

  1. 队列 / Topic 持久化
  2. 消息持久化
  3. 多副本同步
  4. 高可用部署
  5. 对磁盘、积压、主从切换做监控和告警

最后看消费者这一段。

消息到了消费者,不等于业务已经真正处理成功。

如果消费者刚收到消息就立刻确认,但后面的数据库写入、远程调用、业务落库失败了,这条消息在中间件视角里已经“处理完”,实际上业务却没有成功。

所以更稳妥的做法通常是:

  1. 业务处理成功后再确认消费
  2. 处理失败时不要错误确认,让消息进入重试或重投
  3. 消费逻辑做好幂等,避免重复投递带来副作用
  4. 对失败消息配置重试队列或死信队列

如果把这些措施压缩成一句工程化建议,可以记:发送要确认,Broker 要持久化和副本,消费要先处理后确认,业务要能幂等。

3.1.2 如果消息还是丢了,该怎么处理

真正成熟的系统不会把“消息不丢”只寄托在中间件本身,而是会同时准备:

  1. 发现机制
  2. 定位机制
  3. 补偿机制

先说发现。

消息丢失最怕的不是“丢了一条”,而是:系统已经丢了,但团队根本不知道。

所以常见做法通常包括:

  1. 发送成功率、消费成功率监控
  2. 消息积压、失败重试、死信数量监控
  3. 关键业务链路对账
  4. 必要时做消息轨迹或审计日志

再说定位。

丢消息时,首先要判断它到底丢在哪一段:

  1. 生产者根本没发成功
  2. Broker 收到了但没持久化成功
  3. Broker 存下来了,但消费者没正常拉到
  4. 消费者拉到了,但业务处理失败却错误确认了

这一步很重要,因为不同环节的补法完全不一样。

最后说补偿。

消息真的丢了之后,常见补偿方式包括:

  1. 业务侧重发
  2. 死信队列重新消费
  3. 定时补偿任务扫表重投
  4. 依靠本地消息表 / Outbox 重放
  5. 通过业务对账结果人工补单、补发、补通知

要注意的是:补偿不是临时救火,而应该是消息系统设计时就提前准备好的兜底手段。

所以工程上更稳的一套组合通常会是:

  1. 生产发送确认
  2. Broker 持久化与高可用
  3. 消费端手动确认
  4. 幂等处理
  5. 重试与死信队列
  6. 本地消息表 / Outbox
  7. 业务对账与补偿任务

3.2 为什么会重复消费

为了尽量避免消息丢失,很多消息系统都会更保守。

更常见的现实不是“绝不重复”,而是:尽量保证消息至少能送达一次,然后让业务自己做幂等。

幂等 说白了就是 同一个动作执行一次和执行多次,最终结果应该保持一致

3.3 为什么顺序难保证

消息队列通常很看重吞吐和扩展能力。
一旦并发处理、分区、重试都参与进来,“所有消息全局有序”就会非常昂贵。

所以工程上更常见的是:只保证某个业务 key 的局部顺序,而不是保证所有消息的全局顺序。


4. 应该先学产品,还是先学通用主线

更推荐的顺序是:

  1. 先理解通用问题
  2. 再看具体产品怎么实现

因为如果没有先理解“消息为什么会重复”“为什么要 ACK”“为什么要死信队列”,后面看 RabbitMQKafka 时,容易只剩术语记忆。


5. 常见消息队列产品怎么分工理解

5.1 RabbitMQ

RabbitMQ 更常被提到的特点是:

  1. 路由灵活
  2. 模型成熟
  3. 比较适合业务消息

它很适合先理解:

  1. Exchange
  2. Queue
  3. Binding
  4. Routing Key
  5. 死信队列

5.2 Kafka

Kafka 更常被提到的特点是:

  1. 吞吐高
  2. 更偏日志流
  3. 适合埋点、日志、流式处理

它的核心概念通常是:

  1. Topic
  2. Partition
  3. Offset
  4. Consumer Group

6. 当前目录下有哪些专题

如果是第一次系统看这块内容,可以按这个顺序往下看:

  1. 看这篇总览,建立“为什么需要 MQ”的整体认知
  2. 再看“两种中间件对比与选型”,建立第一层选型直觉
  3. 再看 RabbitMQ,理解路由模型
  4. 再看 Kafka,理解分区和日志流

7. 一句话总结

消息队列最值得先建立的,不是某个产品的命令细节,而是这条主线:

为什么要异步、为什么要解耦、为什么消息会丢和重复、顺序到底能保证到什么程度,以及不同产品分别擅长解决什么问题。

基于 VitePress 构建的个人技术笔记。