主题
RabbitMQ 入门:Exchange、路由与可靠投递
本文是消息队列系统学习系列的 L1 入门篇。前置:25. 为什么需要 MQ:解耦、削峰与三大核心问题。 学完可以配合面试题食用:11-rabbitmq-architecture-exchange-binding-queue、12-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 + archiveHeaders Exchange
路由规则:匹配消息头(headers)的键值对,支持 x-match=all(全部匹配)或 x-match=any(任一匹配)。不常用,通常用 Topic 替代。
可靠投递:两段拦截
消息从生产者到消费者,中间有两段最容易丢:
- 生产者 → Exchange → Queue:网络闪断、Exchange 不存在、路由无匹配
- 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: correlated 和 publisher-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 条未 ackprefetch 不设的后果:默认无限拉取,几百条消息涌进内存,消费者处理慢导致内存 OOM。生产环境建议根据单条处理耗时和 JVM 堆大小计算,一般 10-50。
死信与 TTL:消息的三种归宿
消息不会永远待在 Queue 里。当以下情况发生时,消息变成"死信"(Dead Letter):
- 消费被拒绝(basicNack 或 basicReject,且
requeue=false) - TTL 过期:消息在 Queue 中存活超过设定时间
- 队列达到最大长度:超出
x-max-length或x-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/