做后端开发这几年,我有个越来越强的感受:很多人在“消息队列重复消费”这类问题上反复踩坑,本质原因不是框架用得少,而是基础知识不够稳。项目一换、版本一升、场景一变,就露馅。所以这篇就是一份老本行角度的基础总结——消息队列到底是什么、有哪些核心概念、为什么会有重复消费、消息丢失和积压又该怎么排查。不敢说能让你原地成为架构师,但看完之后,再去读任何一款消息中间件的文档,你会有种“原来都一回事”的底气。
很多人第一次接触消息队列,是被“异步、解耦、削峰”这几个词吸引来的。但真正上手之后就发现,这几个词说得太轻巧了,落地时全是细节。别急,我从头把这些细节补上。
1. 消息队列到底在解决什么问题
1.1 “消息队列”这三个字拆开看
消息队列这三个字,其实已经说了全部:它是“消息”加“队列”。消息可以理解成一段结构化数据,比如订单号、用户ID、金额、时间戳;队列则是这条消息存放的地方。整条链路由三个角色组成:生产者负责把消息放到队列里,Broker(消息中间件)负责保管和分发,消费者负责从队列里取消息处理。
我用一个很生活化的例子:你给朋友发微信,朋友当时没看,但消息先存在服务器上,他什么时候上线什么时候读。这里面微信服务器就是Broker,你就是生产者,朋友是消费者。同步调用则更像打电话,你必须等我接起来才能说,我没接你就只能干等。
这个比喻能解释很多现象:为什么消息队列天然有缓冲能力?因为Broker把消息暂存起来了,消费方不在线也不影响生产方写入;为什么会出现消费延迟?因为消息虽然到了Broker,但消费方处理能力跟不上。理解了这套最基本的生产-存储-消费模型,后面所有概念都只是它的延伸。
1.2 异步、解耦、削峰,三个词分开说
异步最容易理解。比如下单成功后,系统要发短信、发推送、加积分、更新统计。如果同步调用这些服务,用户要等一两秒才看到“下单成功”。把这些动作丢进消息队列,主流程只要写库然后返回成功,短信、积分由消费者慢慢处理。用消息队列的异步,就是把用户不关心的耗时操作从主线里挪出去。
解耦稍微抽象一点。如果你在订单系统里直接调用库存系统的HTTP接口,订单接口就要知道库存服务的地址、接口签名、异常处理逻辑。一旦库存系统改接口,订单系统也要跟着改。引入消息队列后,订单系统只往队列里发一条“订单已创建”的消息,它不关心谁在监听。库存系统自己订阅这条消息去扣库存,将来就算再加一个风控系统,订单系统也完全不用动。
削峰讲的是瞬时压力。秒杀场景最典型,瞬间有上万人点抢购,数据库每秒能承受的写入量可能只有几千。没有消息队列时,数据库直接被击穿;用消息队列挡住前面,先把请求全部收下来,后端服务按自己能承受的速度慢慢消费,削掉的就是最高峰的那道冲击。这就像水库——上游发大水,先蓄起来,下游按需要放水。
1.3 什么情况下我建议先别用消息队列
不是所有项目都适合上消息队列。小项目、小流量、团队对中间件运维不熟的时候,引入消息队列反而会增加复杂度。本地事务一致性会变成分布式问题,消息重发带来幂等要求,消费失败要考虑重试,Broker挂了要考虑高可用,哪一样都是新增的成本。
我之前见过一个内部管理系统,日活几百人,业务很简单,但团队为了“以后扩展性”硬上了消息队列。结果多了一台服务器要维护,代码里多了一堆异步回调,出了问题还不好调试。后来那把代码重构成同步调用,整个世界都安静了。扩展性不是靠堆组件得来的,业务规模没有大到同步调用撑不住时,别拿消息队列给自己加戏。
如果你遇到的是强一致性要求非常高的业务,比如账户余额扣减,也建议谨慎。消息队列给不了强一致,它保证的是最终一致,中间会有时间差和各种补偿逻辑。存在这类场景,你得先想好幂等和对账方案再上,不然上线后的每一个重试都可能变成账单错误。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 把核心名词装进脑图
2.1 从一条消息的视角看名词
一条消息从发送到被处理,会经过这么几个概念:生产者、Broker、队列或主题、交换机或分区、消费者、消费位点。先看一张最简单的对照,我习惯按“逻辑层”和“物理层”来记:
| 名词 | 一句话解释 | 我踩过的坑 |
|---|---|---|
| Producer(生产者) | 负责发送消息的进程 | 很多人以为消息发出去就算完,其实没确认就存在丢失风险 |
| Broker(服务端) | 真正存储和转发消息的中间件节点 | Broker不等于集群,集群里可能有很多个节点 |
| Queue(队列) | 消息存放的逻辑容器,先进先出 | 多个消费者抢同一个队列不等于消息会发多份 |
| Topic(主题) | 一类消息的集合,类似业务分类 | 不同Topic之间天然隔离,别图省事混用 |
| Consumer(消费者) | 负责处理消息的进程 | 消费者有分组概念,不要只盯着单点 |
| Offset/位点 | 记录消费到哪条消息了 | 重置位点就是让老消息再跑一遍,很容易惹出重复 |
消息本身也有结构,通常是属性(Headers)加消息体(Body),有些中间件还有消息Key。属性里常放业务ID、事件ID,这一块很多人忽略,却是幂等设计和链路追踪的关键。比如RabbitMQ的消息可以带headers和routingKey,Kafka的每条消息有key,RocketMQ有messageKey。可以说,消息体负责“干什么”,消息Key负责“这条消息是谁”。
2.2 队列模型与发布订阅模型
消息中间件主要有两种语义:队列模型和发布订阅模型。队列模型里,一条消息只被一个消费者消费,多个消费者抢同一个队列,谁抢到谁处理。发布订阅模型里,一条消息会被所有感兴趣的订阅者各消费一遍,类比就是公众号推文,谁订阅谁收得到。
| 对比项 | 队列模型 | 发布订阅模型 |
|---|---|---|
| 消息去向 | 一个消费者 | 所有订阅者 |
| 典型场景 | 任务分发、异步处理 | 事件通知、多系统同步 |
| 代码模型 | Queue | Topic + ConsumerGroup |
| 示例 | RabbitMQ的普通队列 | Kafka的Topic、RocketMQ的Topic |
但要注意,这两者不是对立的。在大多数现代中间件里,通过“消费组”可以把发布订阅模型做成队列模型:让多个消费者加入同一个消费组,一条消息只会被组内一个实例消费;不同的消费组各订阅同一主题,消息被重复发给不同组,这就是发布订阅。想明白这个,你在RabbitMQ里设置队列绑定关系、在Kafka里配消费组的时候,思路会清晰得多。
2.3 消费组、位点与“至少一次”交付
消费组(Consumer Group)是很多中间件里的核心概念。同一消费组里的多个消费者实例共同消费一个Topic的消息,组内分摊任务;不同消费组之间逻辑上互相独立,同一条消息可以被不同组分别消费。这就像一个大任务分给多个同事干,但多个小组之间可以并行做不同的综合任务。
位点(Offset)则记录了消费进度。消费者每处理完一条消息,会提交位点。下次重启后继续从位点往下消费。问题就在这:如果你的代码先业务操作再提交位点,操作成功但位点提交失败,重启后会重复处理;如果你先提交位点再执行业务,业务中途异常,消息就丢了。这正是“至少一次”(At Least Once)和“最多一次”(At Most Once)两种语义的区别。大多数业务系统选择的是“至少一次”,保证不丢,但由消费方解决重复。
还有一个很少被讲透的点:重平衡。当消费组里有实例挂掉或新实例加入时,组内分区会重新分配,这时候消费位点可能短暂失效或跳跃,也会触发重复消费。所以写消息队列代码,第一行想的不该是“处理业务”,而是“这条消息可能之前已经处理过一次了”。
3. 重复消费问题:绕不开也躲不掉,只能正面干掉
3.1 重复消费为什么是必然存在的
先用一个顺序图似的流程想清楚。Broker把消息投递给消费者,消费者处理完后返回一个确认信号(Ack),Broker收到确认后,才认为这条消息处理成功,于是不再投递。问题来了:消费者处理完,正准备回Ack时,进程重启了,Ack没发出去;或者Ack发出去了但网络抖动,Broker没收到。Broker那边只知道“超时没收到Ack”,它走重发机制,于是同一条消息再次被投递。
这不是中间件的Bug,而是分布式系统里“网络不可靠,进程可能随时挂掉”这一现实带来的必然。为了不丢消息,Broker选择了重发;为了让业务能够承受重发,你只能做幂等。凡是宣称“绝对不重复”的消息队列产品,只能说在特定条件下把概率压得非常低,不可能做到数学意义上的零重复。
所以别和重复消费斗争了,接受它,然后设计消费方在“收到重复消息”时表现和“收到一次”完全一致,这才是正路。
3.2 幂等消费的标准动作
幂等实现五花八门,但核心套路就几类。
第一类是数据库唯一键约束。业务表里建一个唯一索引,比如订单事件表用“订单ID + 事件类型”作为唯一键。消费者收到消息就往表里插,插入成功说明第一次处理;插入报唯一键冲突就说明重复,直接跳过。这个方案简单可靠,最推荐。
第二类是状态校验。更新订单状态前,先查当前状态,只处理符合期望状态的变更。比如处理“已支付”消息前,确认订单当前状态是“待支付”。如果已经是“已支付”,说明重复投递或者顺序乱掉了,不需要做处理。
第三类是Redis去重。以“业务前缀 + 幂等键”作为Redis key,用SETNX加过期时间来标记已处理。要注意的是,Redis去重不适用于必须和DB事务保持强一致的场景,比如你更新了数据库但设置过期标记失败,就可能导致重复处理。它更适合缓存类、触发类场景。
我工作里最常用的是第一类,简单粗暴有效。比如消费订单消息时,专门建一张event_record表,事件ID做主键或唯一键,消费入口先插入这条记录,成功才算“真正开始处理”。
java复制@RabbitListener(queues = "order.queue")
public void handle(OrderCreatedEvent event) {
EventRecord record = new EventRecord();
record.setEventId(event.getEventId());
record.setBizId(event.getOrderId());
record.setCreateTime(new Date());
try {
eventRecordMapper.insert(record);
} catch (DuplicateKeyException e) {
log.info("重复消息已忽略, eventId={}", event.getEventId());
return;
}
// 真正执行的业务逻辑
doBiz(event);
}
别小看这么一层,它解决的是你深夜被叫起来处理“订单重复创建”问题的概率。
3.3 顺序问题与消费并发之间的纠葛
重复消费已经够麻烦了,如果业务还要求消息严格有序,难度再上一层。比如针对同一个订单,有“创建”、“支付”、“发货”三条消息,如果“支付”先被消费,“创建”还没执行,处理逻辑就崩了。
要让顺序有保障,首先要理解顺序的范围。大多数消息队列保证的是“分区内有序”或者“单队列内有序”,而不是全局有序。Kafka里消息是按key哈希到分区的,同一个key的消息会进入同一个分区,分区内按写入顺序存储。消费端如果只用单线程消费某个分区,就可以保证这个分区内的消息按顺序处理。但一旦消费者开多个线程,同一分区的多条消息被并发消费,顺序就无法保证了。
我的实际经验有两个:一是发送端必须保证业务主键相同的消息进同一个分区或队列,这样顺序的“地基”才存在;二是消费端处理完一条再取下一条,简单粗暴但有效。对应到Kafka,就是单分区配单线程;对应到RabbitMQ,就是单队列配单消费者,配合手动确认串行处理。
也要想清楚,很多业务其实不需要全局顺序。比如你只关心订单最终状态,或者最后统计结果正确就行,那中间几条消息乱序无妨。能用版本号或者状态机兜底的,就别硬追求全局顺序,那会牺牲掉大量吞吐性能。
4. 实操链路:从零把一条消息发出去再收回来
4.1 本地用 Docker 把 Broker 拉起来
光讲概念很容易飘,我建议你本地跑一个最小环境,亲手发一条消息体会一下。我给你挑个上手成本最低的RabbitMQ,一条命令就能起来:
bash复制docker run -d --name rabbitmq \
-p 5672:5672 \
-p 15672:15672 \
-e RABBITMQ_DEFAULT_USER=guest \
-e RABBITMQ_DEFAULT_PASS=guest \
rabbitmq:3.12-management
5672是AMQP协议端口,15672是管理控制台。启动后浏览器打开 http://localhost:15672 ,输入guest/guest,你能看到队列、连接、消息速率等状态。我习惯先把控制台开着,发消息时盯着看,你能直观看到消息积压、入队出队的数字变化,这比看文档管用。
如果你公司用的是Kafka或RocketMQ,启动方式不同,但概念是同一套。Kafka要额外配合ZooKeeper或者KRaft模式,RocketMQ有NameServer和Broker两个进程,这些都属于Broker侧实现细节,不影响你对“生产-存储-消费”模型的理解。
4.2 生产者与消费者的最小代码
我用Spring Boot加spring-boot-starter-amqp这套组合演示。先在application.yml里配好连接:
yaml复制spring:
rabbitmq:
host: 127.0.0.1
port: 5672
username: guest
password: guest
listener:
simple:
acknowledge-mode: auto
prefetch: 1
template:
exchange: order.exchange
routing-key: order.created
生产端,注入RabbitTemplate直接发消息:
java复制@Autowired
private RabbitTemplate rabbitTemplate;
public void publish(OrderCreatedEvent event) {
String json = objectMapper.writeValueAsString(event);
rabbitTemplate.convertAndSend("order.exchange", "order.created", json);
log.info("消息已发送, orderId={}", event.getOrderId());
}
消费端,用监听器收消息:
java复制@RabbitListener(queues = "order.queue")
public void onOrderCreated(String message) {
log.info("收到消息: {}", message);
// 业务逻辑
}
注意这里的两个组件:交换机(Exchange)和队列(Queue)是通过绑定关系联系起来的。发送时用routingKey指定路由,交换机按规则把消息投递到匹配的队列。这个规则抽象非常重要,它让生产端完全不用知道队列的名字,将来队列换名、拆分,生产端代码改动为零。
4.3 真正需要花时间理解的那几个配置
配置项里最容易踩坑的是三个:确认模式、预取数量、持久化。
确认模式控制消费者的Ack时机。自动确认模式下,消费者还没执行代码就告诉Broker“我收到了”,如果业务逻辑抛异常,消息已经标记成已消费,直接丢。手动确认模式需要显式调用确认方法,处理成功才Ack,失败则回复拒绝或重新入队,但代价是代码更繁琐。我的默认建议:业务可靠性强的场景,一律手动确认,acknowledge-mode改成manual,在处理成功后再确认。
预取数量(Prefetch)控制消费者一次从Broker取多少条消息到本地内存。默认数值在一些客户端里是无限的,消费者本地缓存了几百条消息,其他消费者实例一直闲着,积压下反而更快。把prefetch设成1,让每个消费者一次只取一条,处理完再取下一条,这样负载更均衡,也天然降低了重复和忙死一个节点的情况。
持久化决定Broker重启后消息还在不在。RabbitMQ里队列和消息都要设置持久化才可靠,这两步少了任何一步,重启丢消息都没得商量。在Kafka里对应的是副本数和acks参数。配置的时候就多想一步:这条消息丢了,业务能不能接受?能接受就怎么快怎么来,不能接受就把持久化开满,但也要做好性能变慢的心理准备。
5. 消息丢失与积压:老板最怕的两种故障怎么排查
5.1 丢消息,说的是哪一端的丢
处理丢失问题,很多人上来就怀疑Broker,其实消息丢失发生三处,排查思路完全不同。
| 丢失位置 | 常见原因 | 最直接的排查手段 |
|---|---|---|
| 生产端 | 发送时报错未确认;业务失败了还当成功 | 查看生产日志,确认是否收到Broker的发布确认 |
| Broker端 | 队列非持久化;单节点无副本,宕机丢数据 | 检查队列持久化标记;检查集群副本数和同步状态 |
| 消费端 | 自动确认开启;业务异常被吞掉;位点误提交 | 看消费日志,留意Ack时机;关闭自动确认 |
生产端有个很隐蔽的问题:RabbitTemplate默认的发送方法在没有Broker确认时也会正常返回。你以为发送成功了,其实消息半路丢了。解决方案是开启Publisher Confirms模式,发送后等待Broker确认。总有人觉得这影响性能,实测下来在合理主集群下这点开销远小于半夜处理丢单问题的成本。
消费端则要记住一个反直觉规律:先确认,后业务逻辑,消息会丢;先业务逻辑,后确认,逻辑到位但确认失败,会重复。这两者的选择其实就是“丢”和“重”的取舍,绝大多数场景选后者,也就是“宁可重复,不可丢失”。
5.2 积压:消费者跑不动,先救人再查因
消息积压是系统性的红灯警报。最常见的原因有这几类:消费者代码出现异常,Ack一直不返回;消费者实例数量太少,处理能力不足;下游数据库或接口变慢,拖住整个消费链路。
接到报警我的处理顺序是这样的。第一步,先看监控,找到积压最严重的队列和Topic,确认积压数量级。第二步,查消费者日志里有没有大量异常,有异常就先定位修复,别急着扩容——你扩一百个消费者,每个都在抛异常,积压只会越滚越大。第三步,如果业务逻辑正常,纯粹是消费能力不足,最直接的手段是增加消费者实例,注意消费组内实例数和分区数的关系,比如Kafka里分区只有3个,你开10个消费者实例,其中7个是白开的,它们分配不到任何分区。
还有一招容易被忽略,就是隔离优先。如果整条消费链路里只有某一类特殊消息导致处理卡住,常规方案是临时把它分流到单独的队列或Topic,先保证主链路泄洪,再慢慢修特殊分支。很多技术人员遇到积压就慌,其实慌张才是最可怕的——对着正常消费者乱改一通,可能把原来没问题的链路搞挂。
5.3 排查复盘速查表
把最三种高频问题的排查动作压成一张表,方便你关键时刻照着来:
| 症状 | 初步怀疑 | 必查动作 | 常用止血方案 |
|---|---|---|---|
| 消息反复消费 | 确认机制超时/重复投递 | 查消费日志的异常堆栈;查消费位点提交是否频繁失败 | 引入幂等表;调整确认超时时间 |
| 消息丢失 | 自动确认/持久化未开 | 检查ack-mode;检查队列持久化和发送确认模式 | 关闭自动确认;开启发布确认 |
| 消息积压 | 消费异常或容量不足 | 查消费者日志;查分区数和消费者实例数 | 修复消费异常;扩容实例;隔离处理慢消息 |
我每次排查完都会把时间线和证据贴到复盘文档里,而不是只说“解决了”。因为消息队列的问题很少只出现一次,只要环境不变,下次它还会以类似的形式出现。留有排查路径,下次可能十分钟就定位完毕。
关于消息队列,我还有一个个人的切身体会:任何一份讲“基础知识”的文档,最终都要落到“消息一定会重复,系统必须设计成能接受重复”这个共识上。你越早接受这一点,越早把幂等设计敷衍成基础设计,后面的路越顺畅。
最后顺手说一个小技巧:你的消费逻辑入口处,永远给每条消息打一个traceId并记录到来消息的时间。线上出现问题,你不需要猜是谁发的、什么时候发的、处理到哪一步,日志一拉清清楚楚。这个习惯帮我省过太多通宵查问题的时间,值得你从现在开始养成。
