消息队列到底解决了什么:从异步、削峰到消息可靠性
在系统架构设计中,消息队列经常与“异步、削峰、解耦”同时出现。
例如,用户提交订单后,系统还需要:
- 扣减库存;
- 发送通知;
- 增加积分;
- 记录审计日志;
- 更新统计数据;
- 触发配送流程。
如果所有步骤都在一个同步请求中依次完成,任何一个下游服务变慢,用户都需要继续等待。
引入消息队列后,订单服务可以先完成核心操作,再发送一条“订单已创建”消息,由其他服务分别处理后续工作。
不过,消息队列并不会自动让系统可靠。它还会带来消息丢失、重复消费、顺序错乱、队列积压和数据一致性等新问题。
一、消息队列的基本结构
一套消息系统通常包含三类角色:
生产者
→ 消息代理
→ 消费者
生产者
生产者负责创建和发送消息。
例如,订单服务创建订单后发送:
{
"event_id": "evt_202609210001",
"event_type": "order_created",
"order_id": "order_90021",
"user_id": "user_1024",
"created_at": "2026-09-21T10:30:00Z"
}
消息代理
消息代理负责接收、存储和转发消息。
它通常需要处理:
- 消息持久化;
- 消息分区;
- 消费进度;
- 重试;
- 副本;
- 故障恢复。
消费者
消费者从队列中获取消息并执行具体业务。
例如:
库存服务:预留库存
积分服务:增加积分
通知服务:发送消息
统计服务:更新报表
同一业务事件可以被多个不同消费者独立处理。
二、异步处理为什么能缩短响应时间
同步处理流程可能是:
创建订单:200毫秒
扣减库存:300毫秒
增加积分:150毫秒
发送通知:500毫秒
更新统计:200毫秒
如果串行执行,总耗时约为:
200 + 300 + 150 + 500 + 200 = 1350毫秒
使用消息队列后,订单接口可以只完成核心操作:
创建订单
→ 写入消息
→ 返回结果
积分、通知和统计等操作由后台消费者异步执行。
这样可以降低用户请求的等待时间,但要注意:
异步只是把工作移到了请求之后,并没有让这些工作消失。
如果消费者处理能力不足,任务仍然会在队列中积压。
三、消息队列如何实现解耦
没有消息队列时,订单服务可能直接调用多个下游:
订单服务
├── 调用库存服务
├── 调用积分服务
├── 调用通知服务
└── 调用统计服务
订单服务需要知道每个下游的调用方式、超时时间和错误处理策略。
如果增加一个新的风控服务,订单服务也需要修改代码。
使用事件后,订单服务只表达一个事实:
订单已经创建
哪些系统关心这个事件,由消费者自行决定。
这样可以降低服务之间的直接依赖:
订单服务
→ 订单已创建事件
├── 库存消费者
├── 积分消费者
├── 通知消费者
└── 统计消费者
不过,解耦不等于没有依赖。生产者和消费者仍然通过消息结构、字段含义和业务语义形成契约。
如果生产者随意删除字段或改变含义,消费者仍然可能出现故障。
四、消息队列如何实现削峰
假设某个系统平时每秒收到100个请求,活动开始时突然增加到每秒5000个。
数据库只能稳定处理每秒1000个任务。如果所有请求直接进入数据库,数据库可能迅速过载。
消息队列可以先接收任务:
每秒进入5000条
→ 队列暂存
→ 消费者每秒处理1000条
短时间流量高峰被转化为一段更长、相对平稳的处理过程。
这就是削峰。
但消息队列不能突破系统总处理能力。
如果生产速度长期高于消费速度:
生产速度:每秒5000条
消费速度:每秒1000条
那么每秒都会积压4000条。只要这种状态持续,队列最终仍然会被填满,或者任务延迟会越来越高。
因此,削峰适合处理短时间突发流量,不能掩盖长期容量不足。
五、队列与发布订阅有什么区别
工作队列
一条任务通常只需要被某个消费者实例处理一次。
例如:
生成一份报告
压缩一个文件
处理一张图片
多个消费者可以共同分担队列中的任务。
发布订阅
同一条事件需要被多个不同业务分别处理。
例如,订单创建事件可能同时触发:
- 库存处理;
- 用户通知;
- 积分计算;
- 经营统计。
这些消费者不是互相竞争同一份任务,而是分别获得该事件。
具体消息系统对“队列、主题、消费组、订阅”的定义可能不同,但核心问题都是:
一条消息应该由谁接收,以及每类消费者需要处理几次。
六、消息投递的三种常见语义
最多一次
消息最多被处理一次,但可能丢失。
典型流程是:
先确认消息
→ 再处理业务
如果确认后消费者立即崩溃,消息已经被认为完成,但业务并未执行。
这种模式适合允许少量丢失、但不希望重复处理的非关键场景。
至少一次
消息保证被尝试处理,但可能重复。
典型流程是:
先处理业务
→ 成功后确认消息
如果业务已经完成,但消费者在发送确认前崩溃,消息系统会重新投递,于是同一消息可能再次处理。
这是实际系统中非常常见的语义。
精确一次
理想状态是每条消息只产生一次业务效果。
但在跨网络、跨数据库和跨服务场景中,实现真正的端到端精确一次非常困难。
某些消息系统可以在自己的边界内提供事务或去重能力,但如果消费者还要写数据库、调用外部接口,就仍然需要考虑重复执行。
工程上更常见的目标是:
消息允许重复投递
+ 消费者幂等处理
= 业务效果只发生一次
这也常被称为“效果上的精确一次”。
七、消息为什么会丢失
消息可能在多个阶段丢失。
1. 生产者发送失败
生产者认为消息已经发送,但消息代理实际上没有收到。
可能原因包括:
- 网络中断;
- 请求超时;
- 消息代理不可用;
- 发送缓冲区未及时刷新;
- 应用进程突然退出。
生产者需要使用发送确认机制,不能只调用一次发送函数就默认成功。
2. 消息代理故障
消息到达代理后,如果只保存在内存中,代理崩溃可能导致消息丢失。
关键消息通常需要:
- 持久化存储;
- 多副本;
- 写入确认;
- 故障恢复机制。
3. 消费者提前确认
消费者收到消息后立即确认,再执行业务。如果业务处理过程中崩溃,消息不会重新投递。
对于关键任务,应在业务成功后再确认消息。
4. 错误处理直接丢弃
消费者解析失败或业务异常后,如果既不重试,也不保存失败消息,该消息就可能永久消失。
八、消息为什么会重复
即使生产者和消息代理都正常,重复消息仍然无法完全避免。
假设消费者执行:
写入数据库成功
→ 准备确认消息
→ 进程突然崩溃
消息代理没有收到确认,于是重新投递消息。
但数据库写入已经生效,因此消费者会再次执行相同操作。
网络超时也会产生类似问题:
生产者发送消息
→ 消息代理成功保存
→ 确认响应在网络中丢失
→ 生产者认为失败并重新发送
从生产者角度看,它无法确定第一次发送是否成功,只能重试。
因此,重复不是消息系统的罕见异常,而是可靠消息机制必须面对的正常情况。
九、消费者如何实现幂等
幂等表示同一消息处理一次或多次,最终业务结果相同。
方法一:使用唯一事件编号
每条事件拥有唯一event_id。
消费者处理前检查该事件是否已经完成:
未处理
→ 执行业务
→ 记录已处理
已处理
→ 直接返回成功
去重记录和业务修改最好处在同一个数据库事务中,否则可能出现:
业务成功
→ 去重记录写入失败
下次重试仍会重复执行。
方法二:使用数据库唯一约束
例如,一个订单只能获得一次积分,可以建立唯一约束:
UNIQUE(order_id, reward_type)
即使消息重复消费,数据库也只允许第一次插入成功。
方法三:使用状态机
业务对象只能按照合法状态转移:
pending
→ paid
→ shipped
→ completed
如果对象已经处于paid状态,再次收到相同支付成功事件时,可以直接忽略。
方法四:使用条件更新
例如:
UPDATE orders
SET status = 'paid'
WHERE order_id = ?
AND status = 'pending';
根据受影响行数判断状态是否首次改变。
十、为什么“先检查再处理”仍可能重复
下面的代码看似实现了去重:
查询event_id是否存在
→ 不存在
→ 执行业务
→ 写入event_id
两个消费者可能同时执行:
消费者A查询:不存在
消费者B查询:不存在
消费者A执行业务
消费者B执行业务
因此,去重必须依赖原子机制,例如:
- 数据库唯一约束;
- 同一事务中的条件写入;
- 原子状态更新;
- 可靠锁机制。
单纯在应用代码中先查询再判断,无法可靠应对并发。
十一、消息为什么会乱序
假设同一个订单连续产生两条消息:
消息1:订单已支付
消息2:订单已取消
如果它们被发送到不同分区,或者由不同消费者并行处理,可能出现:
先处理取消
后处理支付
最终状态就可能错误。
消息乱序可能来自:
- 多个生产者并发发送;
- 多个分区分别存储;
- 多个消费者并行处理;
- 某条消息处理失败后重试;
- 网络延迟不同;
- 批量发送和异步确认。
十二、如何保证业务顺序
全局顺序通常代价很高,因为它会限制并行能力。
实际系统更常见的是保证同一业务对象内部有序。
例如,将相同订单编号作为分区键:
order_1001的消息
→ 始终进入同一分区
→ 由同一消费顺序处理
不同订单仍然可以并行:
order_1001 → 分区A
order_1002 → 分区B
order_1003 → 分区C
消费者还可以检查事件版本:
订单当前版本:5
收到事件版本:4
→ 旧事件,忽略
收到事件版本:6
→ 按规则处理
需要注意,失败重试可能打乱后续消息处理。如果某条消息必须成功后才能继续,消费者就不能简单跳过它处理下一条。
十三、重试不是越多越好
消息处理失败后,系统通常会重试。
但如果失败原因是永久性的,例如:
- 消息格式错误;
- 必填字段缺失;
- 业务对象不存在;
- 数据违反约束;
- 代码存在固定缺陷;
无论重试多少次,结果都不会改变。
无限重试会导致:
- 消费线程被长期占用;
- 后续消息无法处理;
- 下游服务被重复请求;
- 日志大量增长;
- 队列持续积压。
合理的重试策略通常包括:
- 区分临时错误和永久错误;
- 限制最大重试次数;
- 使用指数退避;
- 加入随机抖动;
- 记录最后失败原因;
- 超过次数后转入死信队列。
十四、什么是死信队列
无法正常处理的消息可以转移到专门的死信队列。
常见进入原因包括:
- 重试次数耗尽;
- 消息过期;
- 格式无法解析;
- 业务规则拒绝;
- 消费者主动拒绝。
死信队列不是垃圾桶。
系统需要对其进行:
- 告警;
- 分类;
- 原因分析;
- 人工修复;
- 修改后重新投递;
- 过期清理。
如果消息只是被移动到死信队列,却没有任何人处理,业务仍然处于失败状态。
十五、什么是毒性消息
毒性消息是指每次被消费都会导致相同错误的消息。
例如,它可能包含:
- 程序无法解析的字段;
- 超出范围的数据;
- 触发代码缺陷的特殊内容;
- 无法满足业务约束的状态。
如果消息系统始终重新投递这条消息,它可能阻塞整个分区或消费队列。
处理毒性消息需要:
- 限制重试;
- 记录原始消息;
- 隔离到死信队列;
- 保留错误上下文;
- 避免影响后续正常消息。
十六、消息积压是怎样形成的
当生产速度大于消费速度时,队列就会积压。
设:
生产速度 = P
消费速度 = C
当:
P > C
每秒积压量约为:
P - C
例如:
生产速度:每秒2000条
消费速度:每秒1500条
则每秒积压约500条。
持续一小时后,理论上会新增:
500 × 3600 = 1800000条
因此,积压不是单纯看队列中有多少消息,还要看:
- 积压增长速度;
- 最老消息等待时间;
- 消费者处理速度;
- 消费失败率;
- 每条消息平均耗时。
十七、如何处理消息积压
增加消费者
如果不同消息可以并行处理,可以增加消费者实例。
但消费者数量受以下因素限制:
- 分区数量;
- 数据库容量;
- 下游接口限流;
- CPU与内存;
- 网络带宽;
- 锁竞争。
增加消费者可能只是把压力转移到数据库或下游服务。
提高单条处理效率
可以通过以下方式优化:
- 批量读取;
- 批量写入;
- 减少外部调用;
- 优化SQL;
- 使用连接池;
- 避免重复计算;
- 缩短事务;
- 并行处理独立步骤。
暂停或限制生产者
如果下游已经过载,应通过背压、限流或降级减少新任务进入。
区分优先级
关键任务与非关键任务可以进入不同队列,避免低优先级任务淹没核心业务。
临时扩容后再缩容
突发积压时可以增加消费者,但需要监控下游是否能够承受扩容后的并发量。
十八、数据库更新与发送消息的一致性问题
假设订单服务执行:
写入订单数据库
→ 发送订单创建消息
可能出现两种错误。
数据库成功,消息发送失败
订单已经存在,但库存、积分和通知服务不知道它已经创建。
消息发送成功,数据库事务回滚
消费者收到“订单已创建”消息,但数据库中实际上没有这个订单。
这说明:
数据库事务
和:
消息发送
属于两个独立系统,不能通过普通本地事务自动保证同时成功。
十九、Outbox模式如何解决双写问题
Outbox模式的核心思想是:在同一个数据库事务中,同时写入业务数据和待发送事件。
例如:
开始数据库事务
→ 写入订单表
→ 写入outbox事件表
→ 提交事务
因为两次写入位于同一个本地事务中,所以会一起成功或一起失败。
后台程序再读取outbox表,将事件发送到消息队列:
读取未发送事件
→ 发送消息
→ 标记已发送
即使后台程序在发送后、标记前崩溃,也可能造成重复发送,但不会轻易丢失消息。
消费者仍然需要实现幂等。
Outbox模式把难以保证的“数据库与消息队列双写一致性”,转换为:
数据库本地事务
+ 可重试消息发送
+ 消费者幂等
二十、消费者也可以使用Inbox模式
消费者可以维护一个已处理事件表,也称为Inbox。
处理消息时:
开始数据库事务
→ 检查event_id
→ 执行业务修改
→ 记录event_id
→ 提交事务
→ 确认消息
唯一约束可以防止相同事件重复生效。
Outbox负责减少生产端消息丢失,Inbox负责降低消费端重复处理风险,两者可以配合使用。
二十一、消息结构为什么需要版本管理
消息一旦发布,就可能被多个消费者长期使用。
假设原消息为:
{
"order_id": "90021",
"amount": 199
}
后续生产者直接把amount改成字符串,或者删除order_id,旧消费者可能立即失败。
更安全的演进方式包括:
- 新增字段时提供默认语义;
- 尽量不删除旧字段;
- 不随意改变字段类型;
- 明确消息版本;
- 同时兼容新旧格式;
- 消费者忽略不认识的可选字段;
- 在淘汰旧版本前确认所有消费者已升级。
消息是服务之间的长期契约,而不只是一次临时参数传递。
二十二、消息中应该包含什么
一条业务事件通常可以包含:
- 唯一事件编号;
- 事件类型;
- 事件版本;
- 业务对象编号;
- 事件发生时间;
- 生产者信息;
- 链路标识;
- 必要业务字段。
例如:
{
"event_id": "evt_202609210001",
"event_type": "order_paid",
"event_version": 2,
"order_id": "order_90021",
"occurred_at": "2026-09-21T10:30:00Z",
"trace_id": "trace_71af",
"data": {
"payment_amount": 199,
"currency": "CNY"
}
}
不要在消息中放入密码、完整凭据等敏感信息。
消息也不宜无限扩大。过大的消息会增加网络、内存、存储和重试成本。大型文件通常更适合保存到专门存储中,消息只携带文件标识和必要元数据。
二十三、事件与命令有什么区别
事件
事件描述已经发生的事实:
订单已创建
支付已完成
文件已上传
事件通常使用过去式,生产者只陈述事实,不指定必须由哪个具体服务处理。
命令
命令表达希望某个处理者执行的操作:
发送订单通知
生成月度报告
冻结用户账户
事件可能被多个订阅者处理,命令通常面向某类明确的执行者。
区分二者有助于明确消息语义,避免所有消息都使用含糊的“处理任务”名称。
二十四、消息队列不适合什么场景
必须立即得到结果
如果后续步骤的结果决定当前请求能否成功,简单异步化可能无法满足要求。
极低延迟的同步查询
读取用户资料或商品详情时,直接查询缓存或数据库通常比经过消息队列更合适。
任务量很小且系统简单
为少量后台任务引入完整消息系统,可能增加不必要的部署和维护成本。
无法接受最终一致性
消息处理通常存在延迟。如果业务要求所有系统在同一时刻立即一致,需要使用更合适的事务或协调机制。
没有能力处理重复与积压
如果消费者不具备幂等性,也没有监控和死信处理流程,引入消息队列可能让故障更难发现。
二十五、消息系统需要监控什么
至少应关注:
- 生产消息速率;
- 消费消息速率;
- 队列积压数量;
- 最老消息等待时间;
- 消费延迟;
- 消费成功率;
- 重试次数;
- 死信数量;
- 消费者在线数量;
- 消息代理磁盘使用量;
- 消息代理内存使用量;
- 分区负载是否均衡;
- 生产和消费错误率;
- 消费者处理耗时;
- 消费进度是否长时间不动。
只监控“消息代理进程是否存活”远远不够。进程正常运行时,消息也可能已经积压数小时。
二十六、常见错误做法
1. 认为消息绝不会重复
网络超时、消费者崩溃和确认丢失都可能造成重复投递。
2. 消费成功前就确认消息
消费者崩溃后,未完成业务可能无法重新处理。
3. 失败后无限重试
永久性错误会持续占用消费者并阻塞正常消息。
4. 把死信队列当垃圾桶
没有告警和修复流程的死信队列,只是隐藏业务失败。
5. 认为增加消费者一定能解决积压
下游数据库或接口可能先被压垮。
6. 假设消息全局有序
多数高吞吐消息系统只能在一定范围内保证顺序。
7. 数据库写入与消息发送直接双写
任意一步失败都可能造成状态不一致。
8. 随意修改消息字段
消费者可能没有同步升级,导致大量消息处理失败。
9. 队列没有容量与保留限制
积压可能耗尽磁盘,并让任务延迟无限增长。
二十七、消息队列设计检查清单
引入消息队列前,可以检查以下问题:
- 业务真正需要异步、削峰还是发布订阅;
- 数据是否允许最终一致;
- 每条消息是否具有唯一编号;
- 生产者如何确认发送成功;
- 消息是否持久化和复制;
- 消费者何时确认消息;
- 消费者是否具备幂等性;
- 重复消息如何识别;
- 是否需要同一业务对象内有序;
- 分区键如何选择;
- 哪些错误允许重试;
- 最大重试次数是多少;
- 是否使用退避和随机抖动;
- 死信消息由谁处理;
- 队列积压到什么程度需要告警;
- 消费速度不足时如何扩容或限流;
- 数据库和消息发送如何保持一致;
- 消息格式如何进行版本管理;
- 服务重启后未完成任务如何恢复;
- 消息中是否包含敏感数据。
结语
消息队列主要解决三个问题:
异步:把非核心工作移到请求之后
削峰:把短时间高流量转换为平稳处理
解耦:让生产者不必直接调用所有消费者
但它同时引入了新的工程问题:
消息可能丢失
消息可能重复
消息可能乱序
消息可能积压
多个系统可能暂时不一致
可靠的消息系统并不是保证每条消息永远只出现一次,而是通过生产确认、持久化、有限重试、死信处理、消费者幂等和状态监控,让消息即使重复、延迟或暂时失败,最终仍能产生正确的业务结果。
引入消息队列的目的不是让架构看起来更复杂,而是用可控的异步过程,换取系统在高流量和局部故障下更稳定的表现。
- 点赞
- 收藏
- 关注作者
评论(0)