RabbitMQ 原理与实战
RabbitMQ 概述
什么是 RabbitMQ
RabbitMQ 是一个开源的消息代理(Message Broker)和队列服务器,它使用 Erlang/OTP 语言编写,实现了 AMQP(Advanced Message Queuing Protocol,高级消息队列协议) 标准。它能够在高并发、分布式系统中作为消息中间件,实现应用解耦、流量削峰、异步通信等功能。
技术栈:Erlang/OTP
RabbitMQ 构建于 Erlang/OTP 平台之上,这赋予了它几个关键特性:
- 轻量级进程模型:Erlang 的 Actor 模型让每个队列、每个连接都可以运行在独立的轻量级进程中,天然支持高并发。
- 容错与热更新:OTP 框架提供了 Supervisor 树和热代码替换能力,使 RabbitMQ 具备生产级的高可用性。
- 跨节点通信:Erlang 的分布式能力让 RabbitMQ 集群管理变得简单。
AMQP 协议
AMQP 是一个应用层协议,专注于面向消息的中间件(MOM)。它的核心设计理念是:
- 线级协议(Wire-level Protocol):无论客户端使用什么语言,只要遵循 AMQP 规范就可以与 Broker 通信。
- 可互操作:不同厂商的 AMQP 实现可以相互通信。
- 模型清晰:定义了 Exchange、Queue、Binding 等抽象概念,逻辑模型与物理拓扑分离。
AMQP 0-9-1 是 RabbitMQ 广泛支持的主流版本(而非 AMQP 1.0,后者是截然不同的协议)。
历史背景
- 2004 年:JPMorgan Chase 开始设计 AMQP 协议,旨在创建一个标准化的金融消息传输协议。
- 2006 年:Rabbit Technologies 成立,开始开发 RabbitMQ。
- 2010 年:VMware 收购 Rabbit Technologies。
- 2013 年:Pivotal Software 接手(由 VMware 与 EMC 合资成立)。
- 2019 年:VMware 再次收购 Pivotal,RabbitMQ 回归 VMware 旗下。
- 至今:RabbitMQ 是 ASF(Apache Software Foundation)旗下项目,但仍由 VMware(现属 Broadcom)主导维护,是目前最流行的开源消息队列之一。
基本工作流程
Publisher → Exchange → Queue → Consumer生产者将消息发送到 Exchange,Exchange 根据 Binding 规则将消息路由到一个或多个 Queue,消费者从 Queue 中拉取或订阅消息。
核心概念
Connection
Connection 是客户端与 RabbitMQ 服务器之间的 TCP 长连接。它是一个物理连接,负责管理底层套接字、协议版本协商、身份认证等。
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setUsername("guest");
factory.setPassword("guest");
// 创建一个 TCP 连接
Connection connection = factory.newConnection();Channel
Channel 是在 Connection 基础上建立的虚拟连接(Multiplexing)。一个 TCP 连接可以创建数百个 Channel,每个 Channel 独立处理一个会话任务。RabbitMQ 的所有 AMQP 指令(如声明队列、发布消息、消费消息)都通过 Channel 进行。
// 从连接中创建 Channel
Channel channel = connection.createChannel();
// 通过 Channel 声明队列
channel.queueDeclare("my.queue", true, false, false, null);
// 通过 Channel 发布消息
channel.basicPublish("", "my.queue", null, "Hello".getBytes());
// 通过 Channel 消费消息
channel.basicConsume("my.queue", true, (consumerTag, delivery) -> {
String body = new String(delivery.getBody());
System.out.println("收到消息: " + body);
}, consumerTag -> {});为什么需要 Channel? 创建 TCP 连接代价高昂(三次握手 + TLS 协商)。通过 Channel 复用同一个 TCP 连接,可以大幅降低连接开销,提升吞吐量。
Queue
Queue 是消息的存储容器,本质是一个有序的字节数组缓冲区。队列具有以下属性:
| 参数 | 类型 | 说明 |
|---|---|---|
name | String | 队列名称,空字符串由服务端自动生成 |
durable | boolean | 是否持久化(服务重启后队列是否存在) |
exclusive | boolean | 是否排他(仅当前 Connection 可见,断开后自动删除) |
autoDelete | boolean | 是否自动删除(最后一个消费者断开后删除) |
arguments | Map | 扩展参数(TTL、优先级、死信等) |
// 声明一个持久化、非排他、非自动删除的队列
channel.queueDeclare("order.pay.queue", true, false, false, null);
// 声明一个临时队列(排他 + 自动删除),常用于 RPC 回调
String tempQueue = channel.queueDeclare().getQueue();Exchange
Exchange 是消息的路由处理器。生产者发送消息时,消息不会直接进入 Queue,而是先到达 Exchange,由 Exchange 根据 Routing Key 和 Binding 规则决定消息的去向。
// 声明一个 Direct 类型 Exchange
channel.exchangeDeclare("order.exchange", BuiltinExchangeType.DIRECT, true);Binding
Binding 是 Exchange 与 Queue 之间的关联规则。通过 Binding,一个 Queue 可以绑定到某个 Exchange 并指定 Routing Key,从而决定哪些消息应该进入该队列。
// 将队列绑定到 Exchange,路由键为 "order.pay"
channel.queueBind("order.pay.queue", "order.exchange", "order.pay");关系示意:
┌──────────────┐
│ Exchange │
│ (order.ex) │
└──────┬───────┘
┌───────┴───────┐
Binding │ order.pay │ order.refund
▼ ▼
┌──────────┐ ┌──────────┐
│Queue A │ │Queue B │
│(pay) │ │(refund) │
└──────────┘ └──────────┘Routing Key
Routing Key 是消息携带的路由标签,Exchange 根据这个标签匹配 Binding 规则,决定消息转发目标。Routing Key 的长度限制为 255 字节。
// 发布消息时指定 Routing Key
channel.basicPublish("order.exchange", "order.pay", null, msg.getBytes());Virtual Host (vhost)
Virtual Host 是 RabbitMQ 的多租户隔离机制。一个 vhost 拥有独立的 Exchange、Queue、Binding 命名空间,不同 vhost 之间完全隔离。一个 RabbitMQ 服务可以创建多个 vhost,常用于区分开发环境、测试环境、生产环境或不同业务线。
# 连接时指定 vhost
spring.rabbitmq.virtual-host=/productionExchange 类型
RabbitMQ 提供四种内置 Exchange 类型,每种类型对应不同的路由策略。
Direct Exchange
路由规则:Routing Key 精确匹配。消息的 Routing Key 必须与 Binding 的 Routing Key 完全一致。
适用场景:点对点通信、任务分发。
// 声明 Direct Exchange
channel.exchangeDeclare("direct.ex", BuiltinExchangeType.DIRECT, true);
// 声明两个队列
channel.queueDeclare("queue.error", true, false, false, null);
channel.queueDeclare("queue.info", true, false, false, null);
// 绑定:error 队列只接收 routingKey = "error" 的消息
channel.queueBind("queue.error", "direct.ex", "error");
// info 队列接收 routingKey = "info" 或 "warn" 的消息
channel.queueBind("queue.info", "direct.ex", "info");
channel.queueBind("queue.info", "direct.ex", "warn");
// 发送消息
channel.basicPublish("direct.ex", "error", null, "严重错误".getBytes()); // → queue.error
channel.basicPublish("direct.ex", "info", null, "普通信息".getBytes()); // → queue.info
channel.basicPublish("direct.ex", "warn", null, "警告".getBytes()); // → queue.info架构示意:
Routing Key: "error" ──► Direct Exchange ──► Queue (error)
Routing Key: "info" ──► Direct Exchange ──► Queue (info)
Routing Key: "warn" ──► Direct Exchange ──► Queue (info)Topic Exchange
路由规则:Routing Key 通配符匹配。Routing Key 使用 . 分隔的单词组成,支持两个通配符:
*:匹配一个单词#:匹配零个或多个单词
适用场景:基于主题的发布/订阅(如日志分类、地理位置事件分发)。
channel.exchangeDeclare("topic.ex", BuiltinExchangeType.TOPIC, true);
channel.queueDeclare("queue.order", true, false, false, null);
channel.queueDeclare("queue.payment", true, false, false, null);
// order 队列接收所有 order 相关消息
channel.queueBind("queue.order", "topic.ex", "order.#");
// payment 队列接收中国区的支付消息
channel.queueBind("queue.payment", "topic.ex", "payment.cn.*");
// 发送消息
channel.basicPublish("topic.ex", "order.created", null, "订单创建".getBytes()); // → queue.order
channel.basicPublish("topic.ex", "order.paid", null, "订单支付".getBytes()); // → queue.order
channel.basicPublish("topic.ex", "payment.cn.alipay", null, "支付宝".getBytes()); // → queue.payment
channel.basicPublish("topic.ex", "payment.us.stripe", null, "Stripe".getBytes()); // → 无匹配Fanout Exchange
路由规则:广播。忽略 Routing Key,将消息发送到所有绑定的队列。
适用场景:全局广播、配置更新通知、缓存刷新。
channel.exchangeDeclare("fanout.ex", BuiltinExchangeType.FANOUT, true);
channel.queueDeclare("queue.cache", true, false, false, null);
channel.queueDeclare("queue.log", true, false, false, null);
// Routing Key 被忽略
channel.queueBind("queue.cache", "fanout.ex", "");
channel.queueBind("queue.log", "fanout.ex", "");
// 所有绑定队列都会收到消息
channel.basicPublish("fanout.ex", "", null, "配置变更".getBytes());
// → queue.cache 收到
// → queue.log 收到Headers Exchange
路由规则:匹配消息的 Headers 属性(键值对),而非 Routing Key。支持 x-match 参数:
x-match = all(默认):所有 header 都匹配x-match = any:任一 header 匹配即可
适用场景:多维度路由、复杂的条件匹配。
channel.exchangeDeclare("headers.ex", BuiltinExchangeType.HEADERS, true);
channel.queueDeclare("queue.a", true, false, false, null);
channel.queueDeclare("queue.b", true, false, false, null);
// 绑定:queue.a 需要 content-type=json 且 encoding=utf-8
Map<String, Object> argsA = new HashMap<>();
argsA.put("x-match", "all");
argsA.put("content-type", "json");
argsA.put("encoding", "utf-8");
channel.queueBind("queue.a", "headers.ex", "", argsA);
// 绑定:queue.b 需要 format=pdf 或 format=doc
Map<String, Object> argsB = new HashMap<>();
argsB.put("x-match", "any");
argsB.put("format", "pdf");
argsB.put("format", "doc");
channel.queueBind("queue.b", "headers.ex", "", argsB);
// 发送消息
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.headers(Map.of("content-type", "json", "encoding", "utf-8"))
.build();
channel.basicPublish("headers.ex", "", props, "数据".getBytes());
// → queue.a 匹配,queue.b 不匹配四种 Exchange 对比
| 类型 | 路由依据 | 通配符 | 性能 | 典型场景 |
|---|---|---|---|---|
| Direct | Routing Key 精确匹配 | 无 | 高 | 点对点任务分发 |
| Topic | Routing Key 通配符匹配 | *、# | 中 | 基于主题的发布订阅 |
| Fanout | 广播所有绑定队列 | 无 | 最高 | 全局广播通知 |
| Headers | 消息 Headers 匹配 | 无 | 低(需解析 Headers) | 复杂条件路由 |
消息确认
消息确认机制是保证消息不丢失的核心能力,分为生产者确认和消费者确认两个维度。
Publisher Confirm(生产者确认)
生产者确认确保消息已成功到达 Broker。RabbitMQ 从 3.6.x 开始支持 Publisher Confirm 和 Publisher Return 两种机制。
工作原理
- 生产者将 Channel 设置为 Confirm 模式(
confirm.select)。 - 每发布一条消息,Broker 返回一个
basic.ack(确认)或basic.nack(否定确认)。 - 生产者通过回调或异步等待方式感知消息是否到达 Exchange。
// 启用 Publisher Confirm
channel.confirmSelect();
// 同步方式(逐个确认)
channel.basicPublish("", "queue.name", null, "msg1".getBytes());
if (channel.waitForConfirms()) {
System.out.println("消息已确认");
}
// 异步方式(推荐生产环境使用)
channel.addConfirmListener((deliveryTag, multiple) -> {
System.out.println("ACK: " + deliveryTag);
}, (deliveryTag, multiple) -> {
System.out.println("NACK: " + deliveryTag + ",需要重发");
});Spring Boot 中的配置
spring:
rabbitmq:
publisher-confirm-type: correlated # 启用 Confirm(correlated 模式)
publisher-returns: true # 启用 Return(消息无法路由时回调)@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
RabbitTemplate template = new RabbitTemplate(connectionFactory);
template.setMandatory(true); // 消息无法路由时返回给生产者
template.setConfirmCallback((correlationData, ack, cause) -> {
if (ack) {
log.info("消息确认成功: {}", correlationData.getId());
} else {
log.error("消息确认失败: {}", cause);
// 重发或记录到数据库
}
});
template.setReturnsCallback(returned -> {
log.warn("消息未路由: {}", returned.getMessage());
});
return template;
}Consumer ACK(消费者确认)
消费者处理完消息后需要向 Broker 发送确认,告知 Broker 可以安全删除该消息。RabbitMQ 提供两种确认模式。
Manual ACK(手动确认)
消费者显式调用 basicAck,Broker 收到后才会标记消息为已处理。
// 消费消息,关闭 autoAck
channel.basicConsume("queue.name", false, (consumerTag, delivery) -> {
try {
// 业务处理
process(delivery.getBody());
// 手动确认:第二个参数 false 表示不批量确认
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
} catch (Exception e) {
// 处理失败:重新入队(true)或进入死信(false)
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true);
}
}, consumerTag -> {});Auto ACK(自动确认)
// 第一个参数 true 表示自动确认
channel.basicConsume("queue.name", true, (consumerTag, delivery) -> {
// 只要消费者收到消息,Broker 立即删除
process(delivery.getBody());
}, consumerTag -> {});注意:自动确认模式下,如果消费者在处理过程中崩溃,消息会丢失。生产环境推荐使用手动确认。
NACK(否定确认)
当消费者处理消息失败时,可以发送 NACK 并决定是否重新入队。
try {
process(delivery.getBody());
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
} catch (Exception e) {
// 重新入队(requeue=true)—— 死循环风险
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true);
// 不入队,进入死信队列(requeue=false)
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
// 拒绝单条消息:basicReject(不支持批量)
channel.basicReject(delivery.getEnvelope().getDeliveryTag(), false);
}Spring Boot @RabbitListener 手动 ACK 配置
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual # manual / auto / none@RabbitListener(queues = "order.queue")
public void handleOrder(Message message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) {
try {
// 业务处理
process(message);
channel.basicAck(tag, false);
} catch (Exception e) {
channel.basicNack(tag, false, false); // 进入死信
}
}确认流程总结
生产者 RabbitMQ 消费者
│ │ │
│── publish ────────────►│ │
│◄── ack ───────────────│ │ ← Publisher Confirm
│ │── deliver ──────────────►│
│ │◄── ack ─────────────────│ ← Consumer ACK
│ │ (消息删除) │死信队列
DLX 原理
死信队列(Dead Letter Queue) 是处理无法被正常消费的消息的机制。当消息满足特定条件变为"死信"时,RabbitMQ 会自动将其重新发布到指定的 Exchange(即 Dead Letter Exchange,DLX),然后路由到对应的死信队列。
┌───────────┐
│ 原 Exchange │
└─────┬─────┘
│
┌────▼────┐ 死信条件触发 ┌──────────┐
│ 业务 Queue│ ──────────────────────────► │ DLX │
└─────────┘ └────┬─────┘
│
┌────▼─────┐
│ 死信 Queue │
└──────────┘死信来源
消息变成死信的三种场景:
- 消息被消费者 NACK(requeue=false)或 Reject
- 消息过期(TTL 到期仍未被消费)
- 队列达到最大长度
配置死信队列
为业务队列指定 x-dead-letter-exchange 和 x-dead-letter-routing-key:
// 1. 声明死信 Exchange 和死信队列
channel.exchangeDeclare("dlx.exchange", BuiltinExchangeType.DIRECT, true);
channel.queueDeclare("dlx.queue", true, false, false, null);
channel.queueBind("dlx.queue", "dlx.exchange", "dlx.rk");
// 2. 声明业务队列,绑定死信配置
Map<String, Object> args = new HashMap<>();
args.put("x-dead-letter-exchange", "dlx.exchange"); // DLX 名称
args.put("x-dead-letter-routing-key", "dlx.rk"); // 死信 Routing Key
channel.queueDeclare("business.queue", true, false, false, args);
channel.queueBind("business.queue", "business.exchange", "business.rk");Spring Boot 配置死信队列
@Configuration
public class RabbitDLXConfig {
// 死信 Exchange
@Bean
public DirectExchange dlxExchange() {
return new DirectExchange("dlx.exchange");
}
// 死信队列
@Bean
public Queue dlxQueue() {
return QueueBuilder.durable("dlx.queue").build();
}
@Bean
public Binding dlxBinding() {
return BindingBuilder.bind(dlxQueue()).to(dlxExchange()).with("dlx.rk");
}
// 业务队列(绑定 DLX)
@Bean
public Queue businessQueue() {
return QueueBuilder.durable("business.queue")
.deadLetterExchange("dlx.exchange") // 指定 DLX
.deadLetterRoutingKey("dlx.rk") // 死信路由键
.build();
}
@Bean
public DirectExchange businessExchange() {
return new DirectExchange("business.exchange");
}
@Bean
public Binding businessBinding() {
return BindingBuilder.bind(businessQueue()).to(businessExchange()).with("business.rk");
}
}TTL 设置
消息 TTL:设置消息在队列中的存活时间,到期后变为死信。
队列 TTL:设置队列未被使用时的自动删除时间。
// 方式一:设置单条消息的 TTL(毫秒)
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.expiration("60000") // 60 秒过期
.build();
channel.basicPublish("exchange", "rk", props, "msg".getBytes());
// 方式二:设置队列中的消息 TTL
Map<String, Object> args = new HashMap<>();
args.put("x-message-ttl", 60000); // 队列中所有消息 60 秒过期
channel.queueDeclare("ttl.queue", true, false, false, args);
// 方式三:Spring Boot 方式
@Bean
public Queue ttlQueue() {
return QueueBuilder.durable("ttl.queue")
.ttl(60000) // 消息 TTL
.deadLetterExchange("dlx.exchange")
.build();
}完整死信流程图
RabbitMQ
┌───────────────────────────────────────────────────────┐
│ │
│ 业务消息 死信 Exchange (DLX) │
│ ─────► 业务 Queue ──(死信)──► ┌──────────────┐ │
│ │ dlx.exchange │ │
│ └──────┬───────┘ │
│ │ │
│ ▼ │
│ ┌──────────────┐ │
│ │ 死信 Queue │ │
│ │ (消费/归档) │ │
│ └──────────────┘ │
└───────────────────────────────────────────────────────┘延迟队列
RabbitMQ 原生不支持延迟队列,但可以通过两种方式实现。
方式一:死信实现延迟(推荐原生方案)
利用死信队列特性:消息设置 TTL 进入业务队列,过期后进入 DLX,再路由到真正的消费队列。
@Configuration
public class DelayQueueConfig {
// 死信 Exchange(充当延迟 Exchange)
@Bean
public DirectExchange delayExchange() {
return new DirectExchange("delay.exchange");
}
// 被延迟消费的实际队列
@Bean
public Queue delayConsumeQueue() {
return QueueBuilder.durable("delay.consume.queue").build();
}
@Bean
public Binding delayConsumeBinding() {
return BindingBuilder.bind(delayConsumeQueue()).to(delayExchange()).with("delay.consume");
}
// 延迟队列(业务发送到此队列,消息带 TTL 过期后进入 delayExchange)
@Bean
public Queue delayQueue() {
return QueueBuilder.durable("delay.queue")
.deadLetterExchange("delay.exchange")
.deadLetterRoutingKey("delay.consume")
.build();
}
@Bean
public DirectExchange delayDirectExchange() {
return new DirectExchange("delay.direct.exchange");
}
@Bean
public Binding delayBinding() {
return BindingBuilder.bind(delayQueue()).to(delayDirectExchange()).with("delay");
}
}// 发送延迟消息(TTL = 30 秒)
rabbitTemplate.convertAndSend("delay.direct.exchange", "delay", message, msg -> {
msg.getMessageProperties().setExpiration("30000");
return msg;
});缺点:同一队列中只有第一个消息到期后,后续消息才能被检查到(排队机制)。社区将此问题称为"队列前面的消息过期时间决定了后续消息的检查时机"。
方式二:延迟插件(rabbitmq-delayed-message-exchange)
官方插件实现了真正的延迟消息,支持毫秒级精度且无上述排队问题。
安装插件
# 在 RabbitMQ 服务器上启用插件
rabbitmq-plugins enable rabbitmq_delayed_message_exchange声明延迟 Exchange
// 原生 API
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct"); // 底层 Exchange 类型
channel.exchangeDeclare("delayed.exchange", "x-delayed-message", true, false, args);
channel.queueDeclare("delayed.queue", true, false, false, null);
channel.queueBind("delayed.queue", "delayed.exchange", "delayed.rk");Spring Boot 集成延迟插件
@Configuration
public class DelayedPluginConfig {
@Bean
public CustomExchange delayedExchange() {
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct");
return new CustomExchange("delayed.exchange", "x-delayed-message", true, false, args);
}
@Bean
public Queue delayedQueue() {
return QueueBuilder.durable("delayed.queue").build();
}
@Bean
public Binding delayedBinding() {
return BindingBuilder.bind(delayedQueue()).to(delayedExchange()).with("delayed.rk").noargs();
}
}// 发送延迟消息(延迟 10 秒)
rabbitTemplate.convertAndSend("delayed.exchange", "delayed.rk", "延迟消息", msg -> {
msg.getMessageProperties().setDelay(10000); // 毫秒
return msg;
});两种实现方式对比
| 对比项 | 死信实现 | 插件实现 |
|---|---|---|
| 精度 | 受队列头部阻塞影响 | 毫秒级精确 |
| 依赖 | 无额外依赖 | 需安装插件(erlang 24+ 兼容) |
| 性能 | 高(原生死信) | 中(需额外校验延迟) |
| 限制 | 同一队列过期时间需一致 | 无特殊限制 |
| 推荐度 | 简单场景适用 | 生产环境强烈推荐 |
优先级队列
原理
RabbitMQ 支持通过 x-max-priority 参数为队列设置优先级范围(0-255,推荐 0-10)。优先级高的消息优先被消费,但前提是队列中有积压消息。
声明优先级队列
Map<String, Object> args = new HashMap<>();
args.put("x-max-priority", 10); // 优先级范围 0-10
channel.queueDeclare("priority.queue", true, false, false, args);Spring Boot 配置
@Bean
public Queue priorityQueue() {
return QueueBuilder.durable("priority.queue")
.maxPriority(10)
.build();
}发送优先级消息
// 原生 API
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.priority(8) // 设置消息优先级
.build();
channel.basicPublish("", "priority.queue", props, "高优先级消息".getBytes());
// Spring Boot
rabbitTemplate.convertAndSend("priority.queue", "高优先级消息", msg -> {
msg.getMessageProperties().setPriority(8);
return msg;
});使用场景
- 付费用户优先处理:免费用户消息优先级 1,VIP 用户消息优先级 10。
- 告警系统:严重告警(CRITICAL)优先级高于 WARNING 和 INFO。
- 订单处理:加急订单优先于普通订单。
注意事项
- 默认优先级为 0,不设置优先级的消息按最低优先级处理。
- 优先级只在消息积压时才有意义。如果消费者处理能力充足、队列无积压,优先级不会产生实际效果。
- 建议优先级范围设置为 1-10,过大的范围会额外消耗内存。
集群模式
普通集群(Classic Cluster)
原理
RabbitMQ 普通集群通过 Erlang 分布式能力互联。元数据(Exchange、Queue、Binding 定义)在所有节点间同步,但消息数据(Queue 内容)只存储在一个节点上。
┌─────────────┐ ┌─────────────┐
│ Node A │ │ Node B │
│ queue.Q │◄────────►│ queue.Q │
│ (主) │ 元数据 │ (备: 指向A)│
└─────────────┘ └─────────────┘
│
┌──────┴──────┐
│ 消息实体 │
│ 仅存于 A │
└─────────────┘- 消费者连接到 Node B 并消费
queue.Q时,Node B 从 Node A 拉取消息。 - 如果 Node A 宕机,
queue.Q的消息丢失(除非配置了镜像队列或 Quorum Queue)。
搭建示例
# 在 Node A 上
rabbitmq-server -detached
rabbitmqctl stop_app
rabbitmqctl reset
rabbitmqctl start_app
# 在 Node B 上
rabbitmq-server -detached
rabbitmqctl stop_app
rabbitmqctl reset
rabbitmqctl join_cluster rabbit@node-a
rabbitmqctl start_app
# 查看集群状态
rabbitmqctl cluster_status镜像队列(Mirrored Queue)
原理
在普通集群基础上,镜像队列将消息数据同步复制到多个节点。每个队列有一个 Master 和多个 Mirrors(Slave)。所有生产/消费请求由 Master 处理,Mirrors 异步或同步复制消息。
┌─────────────┐ ┌─────────────┐
│ Node A │ │ Node B │
│ queue.Q │◄────────►│ queue.Q │
│ (Master) │ 同步 │ (Mirror) │
│ 消息实体 │ 复制 │ 消息副本 │
└─────────────┘ └─────────────┘配置镜像策略(Policy)
# 对所有以 "ha." 开头的队列设置镜像策略
rabbitmqctl set_policy ha-all "^ha\." '{"ha-mode":"all","ha-sync-mode":"automatic"}'| 参数 | 可选值 | 说明 |
|---|---|---|
ha-mode | all / exactly / nodes | all: 所有节点;exactly: 指定数量;nodes: 指定节点名称 |
ha-params | 数字或节点列表 | 配合 ha-mode 使用 |
ha-sync-mode | manual / automatic | 新 Mirror 加入时的同步模式 |
Quorum Queue(仲裁队列)
原理
RabbitMQ 3.8+ 引入的现代替代方案,基于 Raft 共识算法。Quorum Queue 不是镜像队列的升级版,而是重新设计的队列类型。
| 特性 | 镜像队列 | Quorum Queue |
|---|---|---|
| 一致性算法 | Master-Slave 模式 | Raft 共识算法 |
| 节点数 | 任意 | 推荐 3 或 5 节点(奇数) |
| 数据安全 | 可能脑裂导致数据丢失 | 严格一致性 |
| 性能 | 高吞吐 | 略低(需共识协商) |
| 特性 | 支持全部 AMQP 特性 | 不支持临时队列、不支持优先级 |
声明仲裁队列
// 原生 API
Map<String, Object> args = new HashMap<>();
args.put("x-queue-type", "quorum"); // 指定队列类型为 quorum
channel.queueDeclare("quorum.queue", true, false, false, args);
// Spring Boot
@Bean
public Queue quorumQueue() {
return QueueBuilder.durable("quorum.queue")
.quorum() // 仲裁队列
.build();
}三种集群模式对比
| 对比项 | 普通集群 | 镜像队列 | Quorum Queue |
|---|---|---|---|
| 数据分布 | 单节点存储 | 多节点复制 | Raft 共识复制 |
| 高可用 | ❌ 节点宕机数据丢 | ✅ Master 宕机自动切换 | ✅ Leader 宕机自动选举 |
| 一致性 | 最终一致 | 最终一致(可能脑裂) | 强一致(Raft) |
| 吞吐量 | 高 | 中(主写从复制) | 中(需多数派确认) |
| 推荐场景 | 开发/测试 | 传统生产环境 | 生产环境(3.8+) |
Spring Boot 集成
添加依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>配置连接
spring:
rabbitmq:
host: 192.168.1.100
port: 5672
username: admin
password: admin123
virtual-host: /
# 生产者确认
publisher-confirm-type: correlated
publisher-returns: true
# 消费者手动确认
listener:
simple:
acknowledge-mode: manual
prefetch: 1 # 每次只取1条
concurrency: 3 # 并发消费者数
max-concurrency: 10
retry:
enabled: true # 消费失败重试
max-attempts: 3
initial-interval: 1000配置类(声明 Exchange、Queue、Binding)
@Configuration
public class RabbitMQConfig {
// ========== 订单相关 ==========
@Bean
public DirectExchange orderExchange() {
return ExchangeBuilder.directExchange("order.exchange").durable(true).build();
}
@Bean
public Queue orderQueue() {
return QueueBuilder.durable("order.queue")
.deadLetterExchange("order.dlx.exchange")
.deadLetterRoutingKey("order.dlx")
.build();
}
@Bean
public Binding orderBinding() {
return BindingBuilder.bind(orderQueue())
.to(orderExchange())
.with("order.create");
}
// ========== 死信 ==========
@Bean
public DirectExchange orderDlxExchange() {
return ExchangeBuilder.directExchange("order.dlx.exchange").durable(true).build();
}
@Bean
public Queue orderDlxQueue() {
return QueueBuilder.durable("order.dlx.queue").build();
}
@Bean
public Binding orderDlxBinding() {
return BindingBuilder.bind(orderDlxQueue())
.to(orderDlxExchange())
.with("order.dlx");
}
}生产者
@Component
@Slf4j
public class OrderPublisher {
@Autowired
private RabbitTemplate rabbitTemplate;
public void sendOrder(Order order) {
CorrelationData correlationData = new CorrelationData(order.getOrderId());
rabbitTemplate.convertAndSend("order.exchange", "order.create", order, correlationData);
log.info("订单消息已发送: {}", order.getOrderId());
}
}消费者
@Component
@Slf4j
public class OrderConsumer {
@RabbitListener(queues = "order.queue")
public void handleOrder(Message message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) {
try {
String body = new String(message.getBody());
log.info("收到订单消息: {}", body);
// 业务处理...
channel.basicAck(tag, false);
} catch (Exception e) {
log.error("处理订单消息失败", e);
try {
channel.basicNack(tag, false, false);
} catch (IOException ex) {
log.error("发送 NACK 失败", ex);
}
}
}
}消息实体序列化
@Bean
public MessageConverter messageConverter() {
// 使用 Jackson 序列化消息体(推荐)
Jackson2JsonMessageConverter converter = new Jackson2JsonMessageConverter();
converter.setCreateMessageIds(true);
return converter;
}选型建议
RabbitMQ vs RocketMQ vs Kafka
| 对比维度 | RabbitMQ | RocketMQ | Kafka |
|---|---|---|---|
| 开发语言 | Erlang | Java | Scala/Java |
| 协议 | AMQP 0-9-1 | 自研协议 | 自研协议(基于 TCP) |
| 吞吐量 | 万级/秒 | 十万级/秒 | 百万级/秒 |
| 延迟 | 微秒级 | 毫秒级 | 毫秒级 |
| 消息可靠性 | 高(Confirm + ACK) | 高(同步刷盘) | 高(ISR 副本同步) |
| 消息顺序 | 单队列有序 | 单队列有序 | 分区内有序 |
| 消息回溯 | ❌(消费后删除) | ✅(支持时间戳回溯) | ✅(基于 offset 回溯) |
| 延迟消息 | ✅(插件或死信实现) | ✅(内置 18 个等级) | ❌(需自研) |
| 死信队列 | ✅ 原生支持 | ✅ 原生支持 | ❌(需自建 DLT) |
| 事务消息 | ✅(有限支持) | ✅(完整支持) | ✅(Transactional API) |
| 消息过滤 | Header / Routing Key | Tag / SQL 表达式 | 基于 Topic / 自研 |
| 运维复杂度 | 中 | 中 | 高(需 ZooKeeper/KRaft) |
| 社区生态 | 成熟、文档丰富 | 国内生态好 | 最好(流处理生态) |
场景推荐
| 场景 | 推荐消息队列 | 理由 |
|---|---|---|
| 业务解耦 / 异步通知 | RabbitMQ | 功能完善、易用性好、延迟低 |
| 订单系统 / 交易流水 | RabbitMQ / RocketMQ | 强一致性要求,Confirm + ACK 保障 |
| 日志收集 / 埋点数据 | Kafka | 高吞吐、支持长期存储与回溯 |
| 流处理 / 实时计算 | Kafka | 完美对接 Flink、Spark Streaming |
| 电商秒杀 / 削峰填谷 | RocketMQ | 高吞吐 + 延迟消息支持 |
| IoT 设备消息 | RabbitMQ | 轻量级协议支持(MQTT 插件) |
| 金融交易 / 对账 | RocketMQ | 事务消息 + 高可用 |
| 企业内部异步 RPC | RabbitMQ | 灵活的路由 + 轻重级集成 |
选型决策树
消息中间件选型
│
├─ 需要海量吞吐(百万级/秒)?→ Kafka
│
├─ 需要事务消息 / 严格顺序?→ RocketMQ
│
├─ 需要灵活路由 / 低延迟 / 轻量级?→ RabbitMQ
│
├─ 需要流处理 + 消息队列?→ Kafka
│
└─ 需要延迟消息 + 高吞吐?→ RocketMQ(内置)或 RabbitMQ(插件)参考资料