关于消息队列架构的几点记录

如果任何一个下游服务超时或失败,整个下单流程就卡住了。

场景三:跨系统对接 内部系统要对接第三方支付、物流、风控,第三方接口的响应时间不可控,甚至可能间歇性失败。

为什么需要消息队列

刚开始做微服务的时候,我也觉得直接 HTTP 调用就够了:A 服务调 B 服务,B 服务调 C 服务,链路清晰、调试简单。直到遇到这几个场景:

场景一:电商下单 用户下单成功后,需要同时做:扣库存、生成订单、发短信通知、推送物流系统、记录行为数据。如果任何一个下游服务超时或失败,整个下单流程就卡住了。更糟的是,短信接口慢 5 秒,用户就得等 5 秒。

场景二:秒杀活动 几万人同时抢购,流量瞬间打爆。你的下单服务能抗住,但下游的积分、营销服务扛不住。直接同步调用就是雪崩的前兆。

场景三:跨系统对接 内部系统要对接第三方支付、物流、风控,第三方接口的响应时间不可控,甚至可能间歇性失败。同步调用意味着你的系统稳定性被别人牵着鼻子走。

这时候你才会意识到:有些事情不需要马上做完,有些事情不应该阻塞主流程,有些事情必须允许失败重试。消息队列就是用来解决这些"异步"问题的。

消息队列能做什么

异步解耦

把同步调用改成异步通知,主流程只做最核心的事,剩下的扔给 MQ 慢慢消化。

# 同步调用:下单耗时 = 库存 + 订单 + 短信 + 物流 + 数据
# 理想情况:500ms,最坏情况:5s

# 异步调用:下单耗时 = 库存 + 订单 + 发消息
# 稳定在 200ms,后续流程异步处理

好处很直接:用户体验变好了,系统吞吐量上去了,下游服务挂了不影响主流程。

但代价也很明显:消息不是立刻处理的,用户需要接受"最终一致性";消息可能丢失,需要考虑持久化和重试;消息可能乱序,需要业务容忍或做额外处理。

削峰填谷

把突发的流量先存起来,下游按照自己的节奏慢慢消化,像水库一样。

sequenceDiagram participant User as 用户请求 participant MQ as 消息队列 participant Service as 下游服务 User->>MQ: 瞬间 10,000 请求 MQ->>MQ: 缓冲在队列中 MQ->>Service: 按能力 1,000/s 消费 Service->>Service: 稳定处理

关键点在"缓冲":MQ 本身是高可用的,可以抗住突发流量;下游服务只关心自己的处理能力,不需要为极端场景过度扩容。

但要注意:缓冲不是无限大,如果消费速度持续跟不上,队列就会堆积,最终触发背压或丢消息策略。

数据分发

一个消息被多个消费者消费,典型场景是数据同步、日志采集、事件驱动。

比如订单创建后,需要同时同步到:数据仓库、搜索索引、缓存、报表系统。每个下游独立订阅 MQ,互不影响。

常见 MQ 选型对比

RabbitMQ

基于 AMQP 协议的传统 MQ,特点是稳定可靠、功能齐全。

适合场景

  • 对可靠性要求极高的核心业务(支付、订单)
  • 需要复杂路由和死信队列的场景
  • 团队熟悉 AMQP 协议和运维

优点

  • 成熟稳定,生产经验丰富
  • 支持多种消息模式和路由规则
  • 管理界面友好,运维方便

缺点

  • 吞吐量相对较低(单机几万 TPS)
  • 水平扩展需要集群,运维成本较高
  • 依赖 Erlang,技术栈相对小众

Kafka

基于日志的分布式消息系统,特点是高吞吐、低延迟、可扩展。

适合场景

  • 大数据流处理、日志采集
  • 需要极高吞吐量的场景(百万级 TPS)
  • 消息回溯、重放需求

优点

  • 吞吐量极高,可水平扩展
  • 消息持久化,支持历史消息回放
  • 生态完善,对接大数据组件方便

缺点

  • 功能相对简单,高级特性需要自己实现
  • 实时性不如传统 MQ(毫秒级 vs 微秒级)
  • 运维复杂度较高(ZooKeeper、Broker、Consumer Group)

其他选择

RocketMQ:阿里开源,适合国内业务场景,事务消息、定时消息支持好。

Pulsar:云原生设计,计算存储分离,适合多租户场景。

Redis Stream:轻量级消息队列,适合小规模、对可靠性要求不高的场景。

实践中的坑

消息丢失

第一次在生产环境用 MQ,就遇到了消息丢失的问题。排查半天发现是自动提交 offset 的锅:消息拉下来了,还没处理完就 crash,offset 已经提交了,消息就丢了。

解决

  • 发送端开启 confirm 机制,确保消息到达 broker
  • 消费端关闭自动提交,处理完再手动提交
  • 核心业务考虑数据库 + MQ 本地消息表的方案

重复消费

MQ 至少一次投递的语义意味着重复消费不可避免。下游服务如果不幂等,就会出问题。

解决

  • 唯一 ID + 数据库唯一索引
  • 分布式锁(Redis)
  • 状态机控制,幂等更新

消息积压

某次下游服务升级,消费者全部挂掉,几小时后发现 MQ 堆了千万级消息。

解决

  • 监控消费延迟,及时告警
  • 临时扩容消费者(注意 rebalance 影响)
  • 降级或丢弃非核心消息
  • Kafka 可以临时提高消费并发

顺序问题

同一个用户的先后两个请求,因为并行消费,可能先处理了后面的。

解决

  • 同一个 key 的消息发送到同一个 partition
  • 单线程消费或本地有序队列
  • 业务上接受最终一致性

什么时候不该用 MQ

MQ 不是万能的,有时候不用反而更简单:

数据一致性要求极高:库存扣减这类操作,数据库事务更直接。

低延迟场景:实时通信、高频交易,MQ 的延迟不可接受。

简单业务:就两个服务之间的简单调用,HTTP 够用,引入 MQ 增加复杂度。

团队规模小:维护一套 MQ 集群需要投入人力,小团队可能吃不消。

写在最后

消息队列是分布式系统里的好工具,但它解决的是异步问题,不是所有问题。用之前先问自己三个问题:

  1. 这件事必须立刻做完吗?
  2. 如果做慢了或失败了,系统能接受吗?
  3. 团队有能力维护这套基础设施吗?

回答清楚了再决定要不要上 MQ。上了之后,记得盯着监控,别等到积压了才发现。


这次折腾让我明白:工具的价值不在技术本身,而在它解决的实际问题。MQ 如此,其他技术亦然。

版权声明: 本文首发于 指尖魔法屋-关于消息队列架构的几点记录https://blog.thinkmoon.cn/post/49-mq-architecture-kafka-rabbitmq/) 转载或引用必须申明原指尖魔法屋来源及源地址!