Appearance
消息队列
这里主要整理消息队列这条线,把下面几件事讲清楚:
- 为什么系统会需要消息队列
- 消息队列到底解决什么问题
- 消息为什么会丢、为什么会重复、为什么顺序难保证
RabbitMQ、Kafka分别更适合什么场景
消息队列 不是某一个具体产品,而是一类中间件能力:系统把事件、任务或业务消息写进去,再由下游按约定方式消费处理。
1. 为什么需要消息队列
消息队列最常见的价值有三类:
- 解耦
- 异步
- 削峰
1.1 解耦
如果订单创建后还要同步调用库存、积分、短信、物流等多个系统,主流程会越来越重。
引入消息队列后,更常见的做法是:
- 主流程先完成核心事务
- 再把后续动作写成消息
- 下游系统各自订阅并处理
这样系统之间就不会强耦合在一条同步调用链上。
1.2 异步
很多动作不需要阻塞用户请求,例如:
- 发短信
- 发邮件
- 写审计日志
- 生成报表
把这些动作放进消息队列后,接口响应时间通常会更稳定。
1.3 削峰
削峰 可以理解成 把瞬时涌入的大量请求先缓存到队列里,再按消费者的处理能力慢慢消化。
它解决的是高峰流量直接打爆下游系统的问题。
1.4 Java 线程池和消息队列有什么区别
很多人刚接触异步时,容易把 线程池 和 消息队列 混在一起。它们确实都能把一部分工作从主线程、主请求里拆出去,但解决的问题不是一层:
- 线程池主要解决进程内的并发执行与资源复用
- 消息队列主要解决跨系统异步、解耦、削峰和可靠投递
直接看这张对比表:
| 对比维度 | Java 线程池 | 消息队列 |
|---|---|---|
| 所在层级 | 应用进程内 | 应用进程外 |
| 主要目标 | 复用线程、控制并发、调度任务 | 解耦、异步、削峰、可靠传递 |
| 数据是否跨进程保存 | 通常不会 | 会,消息会先落到 Broker |
| 调用关系 | 更像本服务把任务交给本服务内部线程 | 更像一个服务把事件交给另一个系统或消费者集群 |
| 服务重启后的影响 | 内存队列里的任务可能丢失 | 只要 Broker 做了持久化,消息通常还能继续消费 |
| 适合的任务范围 | 发起并行调用、批处理、后台计算 | 订单事件、通知链路、异步解耦、流量削峰 |
| 典型代价 | 线程竞争、队列堆积、资源配置不当 | 重复消费、顺序控制、消息积压、运维成本 |
1.4.1 为什么说线程池解决的是“进程内异步”
例如一个接口里需要:
- 并行调用 3 个下游服务
- 后台生成报表
- 异步执行一批图片处理任务
这种场景里,任务其实还是属于当前服务自己,只是希望:
- 不要阻塞主线程
- 不要每次都新建线程
- 能控制并发数量
这种场景通常就会用线程池。
线程池重点管理的是:
- 线程资源
- 任务排队
- 执行并发度
它本身并不天然解决:
- 跨系统解耦
- 消息可靠落盘
- 服务重启后的任务恢复
1.4.2 为什么说消息队列解决的是“跨系统异步”
例如订单创建后,后面还要做:
- 扣库存
- 发优惠券
- 发送短信
- 通知物流
如果全部靠当前服务的线程池异步执行,会有几个问题:
- 这些逻辑还是强依赖在当前服务里
- 当前服务挂了,内存里的待执行任务可能一起丢
- 下游处理速度跟不上时,当前服务自己要扛住堆积压力
- 很难自然扩展成多个独立消费者集群
消息队列更像这样:
- 当前服务把业务事件写入 Broker
- 不同下游系统各自订阅并消费
- 消费者可以独立扩缩容
- 消息失败后还能结合重试、死信队列、补偿机制继续处理
所以消息队列真正解决的,不只是“异步执行”,而是把系统间的依赖从同步调用改成基于消息的协作。
1.4.3 怎么选:什么时候用线程池,什么时候用消息队列
🌟 可以直接用一个很实用的判断方式:
- 任务只属于当前服务内部,而且结果不需要跨系统传递,优先考虑线程池
- 任务需要跨服务传播、解耦、削峰,或者希望和主流程彻底拆开,优先考虑消息队列
例如:
- 接口里并行查多个下游接口,适合线程池
- 订单创建后通知积分、营销、物流系统,适合消息队列
- 单机里的批量文件处理,通常先考虑线程池
- 秒杀高峰下把请求缓存下来慢慢消费,更适合消息队列
1.4.4 能不能一起用
可以,而且工程上很常见。
例如:
- 服务 A 把订单消息写入消息队列
- 服务 B 作为消费者拉取消息
- 服务 B 收到消息后,再把本地任务分发到线程池并发执行
这时候两者分工是:
- 消息队列负责系统之间的异步传递与削峰
- 线程池负责单个服务内部的任务并发执行
所以它们不是互相替代,而是经常上下配合。
2. 消息队列里最核心的几个角色
一条典型链路里通常会出现这些角色:
Producer:生产者,负责发送消息Broker:消息中间件本身,负责存储和转发消息Consumer:消费者,负责读取并处理消息
这里的 Broker 直接看成:专门负责接收消息、保存消息、把消息分发给消费者的一层基础设施。
3. 消息队列真正难的地方是什么
很多人一开始会觉得 MQ 很简单,就是“把消息放进去,再拿出来”。但工程上真正麻烦的,通常是这几件事:
- 消息丢失
- 重复消费
- 顺序问题
- 消息积压
- 消费失败后的重试和补偿
3.1 为什么会丢消息
一条消息从发送到消费,至少会经过:
- 生产者发送
- Broker 接收
- Broker 持久化或复制
- 消费者处理
- 消费者确认
其中任何一个阶段出故障,都可能出现消息丢失风险。
3.1.1 如何尽量保证消息不丢失
先记一句最重要的话:消息队列很难承诺“绝对不丢”,工程上真正追求的是“尽量不丢、丢了能发现、发现后能补”。
如果把“防丢”拆开看,通常要同时守住 3 段链路:
- 生产者到 Broker
- Broker 自己的存储与复制
- Broker 到消费者,再到业务处理完成
看生产者到 Broker 这一段链路。
如果生产者把消息发出去了,但并不知道 Broker 是否真的收到了,这里就会有丢失风险。
所以常见做法通常包括:
- 发送确认机制
例如RabbitMQ的Publisher Confirm、Kafka的acks=all - 发送失败重试
网络抖动或瞬时故障时,不要一失败就直接放弃 - 本地消息表 / Outbox
把业务事件可靠写入本地数据库,再异步投递到消息中间件
再看 Broker 自己这一段。
如果 Broker 把消息只放在内存里,或者只有单副本,一旦进程崩溃、机器宕机、主节点损坏,就可能直接丢消息。
所以常见做法通常包括:
- 队列 / Topic 持久化
- 消息持久化
- 多副本同步
- 高可用部署
- 对磁盘、积压、主从切换做监控和告警
最后看消费者这一段。
消息到了消费者,不等于业务已经真正处理成功。
如果消费者刚收到消息就立刻确认,但后面的数据库写入、远程调用、业务落库失败了,这条消息在中间件视角里已经“处理完”,实际上业务却没有成功。
所以更稳妥的做法通常是:
- 业务处理成功后再确认消费
- 处理失败时不要错误确认,让消息进入重试或重投
- 消费逻辑做好幂等,避免重复投递带来副作用
- 对失败消息配置重试队列或死信队列
如果把这些措施压缩成一句工程化建议,可以记:发送要确认,Broker 要持久化和副本,消费要先处理后确认,业务要能幂等。
3.1.2 如果消息还是丢了,该怎么处理
真正成熟的系统不会把“消息不丢”只寄托在中间件本身,而是会同时准备:
- 发现机制
- 定位机制
- 补偿机制
先说发现。
消息丢失最怕的不是“丢了一条”,而是:系统已经丢了,但团队根本不知道。
所以常见做法通常包括:
- 发送成功率、消费成功率监控
- 消息积压、失败重试、死信数量监控
- 关键业务链路对账
- 必要时做消息轨迹或审计日志
再说定位。
丢消息时,首先要判断它到底丢在哪一段:
- 生产者根本没发成功
- Broker 收到了但没持久化成功
- Broker 存下来了,但消费者没正常拉到
- 消费者拉到了,但业务处理失败却错误确认了
这一步很重要,因为不同环节的补法完全不一样。
最后说补偿。
消息真的丢了之后,常见补偿方式包括:
- 业务侧重发
- 死信队列重新消费
- 定时补偿任务扫表重投
- 依靠本地消息表 / Outbox 重放
- 通过业务对账结果人工补单、补发、补通知
要注意的是:补偿不是临时救火,而应该是消息系统设计时就提前准备好的兜底手段。
所以工程上更稳的一套组合通常会是:
- 生产发送确认
- Broker 持久化与高可用
- 消费端手动确认
- 幂等处理
- 重试与死信队列
- 本地消息表 / Outbox
- 业务对账与补偿任务
3.2 为什么会重复消费
为了尽量避免消息丢失,很多消息系统都会更保守。
更常见的现实不是“绝不重复”,而是:尽量保证消息至少能送达一次,然后让业务自己做幂等。
幂等 说白了就是 同一个动作执行一次和执行多次,最终结果应该保持一致。
3.3 为什么顺序难保证
消息队列通常很看重吞吐和扩展能力。
一旦并发处理、分区、重试都参与进来,“所有消息全局有序”就会非常昂贵。
所以工程上更常见的是:只保证某个业务 key 的局部顺序,而不是保证所有消息的全局顺序。
4. 应该先学产品,还是先学通用主线
更推荐的顺序是:
- 先理解通用问题
- 再看具体产品怎么实现
因为如果没有先理解“消息为什么会重复”“为什么要 ACK”“为什么要死信队列”,后面看 RabbitMQ、Kafka 时,容易只剩术语记忆。
5. 常见消息队列产品怎么分工理解
5.1 RabbitMQ
RabbitMQ 更常被提到的特点是:
- 路由灵活
- 模型成熟
- 比较适合业务消息
它很适合先理解:
ExchangeQueueBindingRouting Key- 死信队列
5.2 Kafka
Kafka 更常被提到的特点是:
- 吞吐高
- 更偏日志流
- 适合埋点、日志、流式处理
它的核心概念通常是:
TopicPartitionOffsetConsumer Group
6. 当前目录下有哪些专题
如果是第一次系统看这块内容,可以按这个顺序往下看:
- 看这篇总览,建立“为什么需要 MQ”的整体认知
- 再看“两种中间件对比与选型”,建立第一层选型直觉
- 再看
RabbitMQ,理解路由模型 - 再看
Kafka,理解分区和日志流
7. 一句话总结
消息队列最值得先建立的,不是某个产品的命令细节,而是这条主线:
为什么要异步、为什么要解耦、为什么消息会丢和重复、顺序到底能保证到什么程度,以及不同产品分别擅长解决什么问题。