分布式与消息队列场景
开篇:消息队列在生产环境最容易踩的坑
消息队列是分布式系统的血管。用得好,系统解耦、异步高效、削峰填谷。用不好,消息丢了不知道、重复消费造成脏数据、积压百万条把下游拖垮。
这篇文章从实战角度出发,把 MQ 在生产环境最容易遇到的几个场景拉出来过一遍:消息丢了怎么排查、重复消费怎么做幂等、积压了怎么快速消化、Kafka 吞吐量怎么优化。每个场景都带真实方案。
一、消息丢失场景与排查
消息从生产到消费,要经过三个环节:生产者发送、Broker 存储、消费者消费。每一个环节都可能丢消息。
生产者端丢失
生产者发送消息时,网络抖动导致发送失败,或者发送成功但没收到 Broker 的 ACK 就认为失败了。
解决方案: 开启同步发送 + 重试机制。Kafka 设置 acks=all 确保所有副本都写入成功才返回 ACK。RocketMQ 用同步发送模式(send 而非 sendOneway),发送失败自动重试。
如果重试也不放心,引入本地消息表 -- 消息先落库再发送,发送成功更新状态,定时任务扫描未发送成功的消息做重试。这是最可靠的方案,代价是多了一张表和一个定时任务。
本地消息表的核心字段:
| 字段 | 说明 |
|---|---|
| message_key | 业务唯一标识(做幂等) |
| message_id | MQ 返回的 ID(发送成功后才有) |
| message_type | 消息类型(扫表时执行不同逻辑) |
| topic | 消息主题 |
| message_body | 消息体(序列化存储) |
| state | 待发送 / 已发送 / 已失败 / 已消费 |
| retry_count | 重试次数 |
| next_retry_at | 下一次重试时间 |
| fail_reason | 失败原因 |
索引设计:state 做索引用于扫表;message_key + message_type 建唯一索引防重复插入。
状态流转很简单:待发送是初始状态,已失败和已消费是终态。超过重试阈值(比如 10 次)就标记为已失败,降低重试频率(比如 6 小时一次)或报警人工跟进。
Broker 端丢失
Broker 收到消息后还没来得及刷盘就挂了,重启后消息丢失。或者主从同步延迟,主挂了但消息还没同步到从。
解决方案: Kafka 设置 replication.factor >= 3,min.insync.replicas >= 2,确保至少两个副本写入成功。RocketMQ 用同步刷盘(flushDiskType=SYNC_FLUSH)+ 同步复制(brokerRole=SYNC_MASTER)。
性能和可靠性永远是一对矛盾。同步刷盘比异步慢,但不会丢消息。根据业务对可靠性的要求做取舍。
消费者端丢失
消费者拿到消息后还没处理完就提交了 offset,然后处理失败或者应用崩溃了,消息就丢了。
解决方案: 关闭自动提交 offset(enable.auto.commit=false),改为手动提交 -- 处理成功后再提交。这样即使处理失败,下次拉取还能拿到这条消息。代价是可能重复消费,所以消费端必须做好幂等。
端到端的保障
三个环节都做好了还不够,建议再加一层对账机制。生产者记录发送日志,消费者记录消费日志,定时比对两边的记录,发现缺失的消息触发补偿。
二、消息重复消费与幂等处理
消息重复是分布式系统中无法完全避免的 -- 生产者重试、消费者处理完但 offset 提交失败、网络分区导致的消息重投,都会导致同一条消息被消费多次。
解决方案不是防止重复投递(做不到),而是让消费端做好幂等,即使同一条消息被处理多次也不产生副作用。
幂等号的选择
最常见的方式是通过消息中定义的幂等号做防重判断。这个幂等号应该是约定好的业务字段(比如订单号、支付单号),而不是 MQ 的 msg_id。原因是发送者重复发送多次消息会产生不同的 msg_id,但消息内容相同。如果用 msg_id 做幂等,就拦不住这种情况。
一锁二判三更新
有了幂等号之后,标准的幂等流程是:
注意前面讲过的坑:事务粒度不能大于锁粒度。如果用 @Transactional 包住整个方法,锁释放时事务可能还没提交,另一个线程拿到锁后查不到数据就会重复处理。
数据库唯一索引兜底
即使上面的流程做好了,极端情况下(锁失效、缓存故障)还是可能漏过。最后一道防线是数据库的唯一索引 -- 在业务表上对幂等号建唯一索引,重复插入直接违反约束被拒绝。
三、消息积压排查与解决
某天收到报警:消费者消费速度远低于生产速度,积压了几百万条消息。下游系统延迟飙高。
排查思路
先确认是生产加速了还是消费变慢了。查看生产者的发送量监控,如果没有异常变化,那问题在消费端。
消费变慢的常见原因:
- 消费者代码中有慢操作 -- 比如同步调用了外部接口、执行了慢 SQL、做了大量计算
- 消费者线程数不够 -- 单线程消费跟不上生产速度
- 下游依赖故障 -- 数据库连接池满、第三方接口超时
- 消费者发生了 Rebalance -- 频繁的消费者组重平衡导致消费暂停
解决方案
快速恢复(治标): 临时增加消费者实例数和分区数,用更多的消费者并行消费堆积的消息。如果不能加分区,可以在消费者内部用多线程并发处理。
根本解决(治本): 找到消费慢的根因并修复。如果是慢 SQL 就优化索引,如果是外部接口超时就加熔断降级,如果是线程数不够就调大线程池。
紧急降级: 如果积压实在太严重,可以先把消息转存到另一个 topic 或落库,让当前 topic 的积压快速消化。转存的消息后续慢慢处理。
四、Kafka 吞吐量优化实战
场景:单分区单消费者实例,不能增加分区也不能增加消费者,如何提高吞吐量?
异步消费
收到消息后先本地落库,落库成功立刻提交 offset。然后异步线程处理这条消息。成功了删除本地记录,失败了依赖本地记录重试。
有人问这样做 MQ 还有意义吗?当然有 -- 解耦、异步、削峰填谷这些能力仍然在发挥作用,本地落库只是消费端的一个优化手段。
多线程消费
配置 Kafka 一次性多拉取一些消息(调大 max.poll.records),拉取到之后用线程池并发处理。关键问题是 offset 提交的顺序性 -- 异步线程处理时间不同,提交顺序可能乱掉。
方案一:任务编排。 用 CompletableFuture.allOf() 等待所有线程处理完,再提交这批消息中最大的 offset:
CompletableFuture<Void> allFutures = CompletableFuture.allOf(
records.stream()
.map(record -> CompletableFuture.supplyAsync(() -> {
processMessage(record); // 处理逻辑
return null;
}))
.toArray(CompletableFuture[]::new)
);
allFutures.whenComplete((v, e) -> {
if (e == null) {
long maxOffset = records.stream()
.mapToLong(ConsumerRecord::offset).max().orElse(-1L);
consumer.commitSync(Collections.singletonMap(
partition, new OffsetAndMetadata(maxOffset + 1)));
}
// 有失败的:提交失败消息之前的最大连续 offset,或者落库重试
});如果某条消息处理失败,可以只提交失败消息之前的连续成功 offset,减少重复消费范围。或者把失败消息落库做重试。
方案二:本地缓存记录。 单实例场景下可以用本地缓存记录每条消息的处理状态和 offset。定时检查缓存,找到从最小 offset 开始的最长连续已处理序列,提交对应的 offset。比如处理了 100、101、102、104、105,那就提交 102。103 会在下次拉取中重新消费。104 和 105 通过幂等判断跳过。
本地缓存的缺点是重启会丢,但做好幂等的情况下重复消费问题不大。
消息压缩和组合
Kafka 支持 gzip、snappy、lz4 等压缩算法,设置 compression.type=gzip 减少网络传输数据量。如果生产者可以配合改造,还可以把多条消息组合成一个大消息发送,消费者收到后拆开处理。
调整 Kafka 参数
# 消费者端
enable.auto.commit = false # 手动提交 offset
fetch.max.bytes = 104857600 # 每次拉取最大 100MB
max.poll.records = 50 # 每次 poll 最多 50 条
# Broker 端
num.network.threads = 16 # 网络请求处理线程(CPU 核数 x2)
num.io.threads = 16 # 磁盘 IO 线程(CPU 核数 x2)
socket.send.buffer.bytes = 65536 # 发送缓冲区
socket.receive.buffer.bytes = 65536 # 接收缓冲区参数的合理设置需要结合硬件配置、Broker 负载和网络吞吐量综合考虑。
五、消息乱序的处理
下单后依次有支付消息和发货消息。虚拟商品支付后可能立刻自动发货,网络稍有延迟就可能导致发货消息先于支付消息到达消费者。
四种应对策略
顺序消息: 把有顺序依赖的消息按顺序投递到同一个 Partition,单消费者串行消费。最直接但吞吐量受限。
前置状态判断: 消息体中增加 beforeStatus 字段。消费时判断系统中的当前状态是否等于 beforeStatus,不一致就让消息处理失败等 MQ 重投。要求两个前提:消息能推进单据状态、状态是单向的不会回退。
序列号重排: 消息中附递增序列号,消费端缓冲后按序号排序再处理。增加复杂度,需要设置缓存超时处理丢失的消息。
本地事件表(推荐): 收到消息后做基本校验,校验成功就转成内部事件存数据库,立刻返回消费成功。然后异步线程处理事件,成功则更新为已处理,失败则不改状态、执行次数 +1。定时任务扫描未成功的事件重试。
这个方案的好处是所有消息都有存储,不怕丢失也不怕乱序 -- 因为可以在处理时按业务规则排序。前置状态判断依赖 MQ 重投可能导致消息最终丢失(长时间无法消费会停止重投),而本地事件表完全自主可控。
进一步优化:写入事件表的同时在 Redis 中记录业务单号和事件主键 ID。后续新消息来时检查 Redis 是否有同单号的待处理事件,有的话直接拿出来一起处理,减少定时任务带来的延迟。
六、Spring Event vs MQ 的选型
Spring Event 和 MQ 都是生产者-消费者模式,都支持一对多,都能异步。但在几个关键点上差异明显:
| 维度 | Spring Event | MQ |
|---|---|---|
| 持久化 | 基于 JVM 内存,重启就没了 | 磁盘持久化 |
| 削峰填谷 | 无缓冲 | 队列天然缓冲 |
| 失败重试 | 自行处理 | 自动重投 |
| 特殊功能 | 仅自定义线程池 | 延迟/顺序/事务消息 |
| 跨进程 | 不支持 | 支持 |
| 复杂度 | 极简 | 需要部署维护 |
选型原则: 系统内部做解耦或一发多收(比如订单创建后异步写流水),用 Spring Event + 线程池,再加定时任务做失败补偿。多系统交互、需要高级消息能力、并发量大需要削峰 -- 只能用 MQ。
七、设计一个消息队列
如果面试官让你设计一个 MQ,基于对 Kafka、RocketMQ 的理解做拆解:
基本架构: 生产者、Broker(存储+投递+备份)、消费者。Topic 做分类,Partition 做物理分割提高吞吐量。
存储: 内存存储快但可能丢,磁盘存储慢但可靠。实际中磁盘存储 + 操作系统 Page Cache。
通信: 用成熟的 RPC 框架(Dubbo/Thrift),或自研协议。支持推拉结合,实际用长轮询 -- 消费者发请求,Broker 有消息就返回,没有就等一段时间再返回或超时。Kafka 和 RocketMQ 都是这么做的。
可靠性: 主从复制、集群模式保证高可用。确认机制保证不丢消息。
高性能: 批量操作、顺序写入、零拷贝(参考 Kafka)。
功能扩展: 事务消息、延迟消息、顺序消息、死信队列、消息堆积处理、重平衡等。
推、拉还是长轮询
推: 实时性好,但生产速度大于消费速度时会压垮消费者。适合实时性要求极高的场景。
拉: 消费者自控速度,避免堆积。但需要不断轮询,有延迟且对 Broker 有压力。
长轮询(主流选择): 消费者发请求,Broker 有消息直接返回,没有就 hold 住连接等一段时间。等待期间有新消息就立刻返回,超时则结束。兼顾实时性和资源消耗。
还有一个实际限制:有些生产环境只允许单向通信(消费方 -> Broker),这时推模式的双向长连接不可行,只能用拉模式或长轮询。
为什么不用 BlockingQueue
BlockingQueue 能做简单的生产者-消费者,但有五个硬伤:无分布式特性(不能跨节点共享,可能流量倾斜)、无持久化(重启数据丢失)、无高级特性(不支持确认、重试、死信、延迟消息)、无管理监控、性能受限于单机。只适合应用内部的简单异步处理。
八、分布式架构的取舍
分布式一定比单体好吗
答案当然不是。两种架构各有适用场景,关键是匹配业务需求。
单体架构把所有东西部署在一起,好处很直接:开发测试部署都简单(一个项目、一次部署),不需要网络交互(除了数据库和缓存,代码都在一起,大大降低网络延迟),不需要考虑分布式事务、分布式锁、分布式 ID 这些令人头疼的问题。
但随着业务增长,单体的问题也很明显:单机性能有天花板(硬件资源有限),团队扩大后大家在一个应用上改代码合并起来痛苦万分(代码越来越庞大到最后没人敢改),一个模块出了 OOM 或 CPU 飙高全部受影响(共用 JVM、内存、CPU),技术栈被锁死在旧版本上。
分布式架构的好处是容易水平扩展(加机器抗更高并发)、模块化独立开发部署(你干你的我干我的,互相只依赖 API)、减少单点故障(你挂了我不至于跟着挂)、技术栈多样化(有自主权)。
代价也不小:运维复杂(多服务的部署、监控、日志、故障恢复),CAP 约束(可用性和一致性的权衡,随之而来的分布式事务、SLA 保障),问题排查困难(一个请求经过多个系统,需要分布式链路追踪)。
结论: 项目规模小、团队人数少、开发周期短、技术不复杂 -- 单体就够了。项目规模大、需求复杂、业务增长快、需要高并发高可用 -- 才需要上分布式微服务。
分布式 Session 怎么实现
用户登录信息保存在服务器 A,服务器 B 怎么获取?归根结底就是分布式 Session 的问题。最常用的方案是 Redis -- 把用户登录信息存 Redis,所有服务器都从 Redis 中读取。除了 Redis 也可以用 MySQL、MongoDB 等第三方存储。其他方案包括客户端存储(Cookie/JWT)、粘性 Session(负载均衡绑定)、Session 复制(服务器间同步),但都不如 Redis 方案通用和高效。
加分布式锁之后影响并发了怎么办
面试官问完分布式锁之后,经常追一句:"加了锁不是影响并发度了吗?"
很多人被问蒙了,觉得加锁降低了性能。但这里需要区分两种"并发":
分布式锁防止的并发是什么?是同一个用户的同一个行为的重复操作。比如重复下单、消息重投这些。这种并发本来就不该发生,是异常并发。
系统的 QPS 指的是什么?是所有用户在同一时段的正常请求总量,包括不同用户的请求。
通过加锁拦住了异常并发,系统不需要浪费资源处理这些重复请求,反而有更多资源去响应其他用户的正常请求。所以虽然防止了一部分并发,但拦截的都是不该发生的异常请求,系统整体吞吐反而可能上升。
九、本地消息表的进阶话题
下游失败了上游要回滚怎么办
先说结论:本地消息表不适合需要回滚的场景。它适合的是"不回滚,必须成功"的场景 -- 比如下单后给用户投保运费险,不能因为投保失败就取消订单,只能不断重试直到成功。
如果面试官非要回滚,有一个理论上可行但实际很麻烦的方案:多次重试不成功的消息进死信队列,上游监听死信队列做本地事务回滚。回滚前建议反查下游接口确认真的失败了,避免"实际成功了但回调没收到"的情况。
但说实话,一旦有多个消息监听者都要操作 -- 比如下单后依次运费险投保、给用户发消息、给用户加积分 -- 部分成功部分失败的情况下系统根本没法自动处理。两个系统之间因为回滚互相监听消息、互相依赖接口查询,耦合就很严重了。
务实的方案: 多次投递不成功就记录报警,靠人工介入。或者依靠对账机制做准实时对账,发现不一致的情况人工介入。这种场景发生概率很低,一旦发生往往有特殊原因需要人来判断。
消息表的状态管理
消息表的状态包括待发送、已发送、已失败、已消费、已挂起。其中待发送是初始状态,已失败和已消费是终态。
已失败状态不是完全放弃。可以根据重试次数做分级处理:
- 待发送的消息 10 分钟重试一次
- 已失败的消息 6 小时重试一次
- 超过最大重试次数(比如 20 次)的消息不再自动重试,报警人工跟进
next_retry_at 和 last_retry_time 这两个字段可以灵活使用。next_retry_at 适合需要主动设定下次执行时间的场景(比如特殊消息要 3 小时后重试)。last_retry_time 方便扫表时设过滤条件(只扫没重试过的,或 10 分钟前的)。两个字段有一个就行,都没有也不是不行 -- 取决于实际业务。
十、Kafka 多线程消费的完整方案
回到 Kafka 单分区单消费者的场景,多线程消费是提升吞吐量的关键手段。但 offset 提交的正确性是最大的挑战。这里给出一个完整的方案:
步骤
- 配置 Kafka 一次性多拉取消息(
max.poll.records=50) - 拉取到一批消息后,记录最小和最大 offset
- 用线程池(比如
CompletableFuture)并发处理每条消息 - 每条消息处理成功后,把它的 offset 记录到本地缓存(比如 ConcurrentSkipListSet)
- 从最小 offset 开始,找到第一个不连续的位置,提交之前的最大连续 offset
例如一批消息的 offset 是 100-105,处理完的有 100、101、102、104、105。从 100 开始找到 103 断了,提交 102 的 offset。下次拉取从 103 开始,会拉到 103、104、105、106、107。104 和 105 通过幂等判断跳过,只需处理 103、106、107。
为什么用本地缓存
因为题目限定了单分区单消费者单实例,不存在分布式一致性问题。本地缓存(比如 ConcurrentHashMap)既简单又高效。唯一的缺点是机器重启可能丢数据,但做好幂等后重复消费问题不大。
如果要更可靠,也可以用 Redis 或数据库存储处理状态,但对于单实例场景完全没必要增加这个复杂度。
失败消息的处理
如果某条消息处理失败了,两种策略:
- 宽容策略:不提交失败消息之后的 offset,等待 MQ 重新投递。适合消息量小的场景。
- 落库重试:把失败的消息保存到本地数据库,通过定时任务重试。提交已成功消息的 offset 继续前进。适合消息量大不能容忍重复消费开销的场景。
十一、消息积压的深度优化
消息积压是生产环境最常见的 MQ 故障之一。除了前面提到的基本应对策略,这里再深入展开几个实战中的优化手段。
消费者内部的异步化
很多时候消费慢不是因为 MQ 的问题,而是消费者的业务逻辑本身就慢。最有效的优化是消费者内部异步化:收到消息后先做参数校验和幂等判断,通过后立刻本地落库(或存内存队列),返回消费成功。然后用异步线程池处理实际业务逻辑。
这样做的好处是消费者的消费速度只取决于落库或入队的速度(非常快),而不受业务逻辑复杂度的影响。实际的业务处理失败了也不怕 -- 本地有记录,重试即可。
批量消费
如果消费者的每条消息处理逻辑中有数据库操作,可以考虑批量消费 -- 一次拉取一批消息,合并成一次数据库批量操作。比如 50 条积分增加消息,算出总积分一次性更新。数据库操作次数从 50 次降到 1 次,吞吐量大幅提升。
但批量消费需要处理好部分成功部分失败的情况 -- 一般是整批回滚重试,或者把失败的单独记录。
消费者限流与背压
消费者消费太快也不一定是好事 -- 如果下游数据库或接口承受不了高并发,消费者越快反而越容易把下游打垮。这时候需要在消费者端做限流(比如用信号量控制并发度),或者实现背压机制 -- 下游响应变慢时自动降低消费速度。
消息压缩节省带宽
Kafka 支持 gzip、snappy、lz4、zstd 等多种压缩算法。生产者端设置 compression.type 即可自动压缩。消息内容越大、重复信息越多,压缩效果越明显。snappy 和 lz4 侧重速度,gzip 和 zstd 侧重压缩率。一般推荐 lz4 -- 压缩和解压都很快,CPU 开销小。
如果生产者可以配合改造,还可以把多条消息组合成一个大消息发送(应用层合并),消费者收到后再拆开处理。这样降低了消息条数,减少了 MQ 的元数据开销。
Broker 端的参数调优
除了消费者端的优化,Broker 端的参数也很关键:
num.network.threads = 16 # 处理网络请求的线程数(CPU 核数 x2)
num.io.threads = 16 # 处理磁盘 IO 的线程数(CPU 核数 x2)
socket.send.buffer.bytes = 65536
socket.receive.buffer.bytes = 65536num.network.threads 负责处理客户端请求、响应和数据传输。num.io.threads 负责读写磁盘。两个缓冲区控制发送和接收时的缓冲大小。具体值需要根据硬件配置、负载情况和网络吞吐量综合调整。
十二、分布式事务与 MQ 的配合
可靠消息最终一致性
很多跨系统的数据一致性场景,并不需要强一致性(XA/TCC),最终一致性就够了。这时候 MQ 是最好的搭档。
核心模式是:上游先操作本地事务,成功后发消息通知下游。下游收到消息后处理自己的事务。如果下游处理失败,通过 MQ 的重试机制不断重投,直到成功。
关键保障点:
- 上游的消息不能丢 -- 用本地消息表或 MQ 事务消息保证
- 下游必须做幂等 -- 重复消费不能产生副作用
- 兜底对账 -- 定时比对上下游数据,发现不一致触发补偿
事务消息 vs 本地消息表
RocketMQ 原生支持事务消息:生产者先发半消息(Half Message),Broker 收到但不投递。然后生产者执行本地事务。如果本地事务成功,提交半消息让 Broker 投递。如果失败,回滚半消息。如果生产者迟迟没有提交或回滚,Broker 会定时回查。
本地消息表则更通用 -- 不依赖特定 MQ,任何 MQ 都能用。消息先写入本地数据库(和业务操作在同一个本地事务中),然后异步发送到 MQ。发送成功更新状态,发送失败定时任务重试。
两种方案各有优劣:事务消息实现简单但依赖 RocketMQ,本地消息表通用但多了一张表和定时任务。选择取决于技术栈和团队偏好。
最大努力通知
对于一致性要求最弱的场景(比如发邮件通知),用最大努力通知就够了 -- 发消息通知对方,对方收到就处理,没收到也不影响核心业务。可以重试几次,最终还是不行就放弃。
典型应用:订单创建后给用户发短信确认。短信没发出去不影响订单流程。
十三、MQ 选型与实战经验
Kafka vs RocketMQ vs RabbitMQ
三者各有擅长的场景:
Kafka: 吞吐量最高,适合日志采集、大数据流处理、事件溯源。分区机制天然支持水平扩展。缺点是功能相对基础,延迟消息和事务消息支持不如 RocketMQ。
RocketMQ: 功能最全面,原生支持事务消息、延迟消息、顺序消息、消息轨迹。适合电商、金融等业务场景。吞吐量略低于 Kafka 但足够用。国内社区活跃,阿里系项目的首选。
RabbitMQ: 延迟最低,支持灵活的路由规则(Exchange + Binding)。适合实时性要求极高的场景。缺点是吞吐量不如 Kafka 和 RocketMQ,Erlang 技术栈运维门槛较高。
不要为了用 MQ 而用 MQ
MQ 引入的复杂度不低 -- 消息丢失、重复消费、积压、乱序、延迟,每一个都是生产环境的定时炸弹。在决定使用 MQ 之前,先问自己几个问题:
- 这个场景真的需要解耦吗?直接 RPC 调用是不是更简单?
- 这个场景真的需要异步吗?同步处理能不能满足性能要求?
- 这个场景真的有削峰填谷的需求吗?数据库/下游服务扛得住直接调用吗?
如果以上答案都是"不需要",那就不要引入 MQ。系统内部的简单异步,Spring Event 就够了。跨系统的低频调用,RPC 就够了。只有当解耦、异步、削峰这三个需求中至少有一个是刚需时,才值得引入 MQ。
消费者的 Rebalance 陷阱
Kafka 消费者组在成员变化时会触发 Rebalance(重平衡)-- 暂停所有消费者的消费,重新分配分区。这个过程可能持续几秒到几十秒,期间所有消费者都不消费消息。
频繁 Rebalance 的常见原因:
- 消费者处理时间超过
max.poll.interval.ms,被认为"死掉"而被踢出组 - 消费者频繁重启
- 网络抖动导致心跳丢失
解决方案:调大 max.poll.interval.ms 和 session.timeout.ms,减小 max.poll.records 确保每批消息能在超时前处理完,确保消费者进程稳定不频繁重启。
生产环境的监控清单
MQ 上线后,至少要监控以下指标:
- 消息堆积量 -- 超过阈值立刻报警
- 消费者消费延迟(Lag) -- 持续增长说明消费跟不上
- 发送成功率 -- 下降说明 Broker 有问题
- 消费成功率 -- 下降说明消费逻辑或下游有问题
- 消费者组成员变化 -- 频繁变化说明有 Rebalance
- Broker 磁盘使用率 -- 接近满了需要清理或扩容
监控 + 报警 + 值班 on-call,是 MQ 在生产环境稳定运行的三板斧。
小结
消息队列的生产使用,核心就是三件事:不丢(三端保障 + 对账兜底)、不重(消费端幂等 + 唯一索引兜底)、不堵(多线程消费 + 参数调优 + 紧急降级)。
消息乱序用顺序消息或本地事件表解决。吞吐量优化靠异步消费、多线程并发、压缩组合和参数调整。选型时注意 Spring Event 和 MQ 各有适用边界 -- 系统内部用 Event,跨系统用 MQ。
最后,不管方案做得多完善,永远要有一个兜底机制:对账系统、定时任务扫表、人工介入通道。分布式系统的可靠性,从来不是靠单一环节的完美,而是靠多层防护网的叠加。