Skip to content
第 62 / 250 章后端⏱ 10 分钟阅读

第 62 章:消息队列与 RabbitMQ

学习目标

  • 理解 MQ 的核心价值:异步、解耦、削峰
  • 掌握 RabbitMQ 五大模式与 Spring 集成
  • 学会可靠投递、幂等消费、死信队列

一、为什么需要 MQ?

MQ 的三大核心价值

价值含义
异步写完就走,下游慢慢处理
解耦服务之间通过消息通信,不依赖接口契约
削峰流量高峰先堆积在 MQ,系统按能力消费

二、MQ 选型对比

维度RabbitMQRocketMQKafka
吞吐量万级十万级百万级
延迟微秒级毫秒级毫秒级
消息可靠性中(默认可能丢)
事务消息弱(要事务机制)强(原生支持)
生态成熟阿里系大数据、流计算
适用企业业务、可靠投递电商、金融(阿里主推)日志、流计算

三、RabbitMQ 核心概念

概念说明
Producer消息生产者
Consumer消息消费者
Exchange交换机,接收生产者的消息并路由
Queue队列,真正存储消息的地方
Binding交换机和队列之间的绑定规则
Routing Key路由键,交换机用它决定消息去哪个队列
VHost虚拟主机,多租户隔离

四、Spring Boot 集成

xml
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
yaml
spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: admin
    password: ${RABBITMQ_PASSWORD}
    virtual-host: /taskflow
    publisher-confirm-type: correlated       # ① 发送确认
    publisher-returns: true                  # ② 路由失败回调
    listener:
      simple:
        acknowledge-mode: manual             # ③ 手动 ack
        prefetch: 10                         # ④ 每次拉取 10 条
        concurrency: 4                       # ⑤ 4 个消费者线程
        max-concurrency: 8
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 1s

五、声明交换机和队列

java
@Configuration
public class RabbitConfig {

    // ① 业务交换机(Topic 模式)
    @Bean
    public TopicExchange orderExchange() {
        return ExchangeBuilder.topicExchange("order.exchange")
                .durable(true)
                .build();
    }

    // ② 下单队列
    @Bean
    public Queue orderQueue() {
        return QueueBuilder.durable("order.queue")
                .deadLetterExchange("dlx.exchange")              // 死信交换机
                .deadLetterRoutingKey("order.dlq")
                .ttl(24 * 60 * 60 * 1000)                        // 24 小时过期
                .maxLength(100000)                                // 最大消息数
                .build();
    }

    // ③ 绑定:order.* 路由到 order.queue
    @Bean
    public Binding orderBinding() {
        return BindingBuilder.bind(orderQueue())
                .to(orderExchange())
                .with("order.#");
    }

    // ④ 死信交换机
    @Bean
    public DirectExchange dlxExchange() {
        return ExchangeBuilder.directExchange("dlx.exchange").build();
    }

    @Bean
    public Queue dlq() {
        return QueueBuilder.durable("order.dlq").build();
    }

    @Bean
    public Binding dlqBinding() {
        return BindingBuilder.bind(dlq()).to(dlxExchange()).with("order.dlq");
    }
}

六、生产者:可靠发送

java
@Component
@RequiredArgsConstructor
@Slf4j
public class OrderProducer {

    private final RabbitTemplate rabbitTemplate;

    // ① 发送消息(带 ConfirmCallback)
    public void send(OrderMessage message) {
        // ② 业务唯一 ID(用于幂等)
        String msgId = UUID.randomUUID().toString();

        // ③ 构造消息
        Message msg = MessageBuilder
                .withBody(JSON.toJSONBytes(message))
                .setContentType(MessageProperties.CONTENT_TYPE_JSON)
                .setMessageId(msgId)                             // 幂等 key
                .setCorrelationId(message.getOrderId().toString())
                .build();

        // ④ 设置 ConfirmCallback(发送成功回调)
        rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
            if (ack) {
                log.info("消息发送成功 id={}", msgId);
            } else {
                log.error("消息发送失败 id={} cause={}", msgId, cause);
                // 重试或入库
            }
        });

        // ⑤ 设置 ReturnsCallback(路由失败回调)
        rabbitTemplate.setReturnsCallback(returned ->
                log.error("消息路由失败 exchange={} key={}",
                        returned.getExchange(), returned.getRoutingKey()));

        // ⑥ 发送
        CorrelationData corr = new CorrelationData(msgId);
        rabbitTemplate.convertAndSend("order.exchange", "order.created", msg, corr);
    }
}

本地消息表 + MQ(生产级可靠投递)

java
@Transactional
public void placeOrder(OrderDTO dto) {
    // ① 业务落库
    orderMapper.insert(order);

    // ② 同一事务里插入本地消息表
    LocalMessage msg = new LocalMessage();
    msg.setId(UUID.randomUUID().toString());
    msg.setTopic("order.created");
    msg.setPayload(JSON.toJSONString(order));
    msg.setStatus(0);                                 // 待发送
    localMessageMapper.insert(msg);
}

// ③ 后台任务:扫表发送
@Scheduled(fixedDelay = 1000)
public void sendPendingMessages() {
    List<LocalMessage> msgs = localMessageMapper.selectPending(100);
    for (LocalMessage msg : msgs) {
        try {
            rabbitTemplate.convertAndSend("order.exchange",
                    "order.created", msg.getPayload());
            msg.setStatus(1);                          // 已发送
            localMessageMapper.updateById(msg);
        } catch (Exception e) {
            log.warn("消息发送失败,稍后重试 id={}", msg.getId());
        }
    }
}

七、消费者:手动 ack + 幂等

java
@Component
@RequiredArgsConstructor
@Slf4j
@RocketMQMessageListener  // 伪注解示意,实际用 @RabbitListener
public class OrderConsumer {

    private final InventoryService inventoryService;
    private final RedisTemplate<String, Object> redis;

    // ① 监听队列
    @RabbitListener(queues = "order.queue")
    public void onMessage(Message message, Channel channel) throws IOException {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();
        String msgId = message.getMessageProperties().getMessageId();

        try {
            // ② 幂等检查(关键!)
            String key = "mq:consumed:" + msgId;
            Boolean firstTime = redis.opsForValue().setIfAbsent(key, "1", 24, TimeUnit.HOURS);
            if (Boolean.FALSE.equals(firstTime)) {
                log.warn("重复消息,直接 ack msgId={}", msgId);
                channel.basicAck(deliveryTag, false);    // 重复消息直接 ack
                return;
            }

            // ③ 业务处理
            OrderMessage order = JSON.parseObject(message.getBody(), OrderMessage.class);
            inventoryService.deduct(order.getSkuId(), order.getQuantity());

            // ④ 手动 ack
            channel.basicAck(deliveryTag, false);
            log.info("消息处理成功 msgId={}", msgId);

        } catch (BusinessException e) {
            // ⑤ 业务异常:消息重试(可入死信)
            log.warn("业务异常,重试 msgId={}", msgId, e);
            channel.basicNack(deliveryTag, false, true);

        } catch (Exception e) {
            // ⑥ 系统异常:消息重试或入死信
            log.error("系统异常 msgId={}", msgId, e);
            channel.basicNack(deliveryTag, false, false);   // 不重试 → 进入死信
        }
    }
}

幂等的几种实现

java
// 方案 1:Redis SETNX(最常用)
Boolean ok = redis.opsForValue().setIfAbsent(
        "mq:msg:" + msgId, "1", 24, TimeUnit.HOURS);
if (!ok) return;    // 已消费过

// 方案 2:数据库唯一索引
// 消息处理表加唯一约束 msg_id,重复消息 INSERT 失败即视为已处理
INSERT INTO mq_consumed_log (msg_id, create_time) VALUES (#{msgId}, NOW());
// DuplicateKeyException → 已处理

// 方案 3:业务判断
// 订单已处理过 status=PAID,跳过
if (order.getStatus() == PAID) return;

八、五大消息模式

1. 简单模式(Hello World)

java
// 生产者
rabbitTemplate.convertAndSend("hello.queue", "Hello World");

// 消费者
@RabbitListener(queues = "hello.queue")
public void onMessage(String msg) {
    System.out.println("收到:" + msg);
}

2. Work Queue(任务队列)

多个消费者分摊任务(轮询 / 公平分发)

3. Publish/Subscribe(发布订阅)

java
// Fanout 交换机:忽略 routing key,广播到所有绑定队列
@Bean
public FanoutExchange fanout() { return new FanoutExchange("notify.exchange"); }

4. Routing(路由)

java
// Direct 交换机:完全匹配 routing key
rabbitTemplate.convertAndSend("log.exchange", "error", "数据库连接失败");

5. Topic(主题,最常用)

java
// Topic 交换机:通配符匹配
// * 匹配一个单词,# 匹配多个单词

rabbitTemplate.convertAndSend("order.exchange", "order.vip.created", msg);   // 匹配 # 和 *
rabbitTemplate.convertAndSend("order.exchange", "order.normal.paid", msg);    // 只匹配 #

九、死信队列(DLX)

死信:消息被拒绝、过期、队列满时的归宿。

java
// 声明队列时指定 DLX
@Bean
public Queue orderQueue() {
    return QueueBuilder.durable("order.queue")
            .deadLetterExchange("dlx.exchange")
            .deadLetterRoutingKey("order.dlq")
            .build();
}
java
@Component
@Slf4j
public class DlqConsumer {

    @RabbitListener(queues = "order.dlq")
    public void onDead(Message msg) {
        log.error("死信消息: msgId={} body={}",
                msg.getMessageProperties().getMessageId(),
                new String(msg.getBody()));
        // 告警 / 入库 / 人工处理
    }
}

十、消息可靠性保障总结

阶段风险保障
生产者发送网络抖动、MQ 宕机ConfirmCallback + 本地消息表
MQ 存储宕机丢消息队列持久化 + 镜像队列
消费者接收消息丢失手动 ack + 失败重试
消费处理重复消费幂等(Redis SETNX / DB 唯一索引)

十一、本章小结

要点关键
MQ 价值异步、解耦、削峰
选型RabbitMQ(企业业务)/ RocketMQ(阿里电商)/ Kafka(日志流)
五种模式Hello / Work / Fanout / Direct / Topic
可靠发送ConfirmCallback + 本地消息表
可靠消费手动 ack + 失败重试
幂等Redis SETNX / DB 唯一索引
死信DLX 自动收容失败消息

动手练习

练习 1:基础题

实现订单创建 → 通知库存扣减的简单流程:生产者发送 order.created,消费者扣减库存并打印日志。

练习 2:进阶题

为消费者加上:手动 ack + 幂等(Redis SETNX)+ 失败入死信队列 + 死信告警。模拟消费者抛异常,验证消息最终进入死信队列。

练习 3:思考题

设计一个「用户注册成功后发送欢迎邮件」的功能。用 MQ 异步发送,对比同步发送的优缺点。


下一章第 63 章:Kafka 与流处理

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