Skip to content
第 200 / 250 章中间件⏱ 12 分钟阅读

第 200 章:RabbitMQ 入门与实战

学习目标

  • 理解 RabbitMQ 核心概念
  • 掌握交换机与队列
  • 学会 Spring AMQP
  • 了解延迟、死信、优先级队列

一、RabbitMQ 简介

RabbitMQ 是Erlang 实现的 AMQP 协议消息中间件,适合复杂路由场景。

二、核心概念

概念说明
BrokerRabbitMQ 服务器
VHost虚拟主机,逻辑隔离
Exchange交换机,接收消息并路由
Queue队列,存消息
Binding交换机和队列的绑定规则
RoutingKey路由 key
Message消息(payload + properties)

三、交换机类型

3.1 Direct(直连)

精确匹配 RoutingKey。

exchange: order.direct
routing key "order.created" → queue: order-created
routing key "order.paid"    → queue: order-paid

3.2 Topic(主题,最灵活)

通配符匹配:

java
// 绑定
// queue: order.q   binding: order.*
// queue: created.q binding: *.created
// queue: log.q     binding: #

通配符:

  • *:匹配 1 个单词
  • #:匹配 0 或多个单词

3.3 Fanout(广播)

忽略 RoutingKey,绑定的所有队列都收到。

java
// 日志广播
// queue1 (error-log)  → binding: #
// queue2 (all-log)    → binding: #

3.4 Headers(少用)

按 headers 匹配,不灵活。

四、安装

yaml
# docker-compose.yml
services:
  rabbitmq:
    image: rabbitmq:3.12-management
    ports:
      - "5672:5672"
      - "15672:15672"   # 管理界面
    environment:
      RABBITMQ_DEFAULT_USER: admin
      RABBITMQ_DEFAULT_PASS: admin

访问 http://localhost:15672

五、Spring AMQP 集成

5.1 引入

xml
<dependency>
  <groupId>org.springframework.boot</groupId>
  <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

5.2 配置

yaml
spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: admin
    password: admin
    virtual-host: /
    listener:
      simple:
        acknowledge-mode: manual
        prefetch: 10
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 1000ms

5.3 声明交换机 / 队列

java
@Configuration
public class RabbitConfig {

    @Bean
    public TopicExchange orderExchange() {
        return new TopicExchange("order.exchange", true, false);
    }

    @Bean
    public Queue orderCreatedQueue() {
        return QueueBuilder.durable("order.created.q")
            .withArgument("x-dead-letter-exchange", "dlx.exchange")
            .withArgument("x-dead-letter-routing-key", "dlq")
            .build();
    }

    @Bean
    public Binding orderCreatedBinding(Queue orderCreatedQueue, TopicExchange orderExchange) {
        return BindingBuilder.bind(orderCreatedQueue)
            .to(orderExchange)
            .with("order.created");
    }
}

5.4 Producer

java
@Service
public class OrderProducer {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    public void sendOrderCreated(OrderEvent event) {
        rabbitTemplate.convertAndSend(
            "order.exchange",         // exchange
            "order.created",          // routing key
            event,
            message -> {
                message.getMessageProperties().setMessageId(UUID.randomUUID().toString());
                message.getMessageProperties().setContentType("application/json");
                return message;
            }
        );
    }
}

5.5 Consumer

java
@Component
public class OrderConsumer {

    @RabbitListener(queues = "order.created.q")
    public void handleOrder(OrderEvent event, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
        try {
            log.info("收到订单事件: {}", event);
            // 业务处理
            inventoryService.deduct(event);

            channel.basicAck(tag, false);   // 确认
        } catch (Exception e) {
            log.error("处理失败", e);
            channel.basicNack(tag, false, false);  // 拒绝 + 不重回队列
        }
    }
}

六、消息可靠性

6.1 生产者确认

yaml
spring:
  rabbitmq:
    publisher-confirm-type: correlated   # 开启确认
    publisher-returns: true              # 开启退回
java
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory factory) {
    RabbitTemplate template = new RabbitTemplate(factory);
    template.setConfirmCallback((data, ack, cause) -> {
        if (!ack) {
            log.error("发送失败: {}", cause);
        }
    });
    template.setReturnsCallback(returned -> {
        log.error("消息退回: {}", returned.getMessage());
    });
    return template;
}

6.2 消费者 ACK

yaml
acknowledge-mode: manual
java
// 手动 ACK
channel.basicAck(tag, false);

// 拒绝(重回队列)
channel.basicNack(tag, false, true);

// 拒绝(不重回队列)
channel.basicNack(tag, false, false);

6.3 持久化

java
QueueBuilder.durable("xxx").build();      // 队列持久化
MessageDeliveryMode.PERSISTENT             // 消息持久化

七、死信队列(DLQ)

消费失败 / 过期 / 队列满的消息自动进入 DLQ。

java
@Bean
public Queue orderQueue() {
    return QueueBuilder.durable("order.q")
        .deadLetterExchange("dlx.exchange")
        .deadLetterRoutingKey("dlq")
        .ttl(60_000)
        .build();
}

八、延迟队列

8.1 TTL + DLX(经典方案)

java
// 延迟队列:不直接消费,只让消息过期后转 DLX
@Bean
public Queue delayQueue() {
    return QueueBuilder.durable("order.delay.q")
        .deadLetterExchange("order.exchange")
        .deadLetterRoutingKey("order.created")
        .ttl(30 * 60 * 1000)   // 30 分钟
        .build();
}

// 发到延迟队列,过期后自动转到业务队列
rabbitTemplate.convertAndSend("", "order.delay.q", event);

8.2 插件版(RabbitMQ Delayed Plugin)

bash
# 安装插件
rabbitmq-plugins enable rabbitmq_delayed_message_exchange
java
@Bean
public CustomExchange delayExchange() {
    Map<String, Object> args = new HashMap<>();
    args.put("x-delayed-type", "direct");
    return new CustomExchange("delay.exchange", "x-delayed-message", true, false, args);
}
java
// 发送时指定延迟(毫秒)
rabbitTemplate.convertAndSend("delay.exchange", "order.created", event, message -> {
    message.getMessageProperties().setHeader("x-delay", 30 * 60 * 1000);
    return message;
});

九、优先级队列

java
@Bean
public Queue priorityQueue() {
    return QueueBuilder.durable("priority.q")
        .maxPriority(10)        // 1-10,数字越大越高
        .build();
}
java
message.getMessageProperties().setPriority(9);

十、消息幂等

java
@RabbitListener(queues = "order.created.q")
public void handle(OrderEvent event, Channel channel,
                   @Header(AmqpHeaders.DELIVERY_TAG) long tag) {
    try {
        String dedupKey = "rabbit:order:" + event.getOrderId();
        if (!redisTemplate.opsForValue().setIfAbsent(dedupKey, "1", 24, TimeUnit.HOURS)) {
            channel.basicAck(tag, false);    // 重复消息也 ACK
            return;
        }
        inventoryService.deduct(event);
        channel.basicAck(tag, false);
    } catch (Exception e) {
        channel.basicNack(tag, false, true);
    }
}

十一、Kafka vs RabbitMQ

维度KafkaRabbitMQ
吞吐量百万级/s万级/s
延迟10ms+us 级
路由简单(分区内)丰富(4 种交换机)
消息回溯支持(保留期)不支持(消费即删除)
适用日志、事件流业务消息、RPC

十二、本章小结

概念用途
Exchange路由
Queue存储
Binding绑定关系
RoutingKey路由 key
交换机场景
Direct精确路由
Topic通配符路由
Fanout广播

动手练习

  1. 用 Spring AMQP 实现订单发布/订阅
  2. 用 TTL + DLX 实现订单 30 分钟自动关闭
  3. 用 Direct 实现错误日志分类
  4. 用 Topic 实现多服务订阅同一消息

推荐阅读


下一章:第 201 章:MQ 选型对比与场景

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