如果你在某个项目里见过这样的场面——大促时订单服务扛不住下游通知的流量,或者数据同步任务在凌晨积压了上万条记录,又或者两个系统之间接个接口就得为对方的故障背锅——那说明你们项目已经在考虑上消息队列了。消息队列这个东西,说穿了就是一个用于传递数据的"中间站":生产者把消息放进去,消费者按自己的节奏取出来,中间这个站负责暂存、路由,把双方的时间节奏彻底解耦开。这篇内容不聊某个具体产品的源码细节,而是把消息队列项目里最常用到的那套基础知识系统串一遍,从核心模型、三大经典难题(重复消费、消息丢失、顺序问题)、积压应对,到 Kafka、RabbitMQ、RocketMQ、MSMQ 等主流产品的选型逻辑,适合正在做技术选型、刚接手消息队列项目,或者想补一补基础短板的同学收藏参考。
1. 项目里究竟什么场景才需要消息队列
1.1 消息队列在系统架构中的角色
先把概念拉直:消息队列(Message Queue)本质上是一个基于生产者-消费者模型的中间件组件。生产者发送消息,消费者接收消息,消息队列本身承担暂存、路由、分发、可靠传输这些职责。它解决的问题不是"能不能传数据",而是"如何让数据传得更稳、更灵活、更不过分耦合"。
很多刚接触的人容易把消息队列理解成一个"增强版的网络传输工具",这是不够准确的。RPC 也好、HTTP 调用也好,是请求方和服务方直接打交道,双方必须同时在线、接口协议必须对齐、一个出问题另一个立刻感知。消息队列不同,它像是一个中转仓库,生产者只管把货物放进仓库,消费者什么时候来取、怎么取,完全由消费者自己决定。这种"放进去就不用管"的模型,带来的是时间维度上的解耦。
我见过很多项目,上了消息队列之后代码结构反而变乱,原因就是没想清楚哪一步该用、哪一步不该用。消息队列不是银弹,该用 RPC 直连的强一致性场景(比如扣库存、付款),硬塞一个消息中间件进去,只会把事务边界搞得一团糟。
1.2 异步、削峰、解耦:三个最核心的价值场景
消息队列能解决的问题,归纳起来是三件事:异步处理、流量削峰、系统解耦。这三个词在各种文章里出现频率极高,但落到具体项目里,含义可以被拆得很细。
先说异步。最典型的例子就是用户注册成功后发送短信和邮件。如果同步调短信服务商接口,一次注册接口的耗时可能会从 50ms 涨到 500ms,而用户根本不需要等着这两条消息发完才算注册成功。把"发通知"这个动作丢进消息队列,注册接口立刻返回,消费端再去做慢操作,这就是时间维度上的优化。异步带来的直接收益是接口响应时间变短,系统整体的吞吐量也会因此受益。
再说削峰。做电商活动的同学体会最深:平时订单量每秒几百,活动开始瞬间冲到每秒几万。如果订单服务直接面对这种瞬时流量,数据库连接池、下游库存系统都会被瞬间打穿。这时候消息队列就像一个缓冲区,生产者(网关/客户端)把订单消息快速投递进来,消费端根据自己的处理能力按固定速率消费。瞬时洪峰变成了持续的平滑流,这就是"削峰"的直观效果。顺带一提,削峰不等于降低总请求量,它是把短时间内无法消化的请求延后处理,能真正做到"先收下,慢慢做"。
最后是解耦。假设订单系统需要通知库存系统、积分系统、消息中心。如果走接口直连,每新增一个下游,订单系统就要改代码、加调用、处理新下游的各种异常。改用消息队列之后,订单系统只需要定义好"订单创建"这个消息,谁关心谁订阅,订单系统不再关心下游是谁、有几个、会不会挂。这个场景下的收益不是性能,而是架构演进的自由度。
1.3 什么场景下反而不该用消息队列
这一点重要到值得单独拎出来说。消息队列引入之后,系统复杂度会显著上升:多了一个中间件要运维,消息可能丢失、重复、乱序,排查链路变长。如果你只是 A 调 B 的同步接口,且没有异步、削峰、解耦这三类需求,直接 HTTP/RPC 调用更合适。
尤其是强一致性场景。比如账户扣款,需要确认余额扣减成功之后才能响应,这就没法用消息队列做异步化,因为消息的到达和最终一致都无法满足"立即生效"的要求。记住一个判断原则:消息队列适合的是"最终一致"和"可延后"的业务,不适合"必须立刻返回结果"的事务性操作。
2. 消息模型:从最基础的五个概念到不同产品的落地差异
2.1 消息、生产者、消费者、Broker、Topic:先记住这五个词
不管你用哪个具体产品,消息队列项目的核心概念跑不出这五个:
- 消息(Message):在队列中传递的数据单元,包含消息体和属性(比如业务ID、时间戳、路由键)。
- 生产者(Producer):负责创建、发送消息的一方。
- 消费者(Consumer):负责从队列中拉取或接收消息并进行处理的一方。
- Broker:消息中间件服务器本身,负责存储、转发、管理队列。
- Topic(主题):消息的分类单位,生产者发送到某个 Topic,消费者订阅某个 Topic。同一个 Topic 的消息可以被多个消费者组分别消费。
把这五个概念放在一起看,消息队列的工作流就清晰了:生产者把消息发送到 Broker 上的某个 Topic,Broker 负责存储,消费者启动后从该 Topic 拉取消息,处理完成后告诉 Broker"我已经处理完了"。就这么简单。
2.2 队列模型与发布订阅模型的本质区别
早期消息队列(比如 MSMQ、大部分传统企业集成场景)采用的就是典型的点对点队列模型:一条消息只会被一个消费者取走,消费完成后消息就从队列中移除。这个模型适合"任务是我安排的,谁干都行,但只能干一遍"的业务,比如异步导出报表、批量发送通知。
发布订阅模型则完全不同。生产者把消息发到 Topic,多个消费者组都能订阅它,每个组都能得到一份完整的消息副本。比如订单创建成功这条消息,订单消费者用来更新订单状态,风控消费者用来做风险分析,数据仓库消费者用来采集指标,各自消费各的,互不干扰。
这里有一个新手容易绕晕的点:同一个消费者组内的多个消费者实例是竞争关系,不同消费者组之间是订阅关系。这个概念在 Kafka 和 RocketMQ 里尤其重要。组内竞争意味着一条消息只能被组内的某个实例消费一次,组间订阅则让一条消息可以被多个组重复消费。只有把这个关系搞清楚,做水平扩展的时候才知道该加消费者数量还是加消费者组。
2.3 不同产品如何落地这些模型:Topic、Queue、Exchange、MessageQueue
虽然概念相通,但每个产品的具体命名和实现细节差别很大,做项目的时候要提前适应。
Kafka 的模型是 Topic 下分多个 Partition(分区),消息写入时按 key 做 hash 分布,或者轮询分配。消费者组内的消费者实例与 Partition 之间是一对多的分配关系。这个设计为 Kafka 带来了极高的吞吐,但也意味着它更强调"流式处理",而不是传统的点对点队列。Kafka 里"队列"的色彩已经很淡了,它的默认语义是"一个 Topic 可以在多个 Consumer Group 之间重复消费",这和开源的 Kafka 生态定位(日志处理、流数据管道)一致。
RabbitMQ 的模型更接近传统队列,核心概念是 Exchange(交换机)、Binding(绑定)、Queue。生产者不直接把消息发给队列,而是发给 Exchange,由 Exchange 根据绑定规则路由到对应的 Queue。它有 direct、topic、fanout、headers 等路由模式,路由能力非常灵活,适合复杂路由规则的项目。不过在 Kafka 生态里常见的"分区并行"能力,RabbitMQ 需要用多个 Queue + 一致性哈希之类的方案去模拟,历史上不如 Kafka 自然。
RocketMQ 的模型介于两者之间,它既有传统消息队列的易用性,又引入了 Kafka 式的分区机制。RocketMQ 的 Topic 下也有多个 MessageQueue(消息队列),消费者在消费时按队列分配。它把 Kafka 的所有语义都做进了一个"更像可靠消息队列"的产品里,所以国内很多团队在 Java 技术栈里选择 RocketMQ,既能获得高吞吐,又保留事务消息、顺序消息、定时消息这些偏业务的功能。
MSMQ(微软消息队列)则是 Windows 系统自带的老牌消息队列组件,模型就是最简单的点对点队列,也有事务性队列、日记队列这些企业级功能,但它没有 Topic/分区这套现代分布式模型,跨平台能力也弱。关于 MSMQ 的定位和局限,后面选型章节还会展开。
2.4 ACK 确认与重试机制:消息从发到收的状态流转
消息队列项目里最常出问题的地方,往往不是消息怎么发,而是消息"算不算成功"。几乎所有主流消息队列都有一套 ACK(确认)机制来回答这个问题。
以 Kafka 为例,生产者发送消息时有 acks 参数。acks=0 表示不等 Broker 确认,最快但最可能丢消息;acks=1 表示写入主分区就算成功;acks=all(即 -1)表示所有同步副本都写入才算成功,最安全但吞吐下降。消费者的 ACK 则体现在提交 offset(消费位点)上——Kafka 消费者处理完消息后提交 offset,Broker 才知道这条消息可以继续推进,否则重启后会从之前的 offset 重新消费。
RocketMQ 和 RabbitMQ 也有类似的确认逻辑。RocketMQ 支持在消费成功后返回 CONSUME_SUCCESS,或者主动将消息回置为重试;RabbitMQ 则是手动 ack / nack。这里的核心思想都一样:确认的时机决定了消息的语义边界。如果确认得太早(比如拿到消息就算成功),消息处理失败时就会丢失;确认得太晚(处理完所有后续逻辑才确认),又会拖低吞吐。实际项目里的最佳实践是:在业务逻辑执行成功之后、在事务提交或关键副作用产生之后,再执行 ACK。把你的处理流程设计成"先干活、再确认、失败重试",而不是"先确认、再干活"。
3. 重复消费、消息丢失、乱序:项目里躲不开的三个问题
3.1 重复消费问题:为什么会出现,以及怎么根治
热词里排在第一位的就是"消息队列重复消费问题",可见这是项目实战中遇到最多、最让人头疼的一个。首先要纠正一个观念:重复消费不是某个消息队列产品的 bug,而是分布式系统里无法彻底避免的常态。
重复消费产生的根源在于"消息到达的可靠性"与"确认机制的不完美"之间的矛盾。最常见的情形是:消费者 A 从 Broker 拉取了一条消息,开始执行业务逻辑,处理到一半网络闪断,Broker 迟迟收不到 ACK,于是认为消费者 A 已经挂掉,把这条消息在另一个消费者 B 上重新投递。此时 A 其实已经把业务做完了,B 又做了一遍,就造成了重复。还有一种情形发生在生产者侧:生产者向 Broker 发送消息后网络超时,框架层自动重试,Broker 第一次其实已经收到了,二次发送又收到一条一模一样的消息,消费端同样会重复处理。
明白了产生原因,解决方案也就清晰了:既然重复无法从源头杜绝,那就让"重复处理"没有副作用。这就是幂等设计。
所谓幂等,就是同一个操作执行一次和执行多次的结果完全一致。实现幂等有几种常用套路:
- 数据库唯一约束:消费消息后,把业务ID作为唯一键插入去重表。重复消费时插入会因唯一键冲突而失败,此时直接识别为已处理,跳过业务逻辑。
- Redis SetNX 去重:用业务ID作为 Redis key,setnx 成功表示首次处理,失败说明已经处理过。注意 key 要设置合理的过期时间,配合上锁防止并发重复。
- 版本号或状态机校验:更新时携带版本号,update ... where version = old_version,乐观锁能天然挡掉重复更新。
- 消息内携带业务唯一 ID:这条很多人会忽略。在发送消息时,最好在消息体内带上业务的唯一标识(比如订单号、操作批次号),而不是依赖消息自身的 msgId。因为网络重试可能会生成不同的消息 ID,只有业务标识才能真正区分"同一条业务消息"。
我在实际项目里还踩过一个坑:只做了消息去重表,但没关注表的清理策略。去重表会随时间膨胀,最终影响插入性能,甚至拖慢业务。建议给去重记录设置 TTL,或者定期归档已经超过一定时间的老数据。另外,去重逻辑要放在业务事务的同一个事务里执行,否则会出现"去重记录已提交,业务没成功"或者反过来"业务成功,去重记录未提交"的中间态。
3.2 消息丢失问题:三个环节分别怎么防
消息丢失问题比重复消费更隐蔽,因为一旦发生,往往要过很久才能通过对账发现。一条消息从生产到消费,会经过三个环节:生产端、Broker 端、消费端,每一环都有对应的丢失风险,也都有对应的防护手段。
生产端丢失:最常见的原因是发送模式太过随意,发送成功与否没有确认。Kafka 里如果你用 acks=0,消息发出就不管了,Broker 万一在写入前宕机,消息就永远没了。防护手段是做成同步发送 + 失败重试,或者用带回调的异步发送并在回调里处理失败,同时把发送结果写入日志,方便事后对账。
Broker 端丢失:Broker 收到消息后先放内存再落盘,如果还没落盘就宕机,消息就会丢。防护手段有三个层次:一是开启持久化配置,Kafka 可以设置 log.flush.interval.messages,RabbitMQ 可以设置持久化交换机和持久化队列,RocketMQ 默认刷盘;二是开启副本机制,Kafka 的主题副本数至少设为 3,min.insync.replicas 设为 2,确保多数副本写入成功才确认;三是区分同步落盘和异步落盘的可靠性差异,花钱买性能还是买可靠性,要在这里做清醒的取舍。
消费端丢失:最典型的场景是消费者收到消息后在处理前就提交了 ACK/offset,结果业务逻辑抛异常,消息就再也找不回来了。防护手段很简单:不要用自动提交,改成手动提交,并且先处理业务逻辑、成功后再提交。这条建议听着像废话,但很多项目挂掉的直接原因就是图省事开了自动提交。
3.3 顺序问题:全局有序和部分有序怎么取舍
消息乱序问题同样经典。典型场景:订单创建消息先发出,订单更新消息后发出,如果两个消费者并发处理,更新消息先被处理了,创建消息才被处理,数据库里的订单状态就错了。
先明确一点:全局有序是最理想但也最难的状态。要实现全局有序,要么把所有消息都放进单个队列单消费者处理,牺牲并行度;要么在消费者内部加锁串行化,吞吐量会大幅下降。绝大多数业务真正需要的,只是局部有序——同一类消息、同一个业务维度内有序就够了。
实现局部有序的标准做法是按业务 key 做路由。Kafka 支持指定 key 进行分区:相同 key 的消息会进入同一个 Partition,而一个 Partition 只会被组内一个消费者实例处理,所以天然保持分区内顺序。RocketMQ 把这种能力做成了开箱即用的"顺序消息",支持全局顺序和分区顺序两种级别,实际项目里用分区顺序就足够了。RabbitMQ 没有内置分区,可以用一致性哈希插件,或者为每个业务 key 建独立队列,这两种做法都可行,但复杂度不低。
还有一层容易踩的坑:即使生产者发送顺序正确,消费端的重试机制也可能打乱顺序。比如消息 A 处理失败进入重试,后到的消息 B 没有依赖先执行,B 被处理完,A 重试成功后再落地,顺序照样乱。解决思路是:顺序消息场景下,一个分区内的消息如果某条处理失败,一般需要立即暂停当前分区后续消息的消费(停滞在那里做重试),而不是跳过继续,这样才能保住顺序。
4. 消息积压:当消费者跟不上生产速度时
4.1 积压的典型症状与深层原因
消息积压是生产环境最常出的事故之一。症状很直观:队列里的未消费消息数量(backlog)持续上涨,消息消费不断延迟,业务方开始接到用户投诉"消息怎么还没到"。
积压的常见原因,我归纳为三类。第一类是消费速度天然跟不上生产速度,属于容量规划问题,比如大促期间流量暴涨。第二类是消费者实例异常,比如消费者进程僵死、持有的连接被 Broker 踢掉、消费者启动时崩溃但没被及时发现。第三类是死循环或单条消息处理过慢:某条消息触发了死循环,或者消费端调用了超时极长的下游接口,消费者被一条消息卡住,后面的消息全部排队。
还有一种容易被忽略的积压原因:消费者线程数配置过低。很多框架默认消费者只有一个线程循环拉取消息,如果你的业务处理里有大量 IO 等待(比如调外部接口、查数据库),单线程根本吃不满网络带宽和 CPU,积压是必然结果。这时候正确的姿势是把"拉消息"和"处理消息"彻底分离:消费者只负责快速拉取并扔进本地线程池,业务逻辑在线程池里并发执行。
4.2 排查链路:从查看堆积到定位根因的完整流程
遇到积压,我的排查习惯是固定的,尽量用数据说话:
- 先看消息队列管理端的堆积量趋势,确认是持续上涨还是短暂波峰。持续上涨说明是结构性问题,短暂波峰可能只是瞬时突发。
- 看消费者的核心指标:消费速率(messages/s)、拉取延迟、重复消费次数。如果消费速率接近 0,大概率是消费者实例挂掉或者卡死。
- 看消费者进程的线程状态:是不是有线程长期阻塞在 IO 上,有没有死锁,GC 是否过于频繁。
- 看下游依赖:消费端调用的数据库、外部 API 响应时间是否飙升,有没有慢 SQL 全表扫描。
- 最后看日志里的异常,比如反序列化失败、业务规则变动导致的消息格式不兼容。
这一套走完,90% 的积压都能定位到根因。有一个原则:不要急着重启消费者,先拿到现场数据再动手,否则很容易复现又再次积压。
4.3 积压发生后的紧急应对与根治措施
紧急应对的核心思路是"先恢复业务,再排查根因"。通常按阶梯执行:
- 扩容消费者:水平增加消费者实例/线程。要注意的是,如果消费者组绑定了固定分区,新增实例对 Kafka/分区型队列才有效;RabbitMQ 的多消费者则要保证消息均匀分发。
- 临时关掉不必要的下游逻辑:比如消费端调用的非核心通知、日志记录,先全力消化堆积的业务消息。
- 如果堆积量实在太大,可以先把积压消息转存到临时 Topic/数据库,消费者专心处理一个较小的 backlog,处理完再重放。
- 极端情况下对老消息做抽样丢弃并补偿,但这种操作必须和业务方确认,千万不要自己决定。
根治措施则要回到架构:给消费端加多级缓冲(拉取线程池 + 业务处理线程池);对消费端调用的依赖做超时和熔断,防止下游故障拖垮消费速度;给关键 Topic 设置积压监控报警,积压量超过阈值自动通知。另外,消费端的重试策略也要合理,无限重试会让一条坏消息无限占用消费者线程,正确的做法是设置最大重试次数,超过后把消息转入死信队列,或者降级处理。
5. 选型实测:Kafka、RabbitMQ、RocketMQ、MSMQ 怎么选
5.1 几款主流消息队列的现状梳理
做项目选型时,话题基本围绕几个固定选手:Kafka、RabbitMQ、RocketMQ,以及老牌环境中仍然存在的 MSMQ/ActiveMQ。我把它们摆在一起做个横向对比,方便快速区分:
| 维度 | Kafka | RabbitMQ | RocketMQ | MSMQ |
|---|---|---|---|---|
| 核心模型 | Topic + Partition | Exchange + Queue | Topic + MessageQueue | 点对点 Queue |
| 单机吞吐 | 极高(百万级) | 中等(万级) | 高(十万级+) | 低 |
| 消息可靠性 | 高(副本+ACK) | 高(持久化+ACK) | 高(事务消息+同步刷盘) | 中(受限于 Windows 环境) |
| 路由灵活性 | 弱(依赖分区键) | 强(多种 Exchange 类型) | 中(有 Tag/自定义过滤) | 弱 |
| 顺序消息 | 分区内有序 | 需自己设计 | 支持全局/分区顺序 | 不支持 |
| 事务消息 | 支持(需二次封装) | 弱 | 原生支持 | 支持本地事务队列 |
| 跨平台 | 强 | 强 | 强 | 仅 Windows |
| 运维成本 | 较高(依赖 ZooKeeper/KRaft) | 较低 | 中(NameServer 简单) | 低(系统组件) |
| 典型场景 | 日志管道、大数据、数据同步 | 系统内部异步、灵活路由集成 | 电商、金融、订单类业务 | 老 Windows 企业内部集成 |
5.2 Windows 消息队列 MSMQ 的历史定位与选型建议
热搜词里出现了"windows消息队列"和"msmq消息队列",这里单独说几句。MSMQ 是微软提供的一个老牌组件,不需要额外安装消息中间件,Windows 环境里配置好就能用,甚至很多老企业系统的本地集成还是靠它。它的优点是简单、系统自带、对 Windows 生态下的 .NET 应用友好,支持事务性队列和日记队列,这在当年确实解决了大量实用问题。
但它的短板也相当明显:没有 Topic/分区这套现代模型,扩展性受限于单台 Windows 服务器,跨平台能力弱,官方文档和维护资源越来越少,消息堆积和低吞吐问题在互联网业务场景几乎无解。我的建议是:新项目哪怕跑在 Windows 上,也优先考虑 RabbitMQ 或 Kafka;只有当你维护的是一个跑了很多年的老系统、且现状完全没法动基础设施的时候,MSMQ 才作为历史遗留方案继续存活。如果非要用,尽量把它封装在一个独立的发送/接收组件后面,为将来替换留一条退路。
5.3 选型的判断依据:吞吐、可靠性、团队、成本四个维度
具体到你的项目,选型不必追求最强,而是要最匹配。我总结的四个判断维度:
- 吞吐量要求:如果每天处理几百万条业务消息,RabbitMQ 完全够用;如果要做日志管道、实时数仓、每秒几十万的接入,Kafka 是绕不开的;RocketMQ 则处于两者之间,适合 Java 技术栈的高吞吐业务场景。
- 可靠性与事务要求:金融、订单这类强可靠的业务,RocketMQ 的事务消息和同步刷盘非常顺手;Kafka 需要自己设计事务流程;MSMQ 的可靠性依赖系统环境,新项目不推荐。
- 团队技术栈与运维能力:Java 为主的团队用 RocketMQ 更顺,大数据生态用 Kafka 最自然,中小团队想要省心、文档多、社区活跃,RabbitMQ 是最稳妥的选择。Kafka 的运维成本不低,光集群规划、分区平衡就够喝一壶。
- 成本与基础设施:如果公司已经有一套成熟的消息集群(比如 Kafka),那就别轻易引一个新的中间件进来。新增中间件意味着新增一套监控、备份、告警体系,这部分的隐性成本经常被低估。
6. 项目落地经验:几个值得写进团队规范的实战细节
6.1 幂等设计必须前置,别等出了线上事故再补
幂等这个话题在第三章说了很多,这里想强调一个组织层面的经验:幂等设计要在项目建模阶段就设计好,而不是等重复消费事故出现以后再来补救。一旦线上已经产生重复数据,补幂等要涉及历史数据清洗、去重表初始化、双写在过渡期的一致性校验,成本是前置设计的十倍不止。
落地时做到两点:一是所有消息体里必须有一个业务唯一 ID 字段,保证可以追溯;二是所有消费逻辑统一走一个"消费基类",基类里封装去重、日志、异常捕获、重试上限等通用逻辑,业务开发只管写自己的处理函数。这样新业务接入时天然具备幂等能力,不会因为某个开发忘了处理而埋雷。
6.2 消息体设计:字段、序列化与兼容性
很多人忽略消息体设计,随手塞一个 JSON 字符串就发出去,等需求变更时哭都来不及。我的经验是:消息体里放业务数据的核心字段,不要放整个数据实体的大对象,更不要直接透传数据库表结构。对外部依赖强相关的数据要留版本号字段,后续字段新增删改时便于兼容。
序列化方案上,JSON 是绝大多数场景的首选,可读性好、跨语言友好。追求极致的性能,可以用 Protobuf。但无论选哪种,都要在消费端做异常兜底,反序列化失败的场景不能直接抛异常把线程卡死,而是记录原始消息并转人工处理。
6.3 监控报警:没有监控的消息队列是定时炸弹
消息队列一旦运行起来,没有监控的意识就很危险。消费积压、消费速率下降、大量重试、死信堆积,这些问题往往不是马上爆发的,而是积累一段时间后突然变成事故。至少要监控这些核心指标:堆积量(backlog 数量和消费延迟)、消费速率、ACK 失败次数、消息生产速率、Broker 磁盘和 CPU。
报警阈值要根据业务容忍度设定。比如消息延迟几分钟没关系,那可以设置 10 分钟阈值;如果是实时风控链路,可能要求秒级报警。相关实践是:除了阈值告警,最好加一个"零速率告警",消费速率连续 N 分钟为 0 也要报警——这种情况往往是消费者实例挂了,业务方甚至都没察觉。
6.4 发布变更:消费端和消息格式的兼容策略
最后想专门提一点发布环节的坑,很多项目是在升级时出问题的。当你改了消息体字段,或者改了消费逻辑,假设旧消费者还没全部下线,新消息格式被旧消费者读取到,轻则字段缺失、重则反序列化失败。稳妥做法是"兼容扩展":加字段是安全的,改字段名/删字段必须做双版本兼容,或者先上线新消费者消费新旧两种格式,再在下一个版本彻底去掉旧格式。
消费逻辑变更也要注意顺序。比如你要增加一条消息的消费逻辑,先让新旧代码并存,再逐步切流量,最后去掉旧逻辑。千万别直接一个灰度发布就把消费端代码换了,在消息中间件运作的分布式环境里,这种"小步快跑、逐步切换"的做法是保命的。
最后分享一个个人体会:消息队列项目里,真正费心的事情从来不是怎么发消息和收消息,而是怎么让消息在异常情况下依然被安全、有序、不重复地处理。刚才讲的这些基础知识和实战经验,如果能在项目初期就刻进团队的设计习惯里,后续会让你少熬很多夜。这个内容还有很多延伸方向,比如不同产品的部署调优、消息轨迹追踪链路,这次先把骨架搭好,后面再逐步深入。