Skip to content

RabbitMQ 入门:Exchange、路由与可靠投递

本文是消息队列系统学习系列的 L1 入门篇。前置:25. 为什么需要 MQ:解耦、削峰与三大核心问题。 学完可以配合面试题食用:11-rabbitmq-architecture-exchange-binding-queue12-rabbitmq-confirm-return-reliable-delivery

为什么是 RabbitMQ

Kafka 是日志流的王者,RocketMQ 是电商业务的首选,那 RabbitMQ 呢?它擅长的场景恰好是 Kafka 和 RocketMQ 都不太舒服的地方:灵活的路由逻辑轻量级运维。一个订单系统里,同一个消息需要按不同规则分发到多个下游——物流组只收已发货通知、风控组收所有异常状态、审计组收全部订单快照。用 Kafka 做这种路由需要写一堆 Stream 处理逻辑,而 RabbitMQ 通过 Exchange 加 Binding 就能原生搞定。

RabbitMQ 实现了 AMQP 0-9-1 协议,核心思路和 Kafka 的 Topic 模型完全不同:消息先到 Exchange,路由规则由 Exchange 和 Binding 决定,Queue 只是消息的最终落脚点。生产者不直接写 Queue,它只发消息到 Exchange。

下面直接跑起来。

AMQP 模型:生产者不认识消费者

mermaid
flowchart LR
    P[Producer] --> E[Exchange]
    subgraph Server["RabbitMQ Server"]
        E --> B1[Binding routingKey=order.created]
        E --> B2[Binding routingKey=order.*]
        E --> B3[Binding routingKey=#]
        B1 --> Q1[Queue: order.created.only]
        B2 --> Q2[Queue: all.order.events]
        B3 --> Q3[Queue: everything]
    end
    Q1 --> C1[Consumer: 订单处理]
    Q2 --> C2[Consumer: 物流监控]
    Q3 --> C3[Consumer: 审计归档]

和 Kafka 的对比更直观:

  • Kafka:Producer 写 Topic,Consumer 拉 Topic,路由逻辑在客户端(分区器)
  • RabbitMQ:Producer 只发 Exchange,Exchange 根据 Binding 规则把消息投到符合条件的 Queue,Consumer 从 Queue 消费。生产者和消费者完全不知道对方的存在

四种 Exchange 类型,一图见区别

Exchange 类型决定了消息是"广播给所有人"、"只给指定的人"还是"按规则匹配"。

Direct Exchange

路由规则:routingKey 精确匹配。生产者在发送时指定 routingKey,Queue 绑定时也指定 routingKey,两者完全一致才会投递。

场景:日志级别路由——error 路由到告警队列,info 路由到普通日志队列。

Exchange: logs.direct
  Binding: queue.error ← routingKey="error"
  Binding: queue.info  ← routingKey="info"

消息 routingKey="error" → 只进 queue.error
消息 routingKey="warn"  → 无匹配,丢弃

Fanout Exchange

路由规则:忽略 routingKey,发给所有绑定的 Queue。纯粹广播。

场景:配置变更通知,所有服务都要收到。

Exchange: config.fanout
  Binding: queue.service-a
  Binding: queue.service-b
  Binding: queue.service-c

消息任意 routingKey → 三个 Queue 各复制一份

Topic Exchange

路由规则:routingKey 按通配符匹配* 匹配一个单词,# 匹配零或多个单词。

场景:订单事件分发——order.created 给分析系统,order.shipped 给物流系统,order.# 给全量归档。

Exchange: order.topic
  Binding: queue.analytics  ← routingKey="order.created"
  Binding: queue.logistics  ← routingKey="order.shipped"
  Binding: queue.archive    ← routingKey="order.#"

消息 routingKey="order.created" → 去 analytics + archive
消息 routingKey="order.shipped" → 去 logistics + archive

Headers Exchange

路由规则:匹配消息头(headers)的键值对,支持 x-match=all(全部匹配)或 x-match=any(任一匹配)。不常用,通常用 Topic 替代。

可靠投递:两段拦截

消息从生产者到消费者,中间有两段最容易丢:

  1. 生产者 → Exchange → Queue:网络闪断、Exchange 不存在、路由无匹配
  2. Queue → Consumer:消费方处理到一半崩溃

RabbitMQ 对这两段分别提供了保护机制。

生产端:Confirm + Return

java
// 配置 Confirm 回调—消息送到 Exchange 后触发
channel.confirmSelect();
channel.addConfirmListener(new ConfirmListener() {
    @Override
    public void handleAck(long deliveryTag, boolean multiple) {
        // 消息已到 Exchange
    }
    @Override
    public void handleNack(long deliveryTag, boolean multiple) {
        // 消息未到 Exchange,需要重发
    }
});

// Return 回调—Exchange 无法路由到任何 Queue 时触发
channel.addReturnListener((replyCode, replyText, exchange, routingKey,
    AMQP.BasicProperties properties, byte[] body) -> {
    // 消息被退回,记日志或存死信
});

Spring AMQP 里直接配置 publisher-confirm-type: correlatedpublisher-returns: true,框架自动处理回调。

消费端:手动 ack 与 prefetch

java
// 手动 ack:处理完再确认,没处理完挂了则消息重回队列
channel.basicConsume(queueName, false, new DefaultConsumer(channel) {
    @Override
    public void handleDelivery(String consumerTag, Envelope envelope,
        AMQP.BasicProperties properties, byte[] body) throws IOException {
        try {
            process(body);  // 业务处理
            channel.basicAck(envelope.getDeliveryTag(), false);
        } catch (Exception e) {
            channel.basicNack(envelope.getDeliveryTag(), false, true);  // true=重回队列
        }
    }
});

// prefetch:一次只拉 N 条,防止消费者被突增消息压垮
channel.basicQos(10);  // 每次最多 10 条未 ack

prefetch 不设的后果:默认无限拉取,几百条消息涌进内存,消费者处理慢导致内存 OOM。生产环境建议根据单条处理耗时和 JVM 堆大小计算,一般 10-50。

死信与 TTL:消息的三种归宿

消息不会永远待在 Queue 里。当以下情况发生时,消息变成"死信"(Dead Letter):

  1. 消费被拒绝(basicNack 或 basicReject,且 requeue=false
  2. TTL 过期:消息在 Queue 中存活超过设定时间
  3. 队列达到最大长度:超出 x-max-lengthx-max-length-bytes 的消息

死信去哪?配置 x-dead-letter-exchange 即可,死信自动投到指定 Exchange+RoutingKey。

java
// 声明死信队列
Map<String, Object> args = new HashMap<>();
args.put("x-dead-letter-exchange", "dlx-exchange");
args.put("x-dead-letter-routing-key", "dlx-routing-key");
args.put("x-message-ttl", 60000);  // 60 秒未消费变死信
args.put("x-max-length", 1000);    // 最多 1000 条
channel.queueDeclare("business-queue", true, false, false, args);

典型场景:订单超时未支付。订单消息 TTL 30 分钟,到期进死信队列,消费者从死信队列里拉出来做取消订单操作。

动手实操:Spring AMQP 完整示例

先起 RabbitMQ

bash
docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3.13-management

浏览器打开 http://localhost:15672,guest/guest 登录。管理界面可以直观看到:Connections、Channels、Exchanges、Queues 四个 tab 分别是连接、信道、交换器、队列的实时状态。

Spring Boot 配置

yaml
spring:
  rabbitmq:
    host: localhost
    port: 5672
    publisher-confirm-type: correlated    # 开启 Confirm 回调
    publisher-returns: true               # 开启 Return 回调
    template:
      mandatory: true                     # 路由不到时触发 Return
    listener:
      simple:
        acknowledge-mode: manual          # 手动 ack
        prefetch: 10

完整代码:声明、发送、消费

java
@Configuration
public class RabbitConfig {

    // 1. 声明 Exchange、Queue、Binding
    @Bean
    public TopicExchange orderExchange() {
        return new TopicExchange("order.topic", true, false);
    }

    @Bean
    public Queue orderQueue() {
        // 死信队列配置:超过 60s 未消费进死信
        Map<String, Object> args = new HashMap<>();
        args.put("x-dead-letter-exchange", "dlx.direct");
        args.put("x-dead-letter-routing-key", "order.dead");
        args.put("x-message-ttl", 60000);
        return new Queue("order.queue", true, false, false, args);
    }

    @Bean
    public Binding orderBinding(Queue orderQueue, TopicExchange orderExchange) {
        return BindingBuilder.bind(orderQueue)
            .to(orderExchange).with("order.#");
    }

    // 2. 死信队列
    @Bean
    public DirectExchange dlxExchange() {
        return new DirectExchange("dlx.direct");
    }

    @Bean
    public Queue dlxQueue() {
        return new Queue("dlx.order.queue", true);
    }

    @Bean
    public Binding dlxBinding() {
        return BindingBuilder.bind(dlxQueue())
            .to(dlxExchange()).with("order.dead");
    }
}

@Component
@Slf4j
public class OrderPublisher {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    @PostConstruct
    public void init() {
        // Confirm 回调
        rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
            if (ack) {
                log.info("消息已到 Exchange, id={}", correlationData.getId());
            } else {
                log.error("消息未到 Exchange, id={}, cause={}", correlationData.getId(), cause);
                // 重试或写入本地消息表
            }
        });
        // Return 回调
        rabbitTemplate.setReturnsCallback(returned -> {
            log.warn("消息路由不到 Queue, exchange={}, routingKey={}, replyText={}",
                returned.getExchange(), returned.getRoutingKey(), returned.getReplyText());
            // 存死信或告警
        });
    }

    public void sendOrderCreated(String orderId) {
        CorrelationData cd = new CorrelationData(orderId);
        rabbitTemplate.convertAndSend("order.topic", "order.created",
            "订单创建: " + orderId, cd);
    }
}

@Component
public class OrderConsumer {

    @RabbitListener(queues = "order.queue")
    public void handleOrder(Message message, Channel channel,
                            @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
        try {
            // 业务处理
            String body = new String(message.getBody(), StandardCharsets.UTF_8);
            System.out.println("收到订单消息: " + body);
            // 处理成功,手动 ack
            channel.basicAck(tag, false);
        } catch (Exception e) {
            // 处理失败,requeue=false 进死信
            channel.basicNack(tag, false, false);
        }
    }
}

常见误区与小结

  • Exchange 名写错导致消息丢失:Confirm 回调会报 Nack,但很多人没配回调,消息无声消失。生产环境必须配 Confirm + Return。
  • 没设 prefetch 导致 OOM:默认预取无限,消费者处理慢时内存暴涨。设 prefetch: 10-50 是标准做法。
  • 死信队列忘配:消费端 nack 设置 requeue=true,消息会反复重试形成死循环。消费异常直接进死信,留人工处理。
  • 手动 ack 忘调basicConsume 第二个参数 autoAck=true,消息一拉就算确认,消费方崩溃就丢消息。生产环境必须 autoAck=false

小结:RabbitMQ 的核心价值在于灵活的路由——用 Exchange 加 Binding 组合出各种分发策略,这是 Kafka 的 Topic 模型做不到的。它的可靠投递依赖两端配合:生产端 Confirm/Return 保证消息到 Queue,消费端手动 ack 保证处理完才确认。下一篇 27. Kafka 核心概念 进入 Kafka 的世界,看 Topic-Partition 模型和 Consumer Group 是怎么工作的。

参考

参考:RabbitMQ 官方文档 https://www.rabbitmq.com/documentation.html、Spring AMQP Reference https://docs.spring.io/spring-amqp/reference/

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。
粤ICP备2026104257号-1