消息队列03:RabbitMQ
从 Exchange 路由、Queue 队列、死信队列到延迟消息插件,拆解 RabbitMQ 在企业级消息场景下的路由模型和可靠性机制。
RabbitMQ 不是”吞吐之王”,而是”协议和路由之王”。它的核心场景是企业内部的消息分发、任务队列、RPC 调用和复杂路由需求,支持 AMQP 协议,跨语言生态丰富。
这篇文章围绕四个核心机制:Exchange 路由、Queue 模型、死信队列、延迟消息,把 RabbitMQ 与 Kafka/RocketMQ 的差异和选型逻辑讲清楚。
为什么选 RabbitMQ,而不是别的
先定位清楚:RabbitMQ 的核心场景是企业级消息分发和复杂路由,支持 AMQP 协议,跨语言生态丰富,适合任务队列、RPC、消息广播和需要复杂路由的业务场景。
与 Kafka/RocketMQ 的核心差异:
| 维度 | RabbitMQ | Kafka | RocketMQ |
|---|---|---|---|
| 协议 | AMQP 0-9-1,跨语言客户端丰富 | 自定义协议,Java/Go 生态为主 | 自定义协议,Java 生态为主 |
| 路由能力 | Exchange + Routing Key,支持 fanout/direct/topic/headers | Topic + Partition,路由简单 | Topic + Tag/SQL92 过滤 |
| 吞吐量 | 单机万级 TPS,吞吐较低 | 单机十万级 TPS,吞吐最高 | 单机万级到十万级 TPS |
| 延迟消息 | 原生延迟插件(rabbitmq_delayed_message_exchange) | 不原生支持 | 原生支持 18 个延迟级别 |
| 死信队列 | 原生支持,可配置 DLX | 不原生支持 | 原生支持,自动进入 |
| 消息可靠性 | 消息确认 + 持久化 + 镜像队列 | ISR + acks=all | 同步刷盘 + 同步复制 |
| 适用场景 | 任务队列、RPC、复杂路由、跨语言 | 日志采集、实时计算、数据管道 | 金融业务消息、事务消息 |
简单判断:跨语言、任务队列、复杂路由、消息广播、RPC → 选 RabbitMQ。日志采集、实时计算 → 选 Kafka。金融交易、事务消息 → 选 RocketMQ。
核心架构:Producer、Exchange、Queue、Consumer
RabbitMQ 的架构围绕 AMQP 协议,核心是 Exchange 的路由能力:
- Producer:消息生产者,发送消息到 Exchange(不直接发 Queue)
- Exchange:消息路由器,根据 Routing Key 和绑定规则分发消息到 Queue
- Queue:消息队列,存储消息,Consumer 从 Queue 消费
- Binding:Exchange 和 Queue 的绑定关系,定义路由规则
- Consumer:消息消费者,从 Queue 拉取消息并处理
- Virtual Host:虚拟主机,隔离不同应用的 Exchange/Queue/权限
核心设计思路:Exchange 做复杂路由,Queue 做存储和消费,Binding 定义分发规则。
Producer 不直接发 Queue,而是发到 Exchange,由 Exchange 根据路由规则分发。这是 RabbitMQ 与 Kafka/RocketMQ 最大的差异:Kafka/RocketMQ 的 Producer 直接指定 Topic/Queue,路由简单;RabbitMQ 的路由规则在服务端(Exchange),Producer 只指定 Routing Key。
Exchange 类型:四种路由模式
Exchange 是 RabbitMQ 的核心,四种类型对应四种路由模式:
1. Direct Exchange:精确匹配
消息的 Routing Key 与 Binding Key 完全匹配时,消息进入 Queue:
Producer 发送:Routing Key = "order.create"
Exchange 绑定:Queue A ← Binding Key = "order.create"
Exchange 绑定:Queue B ← Binding Key = "order.update"
结果:消息只进入 Queue A
适用场景:需要精确路由的业务场景,如订单创建、支付成功、库存更新各自路由到不同队列。
2. Fanout Exchange:广播模式
消息进入 Exchange 后,广播到所有绑定的 Queue,忽略 Routing Key:
Producer 发送:Routing Key = "any"
Exchange 绑定:Queue A、Queue B、Queue C
结果:消息进入所有三个 Queue
适用场景:消息广播、事件通知、日志分发到多个下游。
3. Topic Exchange:通配符匹配
Binding Key 支持通配符,* 匹配一个单词,# 匹配零或多个单词:
Producer 发送:Routing Key = "order.create.success"
Exchange 绑定:Queue A ← Binding Key = "order.create.*"
Exchange 绑定:Queue B ← Binding Key = "order.#"
结果:消息进入 Queue A 和 Queue B
适用场景:需要灵活路由的场景,如订单相关消息统一路由、按事件类型分发。
4. Headers Exchange:属性匹配
不依赖 Routing Key,而是根据消息的 Headers 属性匹配。Headers 是一组键值对,Binding 规则可以是”全部匹配”或”任意匹配”。
适用场景:路由规则复杂,需要按多个属性组合匹配的场景。但 Headers Exchange 性能较低,一般不推荐。
Exchange 类型对比
| Exchange 类型 | 路由规则 | 适用场景 |
|---|---|---|
| Direct | Routing Key 精确匹配 | 精确路由,如订单创建、支付成功 |
| Fanout | 广播到所有 Queue | 消息广播、事件通知 |
| Topic | 通配符匹配 | 灵活路由,如 order.*、order.# |
| Headers | Headers 属性匹配 | 复杂属性组合路由(不推荐) |
Queue 模型:消息存储和消费
Queue 是消息的物理存储单元,Consumer 从 Queue 消费消息。
消息确认机制
Consumer 消费消息后,需要发送 ACK 确认,Broker 才删除消息。未确认的消息会重新入队,供其他 Consumer 消费。
确认模式:
| 模式 | 说明 | 特点 |
|---|---|---|
| 自动确认(autoAck=true) | 消息投递后立即删除 | 可能丢消息,不推荐 |
| 手动确认(autoAck=false) | 业务处理成功后手动发送 ACK | 可靠,推荐 |
手动确认代码:
channel.basicConsume(queueName, false, new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope,
AMQP.BasicProperties properties, byte[] body) {
try {
// 业务处理
channel.basicAck(envelope.getDeliveryTag(), false); // 单条确认
} catch (Exception e) {
channel.basicNack(envelope.getDeliveryTag(), false, true); // 拒绝并重新入队
}
}
});
消息持久化
消息、Queue、Exchange 都可以持久化,保证 Broker 故障后数据不丢失:
| 持久化级别 | 说明 |
|---|---|
| Exchange 持久化 | Exchange 声明时 durable=true,Broker 重启后 Exchange 存在 |
| Queue 持久化 | Queue 声明时 durable=true,Broker 重启后 Queue 存在 |
| Message 持久化 | 发送消息时 deliveryMode=2,消息写入磁盘 |
注意:只持久化消息不持久化 Queue/Exchange,Broker 重启后 Queue/Exchange 不存在,消息也无法投递。三者都要持久化才能保证可靠性。
持久化的代价:消息写入磁盘,吞吐降低。如果允许少量丢消息,可以不持久化消息,只持久化 Queue/Exchange。
预取数量(Prefetch Count)
Consumer 可以设置 Prefetch Count,限制未确认消息的最大数量,实现公平分发:
channel.basicQos(prefetchCount); // 限制未确认消息数量
假设 Prefetch Count = 10,Consumer 未确认消息最多 10 条,Broker 不会继续投递,等 Consumer 确认后才继续投递。这避免了某个 Consumer 处理慢,消息堆积在它上面,其他 Consumer 空闲的问题。
死信队列:消费失败的兜底
RabbitMQ 原生支持死信队列(Dead Letter Queue,DLQ)。消息进入死信队列的条件:
| 条件 | 说明 |
|---|---|
| 消息被拒绝 | Consumer 发送 basicNack 或 basicReject,且 requeue=false |
| 消息过期 | 消息 TTL(Time-To-Live)到期 |
| 队列满 | Queue 达到最大长度,无法继续写入 |
死信队列配置:
// 声明死信交换器
channel.exchangeDeclare("dlx.exchange", "direct");
channel.queueDeclare("dlx.queue", true, false, false, null);
channel.queueBind("dlx.queue", "dlx.exchange", "dlx.routing.key");
// 声明业务队列,配置死信队列
Map<String, Object> args = new HashMap<>();
args.put("x-dead-letter-exchange", "dlx.exchange");
args.put("x-dead-letter-routing-key", "dlx.routing.key");
channel.queueDeclare("business.queue", true, false, false, args);
消息进入死信队列后,可以人工排查和补偿,或者重新投递到其他队列。
延迟消息:延迟插件实现
RabbitMQ 不原生支持延迟消息,需要安装 rabbitmq_delayed_message_exchange 插件。
延迟插件原理
插件提供 x-delayed-message Exchange 类型,消息发送时设置 x-delay 头(延迟毫秒数),Exchange 将消息暂存,延迟到期后路由到 Queue。
代码示例:
// 声明延迟交换器
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct"); // 底层路由类型
channel.exchangeDeclare("delayed.exchange", "x-delayed-message", true, false, args);
// 发送延迟消息
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.headers(Map.of("x-delay", 30000)) // 延迟 30 秒
.build();
channel.basicPublish("delayed.exchange", "order.timeout", props, body);
延迟消息暂存在 Exchange 内部,不进入 Queue,延迟到期后才路由。代价是延迟消息存储在内存,大量延迟消息会占用内存,不适合大规模延迟场景。
与 RocketMQ 延迟消息的差异
| 维度 | RabbitMQ 延迟插件 | RocketMQ 延迟消息 |
|---|---|---|
| 原生支持 | 需安装插件 | 原生支持 18 个延迟级别 |
| 延迟精度 | 支持任意毫秒数 | 只支持 18 个固定级别 |
| 存储位置 | Exchange 内存暂存 | Broker 磁盘存储 |
| 大规模延迟 | 内存压力大,不适合 | 磁盘存储,适合大规模 |
简单判断:小规模延迟消息、延迟精度要求高 → RabbitMQ 延迟插件。大规模延迟消息、延迟级别固定 → RocketMQ。
镜像队列:高可用机制
RabbitMQ 的高可用方案是镜像队列(Mirrored Queue):Queue 在多个 Broker 节点上镜像,主节点故障时自动切换到镜像节点。
镜像队列策略:
策略名称:ha-all
模式:all(所有节点镜像)
应用范围:所有 Queue
镜像队列的代价:
- 写入性能降低:主节点写入后,需同步到镜像节点
- 网络带宽消耗:镜像同步占用带宽
- 不适合大规模:大量镜像队列会占用大量资源
Kafka/RocketMQ 的副本机制更成熟,吞吐更高。RabbitMQ 的镜像队列适合小规模、企业内部场景。
常见误区
| 误区 | 真实情况 |
|---|---|
| RabbitMQ 吞吐高 | 单机万级 TPS,吞吐较低,不适合日志管道 |
| Exchange 直接存储消息 | Exchange 只做路由,消息存储在 Queue |
| 延迟消息原生支持 | 需安装插件,不适合大规模延迟 |
| 死信队列 = 分布式事务 | 死信队列只是消费失败的兜底,不是分布式事务 |
| 消息确认后立刻删除 | 消息确认后 Broker 才删除,未确认会重新入队 |
项目判断
引入 RabbitMQ 时,先回答三个问题:
- 跨语言需求:是否有非 Java 的消费者(如 Python、Go、Node.js)?
- 路由复杂度:是否需要复杂路由规则(如按事件类型、按业务属性)?
- 吞吐量:日均消息量是否在百万级以下?
如果跨语言、路由复杂、吞吐不高 → RabbitMQ 合适。如果日志采集、实时计算、吞吐优先 → Kafka 合适。如果金融交易、事务消息 → RocketMQ 合适。
线上排查消息堆积:
- 查看 Queue 的消息数量:
rabbitmqctl list_queues name messages - 查看 Consumer 的预取数量和处理速度
- 检查是否有 Consumer 未确认消息堆积(Prefetch Count 设置过大)
- 增加 Consumer 数量或提高处理速度
- 检查 Queue 是否满,导致消息无法写入