Skip to content
第 12 章 ⏱ 14 分钟阅读

第 12 章:RabbitMQ 进阶 ​

学习目标 ​

  • 掌握 Exchange 类型和路由机制
  • 实现生产者/消费者 Java 代码
  • 处理消息可靠性、幂等、死信
  • 了解 RabbitMQ 集群

一、RabbitMQ 核心概念 ​

概念含义
Producer消息生产者
Consumer消息消费者
Exchange交换机,接收消息并路由
Queue队列,存储消息
Binding交换机和队列的绑定规则
Routing Key路由键
Vhost虚拟主机(隔离命名空间)

二、Exchange 类型 ​

2.1 Direct(精确路由) ​

text
Routing Key = "order.created"  → 绑定了 "order.created" 的 Queue 收到

2.2 Topic(模式匹配) ​

text
Routing Key = "order.paid.cn"
Binding Pattern:
  "order.*"      → 匹配("*" = 一个单词)
  "order.#"      → 匹配("#" = 零或多个单词)

2.3 Fanout(广播) ​

无视 Routing Key,所有绑定的 Queue 都收一份。

2.4 Headers(键值匹配) ​

少用,性能差。

三、安装 ​

bash
docker run -d --name rabbitmq \
  -p 5672:5672 -p 15672:15672 \
  -e RABBITMQ_DEFAULT_USER=admin \
  -e RABBITMQ_DEFAULT_PASS=admin123 \
  rabbitmq:3-management

访问 http://localhost:15672(管理界面)。

⚠️ 坑 1:guest 用户默认只能 localhost 登录,Docker 启动必须改 RABBITMQ_DEFAULT_USER。

四、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: admin123
    virtual-host: /              # ⚠️ vhost 必须已存在

声明 Exchange / Queue ​

java
@Configuration
public class RabbitConfig {

    public static final String EXCHANGE = "order.exchange";
    public static final String QUEUE = "order.queue";
    public static final String ROUTING_KEY = "order.created";

    @Bean
    public TopicExchange exchange() {
        return new TopicExchange(EXCHANGE);
    }

    @Bean
    public Queue queue() {
        return QueueBuilder.durable(QUEUE).build();   // durable = 重启不丢
    }

    @Bean
    public Binding binding() {
        return BindingBuilder.bind(queue()).to(exchange()).with("order.*");
    }
}

⚠️ 坑 2:Queue 默认 durable=true 但不持久化消息;重启 queue 在,消息丢。要 durable + 消息发时 MessageDeliveryMode.PERSISTENT。

五、生产者 ​

java
@Service
public class OrderProducer {

    @Resource
    private RabbitTemplate rabbitTemplate;

    public void send(Long orderId) {
        rabbitTemplate.convertAndSend(
            RabbitConfig.EXCHANGE,
            "order.created",                          // routing key
            orderId,
            msg -> {
                msg.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
                msg.getMessageProperties().setMessageId(UUID.randomUUID().toString());
                return msg;
            }
        );
    }
}

六、消费者 ​

java
@Component
public class OrderConsumer {

    @RabbitListener(queues = RabbitConfig.QUEUE)
    public void handle(Long orderId, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag)
            throws IOException {
        try {
            log.info("处理订单:{}", orderId);
            // 业务逻辑
            channel.basicAck(tag, false);              // 手动 ack
        } catch (Exception e) {
            channel.basicNack(tag, false, false);       // 不重试,直接进死信
        }
    }
}

⚠️ 坑 3:@RabbitListener 默认是 auto-ack,消息发出就认为成功 —— 关掉 auto-ack,业务成功才 basicAck。 ⚠️ 坑 4:@RabbitListener 方法体里抛异常,默认会无限重试,需要配 retryTemplate 或死信队列。

七、消息可靠性 ​

7.1 三种保护 ​

阶段保证手段
生产发送publisher-confirm 异步确认
队列存储队列和消息都 durable + PERSISTENT
消费处理手动 ack + 重试机制

7.2 开启发布确认 ​

yaml
spring:
  rabbitmq:
    publisher-confirm-type: correlated    # 异步回调
    publisher-returns: true
java
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory cf) {
    RabbitTemplate tpl = new RabbitTemplate(cf);
    tpl.setConfirmCallback((data, ack, cause) -> {
        if (!ack) log.error("消息发送失败:{}", cause);
    });
    return tpl;
}

八、幂等消费 ​

RabbitMQ 至少一次(重复消费不可避免),业务必须幂等。

java
// 用 Redis 判重
public boolean isProcessed(String messageId) {
    Boolean first = redis.opsForValue()
        .setIfAbsent("msg:" + messageId, "1", Duration.ofDays(1));
    return Boolean.FALSE.equals(first);   // 已存在 = 重复
}

⚠️ 坑 5:幂等 token 必须全局唯一(用 messageId 或业务 ID),Redis 用 SETNX 保证原子。

九、死信队列(DLX) ​

消息被 nack、rejected、过期时,会进死信交换机,再绑个队列专门收死信:

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

死信处理 = 人工排查 + 重新入队。

十、本章小结 ​

要点关键
ExchangeDirect / Topic / Fanout / Headers
可靠性持久化 + 手动 ack + 发布确认
幂等Redis SETNX 判重,messageId 全局唯一
死信DLX + 死信队列,排查和重投
retryx-dead-letter-exchange + retry-interval
Listener默认无限重试,要配 retry 次数

动手练习 ​

  1. 用 Docker 启动 RabbitMQ,声明 Topic Exchange,绑定 order.* 队列
  2. 生产者发 "order.created"、"order.paid.cn",消费者只收 order.*
  3. 故意抛异常,观察死信队列收到消息

下一章:第 13 章:Kafka 进阶 →

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