Skip to content
第 54 章 后端 ⏱ 10 分钟阅读

第 54 章:消息队列 ​

学习目标 ​

  • 知道为啥用 MQ
  • 掌握 RabbitMQ 发送 / 消费
  • 掌握 Kafka 入门

一、为啥用 MQ ​

解耦 + 异步 + 削峰:

下单 → 同步:扣库存 + 发短信 + 推送,3 个调用 = 300ms
下单 → 异步:扣库存(同步) + 发消息(MQ),1 个调用 = 100ms
                 ↓
          短信服务 / 推送服务 异步消费

二、RabbitMQ ​

2.1 集成 ​

xml
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
yaml
spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest

2.2 核心概念 ​

概念含义
Producer消息生产者
Consumer消息消费者
Queue队列,存消息
Exchange交换机,路由消息
RoutingKey路由 key
Binding交换机和队列的绑定规则

2.3 三种交换机 ​

交换机路由规则
directRoutingKey 完全匹配
topic* 匹配一个单词,# 匹配多个
fanout广播,所有绑定的队列都收

2.4 配置 ​

java
@Configuration
public class RabbitMQConfig {

    @Bean
    public Queue orderQueue() {
        return QueueBuilder.durable("order.queue").build();      // 持久化
    }

    @Bean
    public DirectExchange orderExchange() {
        return new DirectExchange("order.exchange");
    }

    @Bean
    public Binding binding(Queue orderQueue, DirectExchange orderExchange) {
        return BindingBuilder.bind(orderQueue).to(orderExchange).with("order.create");
    }
}

💡 三个 bean 各管什么:

  • Queue = 信箱(消息暂存地)
  • DirectExchange = 邮局分拣机(一个业务通常只建一个)
  • Binding = 分拣规则("RoutingKey = X 的信 → 投到 Y 信箱")

生产者和消费者互不认识:

  • 🟢 生产者:rabbitTemplate.convertAndSend("exchange 名", "RoutingKey", 内容) — 只管扔
  • 🔵 消费者:@RabbitListener(queues = "队列名") — 只管收

中间靠 Exchange + Binding 自动路由。加新业务只需加一个 Queue + Binding,生产者代码完全不用改。

💡 简单场景可以省略 Exchange + Binding:RabbitMQ 有一个默认交换机 ""(自动存在),所有队列自动绑在它上面,且路由 key = 自己的队列名。

java
// 只配 Queue,没配 Exchange / Binding
@Bean
public Queue smsQueue() {
    return QueueBuilder.durable("sms.queue").build();
}
java
// 生产者:1 个 String 参数,Spring 帮你翻译成 (默认交换机, "sms.queue" 作为路由 key)
rabbitTemplate.convertAndSend("sms.queue", msg);   // ✅ 找到 sms.queue

但生产建议显式配 Exchange + Binding:加新业务、加路由规则都清晰,排查问题方便。只有"1 个业务 1 个队列"的玩具场景才用默认交换机。

2.5 发送 ​

java
@Service
@RequiredArgsConstructor
public class OrderProducer {
    private final RabbitTemplate rabbitTemplate;

    public void sendOrderCreated(Order order) {
        rabbitTemplate.convertAndSend(
            "order.exchange",                  // 交换机
            "order.create",                    // 路由 key
            order                              // 消息体
        );
    }
}

2.6 消费 ​

java
@Component
public class OrderConsumer {

    @RabbitListener(queues = "order.queue")
    public void handle(Order order) {
        log.info("收到订单: {}", order);
        try {
            // 业务逻辑
        } catch (Exception e) {
            log.error("消费失败", e);
            throw e;        // 抛出 → 进入重试 / 死信队列
        }
    }
}

2.7 消息可靠性 ​

💡 核心问题:消息会丢。RabbitMQ 用两层确认防丢:生产确认(生产者→broker)和消费者 ACK(broker→消费者)。默认两层都不开,消息丢了都不知道。

💡 broker 是什么? = 装在服务器上的 RabbitMQ 服务进程(docker run rabbitmq 启起来的那玩意儿),生产者不直接找消费者,而是"扔给 broker,broker 转给消费者"。broker 就是 RabbitMQ、Kafka、RocketMQ 等消息系统的统称。

java
@RabbitListener(queues = "order.queue")
public void handle(
        Order order,                                    // ① 消息体(自动反序列化的对象)
        Channel channel,                                // ② AMQP 通道,用来 ACK/NACK
        @Header(AmqpHeaders.DELIVERY_TAG) long tag) {   // ③ 消息编号,告诉 broker 回的是哪条
    try {
        process(order);                                 // ④ 真正的业务逻辑
        channel.basicAck(tag, false);                   // ⑤ 手动 ACK:处理成功,broker 可以删消息
    } catch (Exception e) {
        channel.basicNack(tag, false, false);           // ⑥ 处理失败:第 1 个 false=不批量,第 2 个 false=不重回队列,进死信
    }
}

API 拆解:

API含义
basicAck(tag, false)tag = 这条消息的编号;第二个 false = 单条确认(不批量)
basicNack(tag, false, false)三个参数:①消息编号 ②是否批量 ③是否重新入队,false = 不重试,进死信队列;true = 重新入队再来一遍

acknowledge-mode: manual 必须配:不配这个,框架默认自动 ACK,你写的 basicAck/basicNack 就不生效。

yaml
spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: manual          # ① 关键!手动 ACK 才需要写 basicAck/basicNack
        prefetch: 10                      # ② 一次预取 10 条消息并发处理,默认 250
        retry:
          enabled: true                   # ③ 失败时本地重试 3 次
          max-attempts: 3
          initial-interval: 2s

⚠️ 坑 1:生产者也要开启确认(publisher-confirm-type=correlated),否则 convertAndSend 一调用就返回,你以为消息送到了,其实只是字节流推到 TCP——如果 RabbitMQ 在收到前宕机,消息就丢了。

💡 现实建议:绝大多数项目用默认自动 ACK + 简单重试配置就够(代码里不用写 Channel 那段),本节是严格不丢(订单、支付)的进阶方案,知道有这套即可。

三、Kafka ​

3.1 集成 ​

xml
<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>
yaml
spring:
  kafka:
    bootstrap-servers: localhost:9092
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      acks: all                  # 全部副本都确认
      retries: 3
    consumer:
      group-id: my-app
      key-deserializer: ...
      value-deserializer: ...
      auto-offset-reset: earliest

3.2 发送 / 消费 ​

java
@Service
@RequiredArgsConstructor
public class OrderEventProducer {
    private final KafkaTemplate<String, Order> kafkaTemplate;

    public void send(Order order) {
        kafkaTemplate.send("order-topic", order.getId().toString(), order);
    }
}

@Component
public class OrderEventConsumer {

    @KafkaListener(topics = "order-topic")
    public void handle(Order order) {
        log.info("收到: {}", order);
    }
}

3.3 Kafka vs RabbitMQ ​

场景选
业务消息(订单、通知)RabbitMQ
日志 / 大数据流Kafka
严格顺序Kafka
复杂路由RabbitMQ

四、本章小结 ​

要点关键
RabbitMQ业务消息,可靠投递
Kafka日志 / 流处理,高吞吐
可靠性生产确认 + 消费手动 ACK
重试配置最大重试,失败进死信

动手练习 ​

  1. 用 RabbitMQ 实现"下单后异步发短信"
  2. 用 @RabbitListener 消费消息,故意抛异常看重试

下一章:第 55 章:Spring Security + JWT →

本站基于 VitePress 构建 · 由 StackHub 团队维护