Skip to content

消息队列:削峰填谷的利器

难度:高级 | 预计时间:50 分钟 | 前置:Spring Cloud


📖 秒杀系统的崩溃

故事背景

你负责的电商平台要做秒杀——10 万人在同一秒点击"抢购"。订单服务直接扛不住了——数据库连接池耗尽、CPU 100%、整个服务瘫痪,连正常下单的用户也影响了。运维紧急扩容了 20 台机器——但秒杀已经结束了。

架构师告诉你:下次用消息队列。不用让 10 万个请求直接打到订单服务——先把请求扔到队列里,订单服务按自己的节奏(比如每秒处理 1000 个)慢慢消费。这就是削峰填谷——不让洪峰冲垮系统,让流量平稳流过。消息队列是分布式系统里最常用的解耦和流量整形工具。

💻 代码演示

Spring Boot + RabbitMQ 秒杀场景——生产者(接收请求→发消息)vs 消费者(慢慢处理):

java
// ===== 生产者:Controller 收到秒杀请求,不直接处理,发消息 =====
@RestController
public class SeckillController {
    private final RabbitTemplate rabbitTemplate;

    public SeckillController(RabbitTemplate rabbitTemplate) {
        this.rabbitTemplate = rabbitTemplate;
    }

    @PostMapping("/seckill")
    public ResponseEntity<String> seckill(@RequestBody SeckillRequest req) {
        // 不直接操作数据库,只发消息——毫秒级响应
        rabbitTemplate.convertAndSend(
            "seckill.exchange",    // 交换机
            "seckill.order",       // 路由键
            req                     // 消息体(自动序列化为JSON)
        );
        return ResponseEntity.ok("抢购请求已接收,请等待结果通知");
    }
}

// ===== 消费者:慢慢处理,不受洪峰影响 =====
@Component
public class SeckillConsumer {

    @RabbitListener(queues = "seckill.queue")  // 监听队列
    public void handleSeckill(SeckillRequest req) {
        // 扣库存、创建订单、扣款...再慢也没关系,不影响前端响应
        if (inventoryService.deduct(req.productId(), req.quantity())) {
            orderService.create(req);
            // 处理完通知用户
            notifyService.send(req.userId(), "抢购成功!");
        } else {
            notifyService.send(req.userId(), "已售罄");
        }
    }
}

// application.yml 关键配置
// spring.rabbitmq.host=localhost
// spring.rabbitmq.port=5672
// spring.rabbitmq.username=guest
// spring.rabbitmq.password=guest

运行输出:

# 10万个并发请求 → Controller 秒回 "已接收" → 消息入队
# 订单服务按自己节奏(每秒1000个)从容消费 → 数据库压力平稳
# 即使订单服务短暂宕机,消息在队列里不丢失,重启后继续消费

🎯 三个核心问题

这是什么?

消息队列(MQ)是一种异步通信中间件——生产者发消息到队列,消费者从队列取消息,双方不需要同时在线。主流产品:RabbitMQ(轻量,适合业务解耦)、Kafka(高吞吐,适合日志/流处理)、RocketMQ(阿里系,适合交易场景)。核心概念:Producer(生产者)、Consumer(消费者)、Queue(队列)、Exchange(交换机,负责路由到队列)、Topic(主题,发布/订阅模式)。

为什么需要它?

MQ 解决三个核心问题:①削峰填谷——让请求排队,后端按能力消费(秒杀场景)②解耦——订单服务不需要知道通知服务、积分服务、日志服务的存在,发一条消息,谁关心谁订阅(发布/订阅)③异步——耗时操作(发邮件、生成报表)不用阻塞主流程,发消息后立刻响应用户。缺点:增加系统复杂度、消息延迟、一致性问题(消息丢了怎么办?)。

如果没有它会怎样?

如果没有消息队列——秒杀场景只能靠堆机器硬抗(成本高且不保证不崩);服务间耦合紧密(订单服务需要显式调用通知、积分、日志三个服务——一个挂了全部超时);所有非核心操作都阻塞用户请求(注册后要等发送验证邮件完成才返回"注册成功")。

📝 原理讲解

消息队列核心概念速查:

  • Producer → Exchange → Queue → Consumer:生产者发消息到 Exchange(交换机),Exchange 按路由规则转发到 Queue,消费者监听 Queue 获取消息。这和 Java 里的"发布-订阅"模式本质一样。

  • 消息可靠性三件套:①Producer 确认(Publisher Confirm——消息真的到队列了)②消息持久化(队列和消息都存磁盘——MQ 重启不丢)③Consumer 手动 ACK(处理完才确认——处理到一半宕机,消息重回队列)。

  • 死信队列(DLQ):消息被拒绝、过期、队列满了时,自动进入死信队列——人工排查为啥这条消息消费失败。

  • RabbitMQ vs Kafka:RabbitMQ 适合业务解耦(消息量适中、对延迟敏感、路由规则灵活)。Kafka 适合流式数据(海量日志、点击流、实时计算——每秒百万条消息,消息可回溯)。

  • 消费幂等性:MQ 可能重复投递消息(至少一次投递保证)。消费者必须做幂等处理——同一条消息处理多次结果相同(用唯一 ID 去重)。

🎨 生活类比

类比理解

消息队列像银行叫号系统——你取号(发消息)后可以在等候区玩手机(异步),柜台有空了就叫你(消费者处理)。你不用一直站柜台前等(同步阻塞),柜台也不会被一拥而上的人挤爆(削峰)。 Exchange 是邮局分拣中心——你寄信(发消息)不需要知道收件人在哪,分拣中心(Exchange)根据邮编(路由键)送到正确的邮筒(Queue)。 死信队列是"问题邮件"处理处——地址不详、无人签收的邮件(消费失败的消息)统一送到这里,人工检查出了什么问题。 Kafka 像是 24 小时监控录像——记录所有发生的事情(事件日志),你想回溯昨天的某个时段?拉录像带就行(消息持久化 + 可回溯消费)。

✏️ 动手练习

练习 1

用 Docker 启动 RabbitMQ(或本地安装),在 Spring Boot 项目里实现一个最简单的消息发送和接收:发一条"Hello MQ",消费者打印出来。

<details> <summary>💡 查看提示</summary>

Docker: docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-management。依赖 spring-boot-starter-amqp。用 RabbitTemplate 发送,@RabbitListener 接收。

</details>

练习 2

实现一个"订单创建通知"场景:订单服务创建订单后发消息,通知服务消费消息发短信,积分服务消费消息加积分。做到手动 ACK。

<details> <summary>💡 查看提示</summary>

两个消费者监听同一队列或不同队列。application.yml 里配 spring.rabbitmq.listener.simple.acknowledge-mode=manual,消费者方法参数加 Channel channel, Message message,手动调用 channel.basicAck()。

</details>

练习 3

模拟消费失败转入死信队列:设置队列的 x-dead-letter-exchange,消费者抛异常后不 ACK(或 reject),观察消息进入死信队列。

<details> <summary>💡 查看提示</summary>

在队列声明时加参数:x-dead-letter-exchange=dlx.exchange。消费者抛 AmqpRejectAndDontRequeueException 或 channel.basicNack() 不重新入队。

</details>

✅ 自检站

<details> <summary><strong>消息队列的"削峰填谷"具体是怎么做到的?</strong></summary>

生产者(Controller)收到请求后不发到后端处理,而是发一条消息到 MQ 就立刻响应用户——延迟毫秒级。消费者按自己的处理能力(比如每秒 1000 条)匀速消费——不管多少请求涌进来,后端处理的速率是恒定的。高峰时的消息积压在队列里,低谷时消费者慢慢消化——平滑了流量,保护了后端系统。

</details>

<details> <summary><strong>如何保证消息不丢失?从生产到消费需要哪些措施?</strong></summary>

①Producer 端:开启 Publisher Confirm,只有收到 Broker 确认才认为发送成功,否则重试 ②Broker 端:队列和消息都设置为持久化(durable + persistent),写入磁盘而非内存 ③Consumer 端:手动 ACK——处理完业务逻辑后才 ack,如果处理到一半宕机,消息未 ack 会重新投递给其他消费者。三者缺一不可。

</details>

<details> <summary><strong>为什么消费者需要幂等处理?如何实现?</strong></summary>

MQ 保证至少一次投递(at-least-once),消息可能被重复消费:网络抖动导致 ACK 丢失→消息重投、Consumer 宕机重启→重新消费。幂等方案:①数据库唯一约束(insert ignore / ON DUPLICATE KEY)②Redis 记录已处理的消息 ID(setnx)③业务状态机(已支付→不重复扣款)。原则:业务逻辑本身要对重复执行是安全的。

</details>