欢迎光临
我们一直在努力

消息队列篇——RabbitMQ

RabbitMQ 使用

本文系统讲解 RabbitMQ 消息队列的使用方法,涵盖 AMQP 协议、核心概念、四种 Exchange 类型、消息可靠性保障、死信队列、延迟队列、集群与镜像队列、常用设计模式、性能调优、监控等内容。每个知识点均配有可运行的 Java 代码实例和ASCII 架构图,适合从零开始系统学习。


目录

一、RabbitMQ 概述

二、AMQP 协议与核心概念

三、安装部署

四、管理界面

五、快速入门(Hello World)

六、Work Queue 工作队列模式

七、Fanout Exchange(扇出)

八、Direct Exchange(直连)

九、Topic Exchange(主题)

十、Headers Exchange(头匹配)

十一、Exchange 类型对比

十二、消息可靠性——生产者端

十三、消息可靠性——消费者端

十四、死信队列(DLX)

十五、延迟队列

十六、消息优先级与惰性队列

十七、集群与镜像队列

十八、常用设计模式

十九、Spring Boot 整合实战

二十、性能调优

二十一、监控与运维

二十二、常见问题与最佳实践

二十三、总结与速查表


一、RabbitMQ 概述

1.1 RabbitMQ 简介

RabbitMQ 是一款由 Erlang 语言编写的开源消息代理(Message Broker),实现了 AMQP(Advanced Message Queuing Protocol)协议。它最初由英国公司 RabbitMQ Technologies 开发,目前归属 VMware/Pivotal 旗下。

核心特点:

特点说明
开源免费 Mozilla Public License
跨平台跨语言 支持 Java、Python、Go、C#、PHP、Ruby 等
高可靠 支持消息持久化、ACK 确认、镜像队列
高可用 支持集群、队列镜像、故障自动转移
插件丰富 管理界面、延迟插件、STOMP、MQTT 等
协议兼容 AMQP 0-9-1、STOMP、MQTT、AMQP 1.0
社区活跃 文档齐全,生态成熟

1.2 为什么需要消息队列

没有消息队列时(同步调用):

用户请求 → 订单系统 ──→ 库存系统(耗时 200ms)

├──→ 支付系统(耗时 300ms)

├──→ 通知系统(耗时 100ms)

└──→ 积分系统(耗时 150ms)

总耗时 = 200 + 300 + 100 + 150 = 750ms
问题:耦合度高、同步阻塞、任一系统宕机则失败

有消息队列后(异步解耦):

用户请求 → 订单系统 ──→ [写入 MQ] ──→ 返回(总耗时 ~50ms)

┌──────────┼──────────┐
▼ ▼ ▼
库存系统 支付系统 通知系统 / 积分系统
(各自消费,互不影响)

优势:解耦、异步、削峰填谷、故障隔离

消息队列三大核心作用:

作用说明示例
解耦 生产者与消费者解耦,互不依赖 订单系统只管发消息,库存/积分/短信各自消费
异步 耗时操作异步化,快速响应 用户注册后异步发送邮件/短信
削峰填谷 高峰时缓冲请求,匀速消费 秒杀场景:10万请求 → 队列 → 按 1000/s 消费

RabbitMQ 的独立特性:

特性说明典型场景
灵活路由 四种 Exchange(Fanout/Direct/Topic/Headers)支持复杂路由规则 按业务类型、通配符路由到不同队列
多协议支持 原生支持 AMQP 0-9-1、STOMP、MQTT、AMQP 1.0 兼容 IoT 设备(MQTT)、Web 端(STOMP)
插件生态 延迟队列插件、管理界面、联邦/镜像等开箱即用 快速扩展消息队列能力

1.3 主流消息队列对比

特性RocketMQRabbitMQKafka说明
开发语言 Java Erlang Scala/Java
协议 自研协议 AMQP 自研协议 RocketMQ 支持 HTTP/自定义协议
吞吐量 高(10万级/秒) 中(万级/秒) 极高(百万级/秒) 单机吞吐
延迟 毫秒级 微秒级 毫秒级 RabbitMQ 延迟最低
顺序消息 支持(分区强顺序) 支持(需单队列) 支持(分区内有序) 均为局部有序
事务消息 支持(半消息机制) 有限支持(事务/发布确认) 支持(事务性 API) 实现机制不同
延迟消息 支持(18 个级别) 支持(TTL+死信/插件) 不支持(需自研) RocketMQ 内置延迟级别
消息回溯 支持(按时间/偏移量) 有限支持 支持(按 Offset) 重放历史消息
死信队列 支持(原生) 支持(原生 DLX) 不支持(需自研)
消息轨迹 支持(内置) RocketMQ 原生支持
管理界面 支持(Dashboard) 支持(Management) 支持(第三方工具)
社区活跃度 高(阿里主导) 极高

选型建议:

  • 需要顺序消息 + 事务消息 + 延迟消息,Java 技术栈 → 选 RocketMQ
  • **追求极低延迟、复杂路由规则、简单场景 **→ 选 RabbitMQ
  • 追求极高吞吐、日志/流处理场景 → 选 Kafka

二、AMQP 协议与核心概念

2.1 AMQP 协议

AMQP(Advanced Message Queuing Protocol)是一个应用层标准协议,规定了消息中间件的交互模型。RabbitMQ 实现了 AMQP 0-9-1。

AMQP 模型:

Publisher(生产者)

│ ① Publish 发布消息

┌──────────────────────────────────┐
│ Exchange(交换机) │ ② 根据路由规则分发
│ │
│ ┌─────┐ ┌─────┐ ┌─────┐ │
│ │ Bnd │ │ Bnd │ │ Bnd │ 绑定 │
│ └──┬──┘ └──┬──┘ └──┬──┘ │
└─────┼────────┼────────┼───────────┘
│ ③ │ │
▼ ▼ ▼
┌───────┐ ┌───────┐ ┌───────┐
│Queue 1│ │Queue 2│ │Queue 3│ 消息队列
└───┬───┘ └───┬───┘ └───┬───┘
│ │ │ ④ Consume
▼ ▼ ▼
Consumer Consumer Consumer(消费者)

2.2 核心概念

概念英文说明
生产者 Producer/Publisher 发送消息的程序
消费者 Consumer 接收消息的程序
交换机 Exchange 接收消息,根据路由规则分发到队列
队列 Queue 存储消息的缓冲区,FIFO
绑定 Binding Exchange 和 Queue 之间的关联关系
路由键 Routing Key 生产者发送时指定的路由标识
绑定键 Binding Key 绑定时指定的匹配规则
连接 Connection TCP 长连接
通道 Channel 连接内的虚拟轻量连接,复用 TCP
虚拟主机 Virtual Host (vhost) 逻辑隔离,不同 vhost 相互独立
Broker RabbitMQ 服务节点

2.3 消息流转全流程

① 生产者创建消息


② 生产者连接到 Broker,打开 Channel


③ 生产者将消息发送到 Exchange
│ 消息 = Properties(headers) + Body(payload)

④ Exchange 根据 Type 和 Binding 规则,决定路由到哪些 Queue


⑤ Queue 接收消息,存储(内存/磁盘)


⑥ Consumer 从 Queue 获取消息(push 或 pull)


⑦ Consumer 处理完成后发送 ACK


⑧ Broker 收到 ACK,删除该消息


(如果未收到 ACK 且消费者断开 → 消息重新入队)

2.4 Connection 与 Channel

┌─────────────────────────────────────┐
│ Connection(TCP) │
│ │
│ ┌────────┐ ┌────────┐ ┌────────┐ │
│ │Channel1│ │Channel2│ │Channel3│ │ ← 多个 Channel 复用一个 TCP 连接
│ └────────┘ └────────┘ └────────┘ │
└─────────────────────────────────────┘

为什么需要 Channel?
– TCP 连接创建/销毁开销大(三次握手 + TLS)
– Channel 是轻量的虚拟连接
– 同一连接上可创建多个 Channel 并发操作
– AMQP 的所有操作(声明队列/发布/消费)都在 Channel 上执行


三、安装部署

3.1 Docker 安装(推荐)

# ① 拉取镜像(含管理界面)
docker pull rabbitmq:3.13-management

# ② 启动容器
docker run -d \\
–name rabbitmq \\
-p 5672:5672 \\ # AMQP 通信端口
-p 15672:15672 \\ # 管理界面端口
-p 25672:25672 \\ # 集群通信端口
-e RABBITMQ_DEFAULT_USER=admin \\
-e RABBITMQ_DEFAULT_PASS=admin123 \\
-v rabbitmq_data:/var/lib/rabbitmq \\
–restart=unless-stopped \\
rabbitmq:3.13-management

# ③ 访问管理界面
# 浏览器打开 http://localhost:15672
# 用户名: admin 密码: admin123

3.2 Linux 安装(RPM 方式)

# ① 安装 Erlang(RabbitMQ 依赖)
# CentOS 7/8 使用 Erlang Solutions 仓库
rpm –import https://packages.erlang-solutions.com/rpm/erlang_solutions.asc
cat > /etc/yum.repos.d/erlang.repo << 'EOF'
[erlang-solutions]
name=erlang-solutions
baseurl=https://packages.erlang-solutions.com/rpm/centos/$releasever/$basearch
gpgcheck=1
gpgkey=https://packages.erlang-solutions.com/rpm/erlang_solutions.asc
enabled=1
EOF

yum install erlang -y

# ② 下载并安装 RabbitMQ
wget https://github.com/rabbitmq/rabbitmq-server/releases/download/v3.13.0/rabbitmq-server-3.13.0-1.el8.noarch.rpm
rpm –import https://www.rabbitmq.com/rabbitmq-release-signing-key.asc
rpm -ivh rabbitmq-server-3.13.0-1.el8.noarch.rpm

# ③ 启动
systemctl start rabbitmq-server
systemctl enable rabbitmq-server

# ④ 开启管理界面插件
rabbitmq-plugins enable rabbitmq_management

# ⑤ 创建管理员用户
rabbitmqctl add_user admin admin123
rabbitmqctl set_user_tags admin administrator
rabbitmqctl set_permissions -p / admin ".*" ".*" ".*"

3.3 核心端口说明

端口用途说明
5672 AMQP 客户端通信(生产者/消费者连接)
15672 HTTP API / 管理界面 Web 管理控制台
25672 集群通信 节点间通信
61613 STOMP STOMP 协议客户端
1883 MQTT MQTT 协议客户端
15674 WebSocket STOMP over WebSocket

3.4 常用命令

# 服务管理
rabbitmq-server # 前台启动
rabbitmq-server -detached # 后台启动
rabbitmqctl stop # 停止
rabbitmqctl status # 查看状态

# 用户管理
rabbitmqctl add_user <user> <pass>
rabbitmqctl delete_user <user>
rabbitmqctl change_password <user> <newpass>
rabbitmqctl set_user_tags <user> administrator
rabbitmqctl list_users

# vhost 管理
rabbitmqctl add_vhost <vhost>
rabbitmqctl delete_vhost <vhost>
rabbitmqctl list_vhosts

# 权限管理
rabbitmqctl set_permissions -p <vhost> <user> <conf> <write> <read>
# 示例:rabbitmqctl set_permissions -p / admin ".*" ".*" ".*"

# 队列/交换机管理
rabbitmqctl list_queues
rabbitmqctl list_exchanges
rabbitmqctl list_bindings

# 插件管理
rabbitmq-plugins list # 查看已安装插件
rabbitmq-plugins enable <plugin> # 启用插件
rabbitmq-plugins disable <plugin> # 禁用插件


四、管理界面

4.1 管理界面概览

访问 http://localhost:15672,使用创建的账号登录后看到:

┌───────────────────────────────────────────────────────┐
│ RabbitMQ Management [Logout]│
├───────────────────────────────────────────────────────┤
│ Overview | Connections | Channels | Exchanges | │
│ Queues | Admin │
├───────────────────────────────────────────────────────┤
│ │
│ ┌─ Overview ──────────────────────────────────────┐ │
│ │ │ │
│ │ ┌─────────────┐ Rates: Published 500/s │ │
│ │ │ 0 messages │ Delivered 480/s │ │
│ │ │ (ready) │ Acknowledged 480/s │ │
│ │ └─────────────┘ │ │
│ │ │ │
│ │ Nodes: rabbit@hostname Running (uptime 3d) │ │
│ │ Ports: 5672 (amqp), 15672 (management) │ │
│ └───────────────────────────────────────────────────┘ │
│ │
│ ┌─ Queues ─────────────────────────────────────────┐ │
│ │ Name | Messages | Consumers | State │ │
│ │ task_queue | 0 | 3 | running │ │
│ │ dead_letter | 0 | 0 | running │ │
│ └──────────────────────────────────────────────────┘ │
└───────────────────────────────────────────────────────┘

4.2 核心页面

页面功能
Overview 全局概览:消息速率、节点状态
Connections 所有 TCP 连接列表
Channels 所有 Channel 列表
Exchanges 交换机管理:创建/删除/查看绑定
Queues 队列管理:创建/删除/查看消息/Purge
Admin 用户/vhost/策略/限制管理

4.3 常用操作

创建 Queue:

参数说明默认值
Name 队列名
Durable 是否持久化(重启不丢) No
Auto delete 最后一个消费者断开后自动删除 No
Arguments x-message-ttl / x-max-priority 等

创建 Exchange:

参数说明
Name 交换机名
Type direct / fanout / topic / headers
Durable 是否持久化
Auto delete 无绑定后自动删除
Internal 是否内部使用(不能被生产者直接发消息)

五、快速入门(Hello World)

5.1 Maven 依赖

<dependencies>
<!– RabbitMQ Java Client –>
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.21.0</version>
</dependency>

<!– SLF4J(RabbitMQ 客户端依赖日志) –>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId>
<version>2.0.13</version>
</dependency>
</dependencies>

5.2 生产者

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

public class HelloWorldProducer {

// 队列名
private static final String QUEUE_NAME = "hello";

public static void main(String[] args) throws Exception {
// 1. 创建连接工厂
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPort(5672);
factory.setUsername("admin");
factory.setPassword("admin123");
factory.setVirtualHost("/");

// 2. 创建连接和 Channel
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

// 3. 声明队列(如果不存在则创建)
// 参数:队列名 | durable | exclusive | autoDelete | arguments
channel.queueDeclare(QUEUE_NAME, false, false, false, null);

// 4. 发送消息
String message = "Hello World!";
// 参数:exchange | routingKey | props | body
// 不指定 exchange 时,使用默认交换机(空字符串),routingKey 即队列名
channel.basicPublish("", QUEUE_NAME, null, message.getBytes());

System.out.println(" [x] Sent: " + message);
}
}
}

// 输出: [x] Sent: Hello World!

5.3 消费者

import com.rabbitmq.client.*;

public class HelloWorldConsumer {

private static final String QUEUE_NAME = "hello";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setUsername("admin");
factory.setPassword("admin123");

// 连接和 Channel 通常保持长连接,不关闭
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();

// 声明队列(与生产者一致)
channel.queueDeclare(QUEUE_NAME, false, false, false, null);

// 注册消费者(异步回调)
System.out.println(" [*] Waiting for messages. To exit press Ctrl+C");

// 参数:队列名 | autoAck | 回调
channel.basicConsume(QUEUE_NAME, true, (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
System.out.println(" [x] Received: " + message);
}, consumerTag -> {
System.out.println(" Consumer canceled: " + consumerTag);
});
}
}

// 输出:
// [*] Waiting for messages…
// [x] Received: Hello World!

5.4 流程图解

Producer Broker Consumer
│ │ │
│ 1. Connection + Channel │ │
│ ─────────────────────────> │
│ │ │
│ 2. queueDeclare("hello")│ │
│ ─────────────────────────> │
│ │ Queue "hello" created │
│ │ │
│ 3. basicPublish(msg) │ │
│ ─────────────────────────> │
│ │ Message stored in Queue│
│ │ │
│ │ 4. basicConsume("hello")│
│ │<─────────────────────────│
│ │ │
│ │ 5. Push message (async) │
│ │ ─────────────────────────>│
│ │ │
│ │ 6. ACK (autoAck=true) │
│ │<─────────────────────────│
│ │ Message deleted │


六、Work Queue 工作队列模式

6.1 模式说明

┌─── Consumer 1(处理 msg1)
Producer → [Queue]──┼─── Consumer 2(处理 msg2)
└─── Consumer 3(处理 msg3)

特点:一条消息只能被一个消费者消费
作用:将任务分发给多个消费者,提高处理速度
分发策略:轮询(默认)/ 公平分发

6.2 生产者

public class WorkQueueProducer {

private static final String TASK_QUEUE_NAME = "task_queue";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

// 声明持久化队列
boolean durable = true;
channel.queueDeclare(TASK_QUEUE_NAME, durable, false, false, null);

// 发送多条消息,模拟任务
for (int i = 1; i <= 20; i++) {
String message = "Task #" + i;
// 设置消息持久化
channel.basicPublish("",
TASK_QUEUE_NAME,
MessageProperties.PERSISTENT_TEXT_PLAIN,
message.getBytes());
System.out.println(" [x] Sent: " + message);
Thread.sleep(100);
}
}
}
}

// 输出:
// [x] Sent: Task #1
// [x] Sent: Task #2
// …
// [x] Sent: Task #20

6.3 消费者(公平分发)

public class WorkQueueConsumer {

private static final String TASK_QUEUE_NAME = "task_queue";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();

channel.queueDeclare(TASK_QUEUE_NAME, true, false, false, null);
System.out.println(" [*] Waiting for messages…");

// ★ 公平分发:每次只分发1条消息,处理完ACK后才发下一条
int prefetchCount = 1;
channel.basicQos(prefetchCount);

// ★ 手动 ACK(autoAck = false)
channel.basicConsume(TASK_QUEUE_NAME, false, (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
System.out.println(" [x] Received: " + message);

try {
// 模拟处理耗时
doWork(message);
} finally {
System.out.println(" [x] Done: " + message);
// 手动确认
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
}
}, consumerTag -> {});
}

private static void doWork(String task) throws InterruptedException {
// 模拟处理:消息中的 '.' 数量代表耗时
int timeNeeded = task.chars().filter(c -> c == '.').count()... ;
Thread.sleep(task.length() * 100L);
}
}

// Consumer 1 输出:
// [*] Waiting for messages…
// [x] Received: Task #1
// [x] Done: Task #1
// [x] Received: Task #3
// …

// Consumer 2 输出:
// [*] Waiting for messages…
// [x] Received: Task #2
// [x] Done: Task #2
// [x] Received: Task #4
// …

6.4 轮询 vs 公平分发

轮询(Round-Robin,默认):
prefetch = unlimited
Consumer1 ← msg1, msg3, msg5, msg7, …
Consumer2 ← msg2, msg4, msg6, msg8, …
问题:如果 Consumer1 处理慢,Consumer2 处理快 → 负载不均

公平分发(Fair Dispatch):
channel.basicQos(1); ← 每个消费者最多1条未确认消息
Consumer1 ← msg1 (处理中)
Consumer2 ← msg2 (处理中)
Consumer1 完成 ACK → 收到 msg3
Consumer2 完成 ACK → 收到 msg4
谁快谁多收 → 按能力分配

对比轮询公平分发
prefetch 不限 1(或自定义)
分配方式 交替 按能力
吞吐量 中(受 ACK 影响)
负载均衡

七、Fanout Exchange(扇出)

7.1 模式说明

┌──→ Queue1 → Consumer A
Producer → [Fanout]┼──→ Queue2 → Consumer B
└──→ Queue3 → Consumer C

特点:广播模式,不关心 Routing Key
所有绑定到该 Exchange 的 Queue 都会收到消息
适用场景:群发通知、日志广播

7.2 生产者

public class FanoutProducer {

private static final String EXCHANGE_NAME = "logs_fanout";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

// 声明 Fanout Exchange
channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT);

String message = "广播消息:系统将于今晚23:00维护";
// Fanout 模式下 routingKey 无意义,可以传空
channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes());

System.out.println(" [x] Broadcast: " + message);
}
}
}

// 输出: [x] Broadcast: 广播消息:系统将于今晚23:00维护

7.3 消费者

public class FanoutConsumer {

private static final String EXCHANGE_NAME = "logs_fanout";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

Connection connection = factory.newConnection();
Channel channel = connection.createChannel();

// 声明 Exchange
channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT);

// ★ 创建临时队列(队列名由 RabbitMQ 随机生成)
// 临时队列特点:exclusive=true,消费者断开后自动删除
String queueName = channel.queueDeclare().getQueue();
System.out.println(" [*] Temp queue: " + queueName);

// ★ 将临时队列绑定到 Exchange
// Fanout 模式下 bindingKey 无意义
channel.queueBind(queueName, EXCHANGE_NAME, "");

System.out.println(" [*] Waiting for broadcast messages…");

channel.basicConsume(queueName, true, (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
System.out.println(" [x] Received: " + message);
}, consumerTag -> {});
}
}

// Consumer A 输出:
// [*] Temp queue: amq.gen-abc123…
// [x] Received: 广播消息:系统将于今晚23:00维护

// Consumer B 输出:
// [*] Temp queue: amq.gen-def456…
// [x] Received: 广播消息:系统将于今晚23:00维护


八、Direct Exchange(直连)

8.1 模式说明

┌──→ Queue (info) → Consumer A (info 日志)
Producer → [Direct] ───┼──→ Queue (error) → Consumer B (error 日志)
└──→ Queue (warn) → Consumer C (warn 日志)

特点:Routing Key 精确匹配 Binding Key
消息只会路由到完全匹配的 Queue
适用场景:按类型路由消息

8.2 生产者

public class DirectProducer {

private static final String EXCHANGE_NAME = "logs_direct";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT);

// 发送不同级别的日志
String[] levels = {"info", "warn", "error", "debug"};
for (String level : levels) {
String message = level + " log message at " + System.currentTimeMillis();
// ★ routingKey 精确匹配
channel.basicPublish(EXCHANGE_NAME, level, null, message.getBytes());
System.out.println(" [x] Sent [" + level + "]: " + message);
}
}
}
}

// 输出:
// [x] Sent [info]: info log message at 1723534…
// [x] Sent [warn]: warn log message at 1723534…
// [x] Sent [error]: error log message at 1723534…
// [x] Sent [debug]: debug log message at 1723534…

8.3 消费者

public class DirectConsumer {

private static final String EXCHANGE_NAME = "logs_direct";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

Connection connection = factory.newConnection();
Channel channel = connection.createChannel();

channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.DIRECT);

// 创建队列并绑定多个 routingKey
String queueName = "error_and_warn_queue";
channel.queueDeclare(queueName, false, false, false, null);

// ★ 同一个队列绑定多个 routingKey
channel.queueBind(queueName, EXCHANGE_NAME, "error");
channel.queueBind(queueName, EXCHANGE_NAME, "warn");
// 不绑定 info/debug → 不会收到

System.out.println(" [*] Waiting for error/warn messages…");

channel.basicConsume(queueName, true, (consumerTag, delivery) -> {
String routingKey = delivery.getEnvelope().getRoutingKey();
String message = new String(delivery.getBody(), "UTF-8");
System.out.println(" [x] Received [" + routingKey + "]: " + message);
}, consumerTag -> {});
}
}

// 输出:
// [*] Waiting for error/warn messages…
// [x] Received [error]: error log message at 1723534…
// [x] Received [warn]: warn log message at 1723534…
// (info 和 debug 被丢弃,因为没有队列绑定它们)

8.4 多队列绑定同一 Key

// 队列 A 和 B 都绑定 routingKey = "error"
channel.queueBind("queue_a", EXCHANGE_NAME, "error");
channel.queueBind("queue_b", EXCHANGE_NAME, "error");

// 此时发送 routingKey = "error" 的消息
// → queue_a 和 queue_b 都会收到(各自独立的副本)
// 注意:Fanout = 全部收到(不管Key),Direct = 按Key精确匹配


九、Topic Exchange(主题)

9.1 模式说明

Producer → [Topic] ─── routingKey = "order.create.vip" ──→ 匹配的队列

特点:支持通配符模式匹配
* = 匹配一个单词
# = 匹配零个或多个单词
适用场景:灵活路由,如日志分类、订单分类

9.2 通配符规则

routingKey 必须用点号 '.' 分隔成多个单词,如 "order.create.vip"

匹配规则:
* (星号) → 匹配恰好一个单词
# (井号) → 匹配零个或多个单词

示例:
routingKey = "order.create.vip"

binding pattern = "order.create.vip" ✅ 精确匹配
binding pattern = "order.create.*" ✅ * 匹配 "vip"
binding pattern = "order.#" ✅ # 匹配 "create.vip"
binding pattern = "order.create" ❌ 不匹配(单词数不同)
binding pattern = "#" ✅ 匹配所有(全捕获)
binding pattern = "*.create.*" ✅ 匹配中间是 create 的
binding pattern = "#.vip" ✅ # 匹配 "order.create"

9.3 生产者

public class TopicProducer {

private static final String EXCHANGE_NAME = "logs_topic";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.TOPIC);

// 发送不同主题的日志
String[] messages = {
"kern.critical", // 内核-严重
"auth.info", // 认证-信息
"kern.info", // 内核-信息
"cron.critical", // 定时任务-严重
"auth.warning", // 认证-警告
"kern.debug" // 内核-调试
};

for (String routingKey : messages) {
String body = "Log: " + routingKey + " at " + System.currentTimeMillis();
channel.basicPublish(EXCHANGE_NAME, routingKey, null, body.getBytes());
System.out.println(" [x] Sent [" + routingKey + "]: " + body);
}
}
}
}

9.4 消费者

public class TopicConsumer {

private static final String EXCHANGE_NAME = "logs_topic";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

Connection connection = factory.newConnection();
Channel channel = connection.createChannel();

channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.TOPIC);

// 创建两个队列,绑定不同的通配符模式
String queue1 = "kern_logs";
String queue2 = "critical_logs";

channel.queueDeclare(queue1, false, false, false, null);
channel.queueDeclare(queue2, false, false, false, null);

// Queue1 只收内核日志
channel.queueBind(queue1, EXCHANGE_NAME, "kern.*");
// Queue2 收所有 critical 日志
channel.queueBind(queue2, EXCHANGE_NAME, "*.critical");

// 还可以绑定全捕获
// channel.queueBind(queueAll, EXCHANGE_NAME, "#");

// 消费 Queue1
channel.basicConsume(queue1, true, (tag, delivery) -> {
System.out.println(" [Kern] " + new String(delivery.getBody(), "UTF-8"));
}, tag -> {});

// 消费 Queue2
channel.basicConsume(queue2, true, (tag, delivery) -> {
System.out.println(" [Critical] " + new String(delivery.getBody(), "UTF-8"));
}, tag -> {});

System.out.println(" [*] Waiting for topic logs…");
}
}

// 输出:
// [*] Waiting for topic logs…
// [Kern] Log: kern.critical at …
// [Critical] Log: kern.critical at … ← kern.critical 同时匹配两个队列
// [Kern] Log: kern.info at …
// [Critical] Log: cron.critical at …
// [Kern] Log: kern.debug at …
// (auth.info 和 auth.warning 不匹配 kern.*,不进入 Queue1)
// (auth.warning 不匹配 *.critical,不进入 Queue2)

9.5 匹配矩阵

routingKeykern.**.critical# (全捕获)
kern.critical
kern.info
kern.debug
cron.critical
auth.info
auth.warning

十、Headers Exchange(头匹配)

10.1 模式说明

特点:不依赖 routingKey,根据消息头的 key-value 匹配
匹配方式:x-match = all(全部匹配)| any(任一匹配)
适用场景:需要多维度路由,routingKey 不够表达时

10.2 生产者

public class HeadersProducer {

private static final String EXCHANGE_NAME = "headers_exchange";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

// 声明 Headers Exchange
Map<String, Object> headersArgs = new HashMap<>();
channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.HEADERS, true, false, false, null);

// 消息1:带 format=json 和 type=report
Map<String, Object> headers1 = new HashMap<>();
headers1.put("format", "json");
headers1.put("type", "report");

AMQP.BasicProperties props1 = new AMQP.BasicProperties.Builder()
.headers(headers1)
.build();
channel.basicPublish(EXCHANGE_NAME, "", props1, "JSON report".getBytes());

// 消息2:带 format=xml 和 type=notification
Map<String, Object> headers2 = new HashMap<>();
headers2.put("format", "xml");
headers2.put("type", "notification");

AMQP.BasicProperties props2 = new AMQP.BasicProperties.Builder()
.headers(headers2)
.build();
channel.basicPublish(EXCHANGE_NAME, "", props2, "XML notification".getBytes());

System.out.println(" [x] Sent messages with headers");
}
}
}

10.3 消费者

public class HeadersConsumer {

private static final String EXCHANGE_NAME = "headers_exchange";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

Connection connection = factory.newConnection();
Channel channel = connection.createChannel();

channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.HEADERS, true, false, false, null);

// Queue1:匹配 format=json AND type=report(全部满足)
String queue1 = "json_report_queue";
channel.queueDeclare(queue1, true, false, false, null);

Map<String, Object> bindArgs1 = new HashMap<>();
bindArgs1.put("x-match", "all"); // 全部匹配
bindArgs1.put("format", "json");
bindArgs1.put("type", "report");
channel.queueBind(queue1, EXCHANGE_NAME, "", bindArgs1);

// Queue2:匹配 format=xml OR type=notification(任一满足)
String queue2 = "xml_or_notification_queue";
channel.queueDeclare(queue2, true, false, false, null);

Map<String, Object> bindArgs2 = new HashMap<>();
bindArgs2.put("x-match", "any"); // 任一匹配
bindArgs2.put("format", "xml");
bindArgs2.put("type", "notification");
channel.queueBind(queue2, EXCHANGE_NAME, "", bindArgs2);

channel.basicConsume(queue1, true, (tag, delivery) -> {
System.out.println(" [Queue1] " + new String(delivery.getBody(), "UTF-8"));
}, tag -> {});

channel.basicConsume(queue2, true, (tag, delivery) -> {
System.out.println(" [Queue2] " + new String(delivery.getBody(), "UTF-8"));
}, tag -> {});

System.out.println(" [*] Waiting for headers messages…");
}
}

// 输出:
// [*] Waiting for headers messages…
// [Queue1] JSON report ← format=json AND type=report ✅
// [Queue2] XML notification ← format=xml OR type=notification ✅


十一、Exchange 类型对比

类型路由规则routingKey 作用适用场景
Fanout 忽略 routingKey,广播 无意义 群发通知、事件广播
Direct 精确匹配 必须完全一致 按类型分发
Topic 通配符匹配 支持 * 和 # 灵活路由、日志分类
Headers 消息头匹配 无意义 多维度路由
默认 (name=“”) 直接路由到 queue 名 = queue 名 简单点对点

选择指南:

┌─ 需要广播所有消息? ──→ Fanout

┌─ 需要按固定类型路由? ──→ Direct

┌─ 需要灵活模式匹配? ──→ Topic

┌─ 需要多维度路由? ──→ Headers

└─ 简单点对点? ──→ 默认 Exchange(无 Exchange)


十二、消息可靠性——生产者端

12.1 消息丢失风险点

Producer → [Exchange] → [Queue] → Consumer
↑ ↑ ↑
风险点1 风险点2 风险点3

风险点1:消息到达 Exchange 后丢失(Exchange 不存在 / 网络问题)
风险点2:消息到达 Queue 后丢失(Queue 未持久化 / Exchange 到 Queue 路由失败)
风险点3:消费者处理失败(消费者宕机 / 处理异常)

12.2 Confirm 模式(确认消息到达 Broker)

public class ConfirmProducer {

private static final int MESSAGE_COUNT = 10;

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

String queueName = "confirm_test";
channel.queueDeclare(queueName, true, false, false, null);

// ★ 开启 Publisher Confirm
channel.confirmSelect();

// 方式一:同步等待(简单但性能差)
for (int i = 0; i < MESSAGE_COUNT; i++) {
String message = "Confirm message #" + i;
channel.basicPublish("", queueName, null, message.getBytes());
// 等待确认(阻塞)
if (channel.waitForConfirms()) {
System.out.println(" [x] Confirmed: " + message);
} else {
System.out.println(" [!] Not confirmed: " + message);
// 重发…
}
}
}
}
}

// 输出:
// [x] Confirmed: Confirm message #0
// [x] Confirmed: Confirm message #1
// …
// [x] Confirmed: Confirm message #9

12.3 异步 Confirm(推荐)

public class AsyncConfirmProducer {

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

String queueName = "async_confirm_test";
channel.queueDeclare(queueName, true, false, false, null);

// ★ 开启 Confirm
channel.confirmSelect();

// ★ 注册异步回调
SortedSet<Long> unconfirmed = Collections.synchronizedSortedSet(new TreeSet<>());

channel.addConfirmListener(
// ack 回调:deliveryTag 之前的消息全部确认
(deliveryTag, multiple) -> {
if (multiple) {
// 批量确认
unconfirmed.headSet(deliveryTag, true).clear();
System.out.println(" [x] Batch confirmed up to: " + deliveryTag);
} else {
unconfirmed.remove(deliveryTag);
System.out.println(" [x] Confirmed: " + deliveryTag);
}
},
// nack 回调:deliveryTag 之前的消息全部未确认
(deliveryTag, multiple) -> {
System.err.println(" [!] Nack: " + deliveryTag + ", multiple=" + multiple);
// 重发逻辑…
}
);

// 批量发送
for (int i = 0; i < 100; i++) {
long seqNo = channel.getNextPublishSeqNo();
unconfirmed.add(seqNo);
String message = "Async confirm #" + i;
channel.basicPublish("", queueName, null, message.getBytes());
}

// 等待所有 confirm 回调完成
while (!unconfirmed.isEmpty()) {
Thread.sleep(100);
}
System.out.println("All messages confirmed.");
}
}
}

12.4 Return 模式(消息不可路由时退回)

public class ReturnProducer {

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

channel.exchangeDeclare("return_test", BuiltinExchangeType.DIRECT);

// ★ 注册 Return 回调(消息无法路由时触发)
channel.addReturnListener(returnMessage -> {
System.err.println(" [!] Return: " + new String(returnMessage.getBody()));
System.err.println(" ReplyCode: " + returnMessage.getReplyCode());
System.err.println(" ReplyText: " + returnMessage.getReplyText());
System.err.println(" RoutingKey: " + returnMessage.getRoutingKey());
// 记录日志 / 存库 / 重发…
});

// ★ 开启 mandatory(消息不可路由时触发 Return 而非丢弃)
boolean mandatory = true;

// 发送到不存在的 routingKey → 触发 Return
channel.basicPublish("return_test", "nonexistent.key", mandatory, null, "test".getBytes());
System.out.println(" [x] Sent message");

Thread.sleep(1000); // 等待回调
}
}
}

// 输出:
// [x] Sent message
// [!] Return: test
// ReplyCode: 312
// ReplyText: NO_ROUTE
// RoutingKey: nonexistent.key

12.5 生产者可靠性总结

机制作用性能影响适用场景
Publisher Confirm 确认消息到达 Broker 低(异步) 大多数场景
Mandatory + Return 不可路由时退回 确保路由成功
事务模式 确认+回滚 高(阻塞) ❌ 不推荐
持久化 消息写入磁盘 防止宕机丢消息

十三、消息可靠性——消费者端

13.1 手动 ACK

public class ManualAckConsumer {

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

Connection connection = factory.newConnection();
final Channel channel = connection.createChannel();

String queueName = "manual_ack_queue";
channel.queueDeclare(queueName, true, false, false, null);

// QoS:每次只推送1条未确认消息
channel.basicQos(1);

// ★ autoAck = false:手动确认
channel.basicConsume(queueName, false, (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
System.out.println(" [x] Received: " + message);

try {
// 模拟处理
doWork(message);

// ★ 正常确认
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
System.out.println(" [x] ACK: " + message);

} catch (Exception e) {
System.err.println(" [!] Processing failed: " + message);

// ★ 处理失败:NACK 并 requeue
// 参数:deliveryTag | multiple | requeue
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true);
// 或使用 basicReject(不支持 multiple)
// channel.basicReject(delivery.getEnvelope().getDeliveryTag(), true);
}
}, consumerTag -> {
System.out.println(" Consumer canceled");
});
}

private static void doWork(String task) throws Exception {
Thread.sleep(1000); // 模拟耗时
}
}

13.2 ACK 方式对比

方法说明multiple适用
basicAck 正面确认,Broker 删除消息 支持 处理成功
basicNack 负面确认,可选择 requeue 支持 处理失败
basicReject 负面确认,可选择 requeue 不支持 处理失败

13.3 消费者自动重连与 QoS

ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPort(5672);
factory.setUsername("admin");
factory.setPassword("admin123");

// ★ 自动恢复配置
factory.setAutomaticRecoveryEnabled(true); // 开启自动重连
factory.setNetworkRecoveryInterval(5000); // 重连间隔 5s
factory.setTopologyRecoveryEnabled(true); // 恢复 Exchange/Queue/Binding

// ★ 心跳
factory.setRequestedHeartbeat(30); // 30秒心跳

// ★ 连接超时
factory.setConnectionTimeout(10000); // 10秒超时

消费者可靠性保障清单:
✅ autoAck = false(手动确认)
✅ basicQos(1)(公平分发,防止消息堆积)
✅ 处理成功 → basicAck
✅ 处理失败 → basicNack + requeue(或转入死信队列)
✅ 开启自动重连(AutomaticRecovery)
✅ 设置心跳检测
✅ 持久化队列 + 持久化消息


十四、死信队列(DLX)

14.1 死信概念

消息变成"死信"(Dead Letter)的条件:

条件说明
消息被消费者拒绝(basicNack/basicReject)且 requeue=false 消费者拒收且不重新入队
消息 TTL 过期 消息存活时间超过设置值
队列达到最大长度 队列满了,新消息被挤出

14.2 死信队列架构

正常流程:
Producer → [Exchange] → [Normal Queue] → Consumer

│ 消息变成死信

死信流程: ↓
[Normal Queue] ──(DLX绑定)──→ [Dead Letter Exchange]


[Dead Letter Queue]


Dead Letter Consumer(告警/存档/重处理)

14.3 完整示例

public class DLXDemo {

// 常量定义
private static final String NORMAL_EXCHANGE = "normal_exchange";
private static final String NORMAL_QUEUE = "normal_queue";
private static final String DLX_EXCHANGE = "dlx_exchange";
private static final String DLX_QUEUE = "dlx_queue";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

// ① 声明死信交换机和队列
channel.exchangeDeclare(DLX_EXCHANGE, BuiltinExchangeType.DIRECT, true);
channel.queueDeclare(DLX_QUEUE, true, false, false, null);
channel.queueBind(DLX_QUEUE, DLX_EXCHANGE, "dlx.routing.key");

// ② 声明正常队列,绑定死信交换机
Map<String, Object> args2 = new HashMap<>();
args2.put("x-dead-letter-exchange", DLX_EXCHANGE); // 死信交换机
args2.put("x-dead-letter-routing-key", "dlx.routing.key"); // 死信路由键
// args2.put("x-message-ttl", 10000); // 可选:消息TTL
// args2.put("x-max-length", 100); // 可选:最大队列长度

channel.exchangeDeclare(NORMAL_EXCHANGE, BuiltinExchangeType.DIRECT, true);
channel.queueDeclare(NORMAL_QUEUE, true, false, false, args2);
channel.queueBind(NORMAL_QUEUE, NORMAL_EXCHANGE, "normal.routing.key");

// ③ 发送消息
String message = "Test DLX message";
channel.basicPublish(NORMAL_EXCHANGE, "normal.routing.key",
MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes());
System.out.println(" [x] Sent: " + message);

// ④ 消费者拒绝消息(requeue=false → 变成死信)
channel.basicConsume(NORMAL_QUEUE, false, (tag, delivery) -> {
System.out.println(" [Consumer] Rejected: " + new String(delivery.getBody()));
// ★ reject 且 requeue=false → 进入死信队列
channel.basicReject(delivery.getEnvelope().getDeliveryTag(), false);
}, tag -> {});

Thread.sleep(1000);

// ⑤ 检查死信队列
// 在管理界面可以看到 DLX_QUEUE 中出现了被拒绝的消息
}
}
}

// 输出:
// [x] Sent: Test DLX message
// [Consumer] Rejected: Test DLX message
// (消息从 normal_queue 进入 dlx_queue)

14.4 TTL + DLX 实现延迟

// 给消息设置 TTL,过期后进入死信队列 → 实现延迟消费
public class TTLDLXProducer {

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

// 正常队列已绑定 DLX(同上示例)

// ★ 给单条消息设置 TTL = 10秒
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.deliveryMode(2) // 持久化
.expiration("10000") // TTL 10秒
.build();

channel.basicPublish("normal_exchange", "normal.routing.key",
props, "Delayed message".getBytes());

System.out.println(" [x] Sent delayed message (TTL=10s)");
// 10秒后消息自动进入死信队列
}
}
}


十五、延迟队列

15.1 方案一:TTL + DLX(通用方案)

Producer → [Queue (TTL=30s)] ──30秒后──→ [DLX Exchange] → [DLX Queue] → Consumer


消费者收到延迟消息(30秒后)

注意:TTL + DLX 方案的陷阱——队头阻塞问题
如果队列中有多条消息,TTL 是按队列头部消息过期来触发的
后面的消息即使先过期,也要等前面的消息过期后才能进入死信队列

15.2 方案二:rabbitmq_delayed_message_exchange 插件(推荐)

# ① 安装延迟插件
# 下载对应版本的 .ez 文件到 plugins 目录
# 下载地址:https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/releases

# Linux:
cp rabbitmq_delayed_message_exchange-3.13.0.ez /usr/lib/rabbitmq/plugins/
rabbitmq-plugins enable rabbitmq_delayed_message_exchange

# Docker:
docker exec rabbitmq rabbitmq-plugins enable rabbitmq_delayed_message_exchange
# 或在启动时挂载插件

15.3 插件方案代码

public class DelayedExchangeProducer {

private static final String EXCHANGE_NAME = "delayed_exchange";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

// ★ 声明延迟交换机(自定义类型 x-delayed-message)
Map<String, Object> exchangeArgs = new HashMap<>();
exchangeArgs.put("x-delayed-type", "direct"); // 底层路由类型
channel.exchangeDeclare(EXCHANGE_NAME, "x-delayed-message",
true, false, false, exchangeArgs);

// 声明队列并绑定
String queueName = "delayed_queue";
channel.queueDeclare(queueName, true, false, false, null);
channel.queueBind(queueName, EXCHANGE_NAME, "delay.key");

// ★ 发送延迟消息
// 消息1:延迟 5 秒
AMQP.BasicProperties props1 = new AMQP.BasicProperties.Builder()
.headers(Map.of("x-delay", 5000)) // 5秒延迟
.build();
channel.basicPublish(EXCHANGE_NAME, "delay.key", props1, "Message 1 (5s)".getBytes());

// 消息2:延迟 10 秒
AMQP.BasicProperties props2 = new AMQP.BasicProperties.Builder()
.headers(Map.of("x-delay", 10000)) // 10秒延迟
.build();
channel.basicPublish(EXCHANGE_NAME, "delay.key", props2, "Message 2 (10s)".getBytes());

// 消息3:延迟 3 秒
AMQP.BasicProperties props3 = new AMQP.BasicProperties.Builder()
.headers(Map.of("x-delay", 3000)) // 3秒延迟
.build();
channel.basicPublish(EXCHANGE_NAME, "delay.key", props3, "Message 3 (3s)".getBytes());

System.out.println(" [x] Sent 3 delayed messages");
}
}
}

// 消费者5秒后收到 Message 1,10秒后收到 Message 2,3秒后收到 Message 3
// (按延迟时间顺序到达,不受发送顺序影响)

15.4 延迟队列方案对比

方案原理优点缺点
TTL + DLX 消息过期进入死信队列 无需插件 队头阻塞、不支持多级延迟
Delayed Exchange 插件 Exchange 层面延迟 精确延迟、无队头阻塞 需安装插件
延迟队列工具(如 Redis ZSet) 外部排序 灵活 非 RabbitMQ 原生
定时任务扫描 轮询扫描 简单 不实时、扫描浪费

十六、消息优先级与惰性队列

16.1 消息优先级

public class PriorityProducer {

private static final String QUEUE_NAME = "priority_queue";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

// ★ 声明队列时设置最大优先级
Map<String, Object> args2 = new HashMap<>();
args2.put("x-max-priority", 10); // 最高优先级 10
channel.queueDeclare(QUEUE_NAME, true, false, false, args2);

// 发送不同优先级的消息
for (int i = 0; i < 10; i++) {
int priority = i % 5; // 优先级 0~4
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.priority(priority)
.build();

String message = "Message " + i + " (priority=" + priority + ")";
channel.basicPublish("", QUEUE_NAME, props, message.getBytes());
}

System.out.println(" [x] Sent 10 messages with priorities");
}
}
}

// 消费者会先收到高优先级的消息
// 注意:优先级仅在消费者跟不上生产者速度时才有效
// 如果消息一到就消费,优先级无意义

16.2 惰性队列(Lazy Queue)

默认队列 vs 惰性队列:

默认队列:
消息尽量保存在内存 → 吞吐量高,但内存占用大
队列过长时可能触发内存换页 → 性能下降

惰性队列(Lazy Queue):
消息直接写入磁盘 → 内存占用极低
消费时从磁盘读取 → 吞吐量低,但支持百万级消息堆积

// 声明惰性队列
Map<String, Object> args2 = new HashMap<>();
args2.put("x-queue-mode", "lazy"); // ★ 设置为惰性模式
channel.queueDeclare("lazy_queue", true, false, false, args2);

// RabbitMQ 3.12+ 默认所有队列都是惰性队列
// 不需要额外配置

对比默认队列惰性队列
消息存储 内存为主 磁盘为主
吞吐量 中低
内存占用 极低
适用场景 实时消息 大量堆积、离线消费

十七、集群与镜像队列

17.1 集群概述

RabbitMQ 集群有两种节点类型:

① 磁盘节点(Disc Node)—— 默认
持久化存储元数据(Exchange/Queue/Binding/User 等)
启动慢,但宕机不丢元数据

② 内存节点(RAM Node)
元数据只存内存,不写磁盘
启动快,但宕机丢元数据

集群要求:至少一个磁盘节点
推荐架构:1~2 个磁盘节点 + 多个内存节点

集群拓扑示例:

┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ node1 (disc) │ ←→ │ node2 (ram) │ ←→ │ node3 (ram) │
│ 10.0.0.1 │ │ 10.0.0.2 │ │ 10.0.0.3 │
└──────────────┘ └──────────────┘ └──────────────┘

Exchange / Binding / User 等元数据在所有节点同步
Queue 默认只在一个节点上创建(可通过镜像队列复制到其他节点)

17.2 搭建集群

# ① 修改 /etc/hosts(每个节点)
# 10.0.0.1 node1
# 10.0.0.2 node2
# 10.0.0.3 node3

# ② 同步 Erlang Cookie(所有节点必须一致)
# 在 node1 上:
scp /var/lib/rabbitmq/.erlang.cookie root@node2:/var/lib/rabbitmq/
scp /var/lib/rabbitmq/.erlang.cookie root@node3:/var/lib/rabbitmq/

# ③ 在 node2、node3 上加入集群
# node2:
rabbitmqctl stop_app
rabbitmqctl join_cluster rabbit@node1
rabbitmqctl start_app

# node3:
rabbitmqctl stop_app
rabbitmqctl join_cluster rabbit@node1
rabbitmqctl start_app

# ④ 查看集群状态
rabbitmqctl cluster_status
# 输出:
# Cluster status of node rabbit@node1 …
# [{nodes,[{disc,[rabbit@node1]},{ram,[rabbit@node2,rabbit@node3]}]},
# {running_nodes,[rabbit@node1,rabbit@node2,rabbit@node3]},
# {cluster_name,<<"rabbit@node1">>},
# …]

17.3 镜像队列(经典镜像)

镜像队列(Mirrored Queue):将队列复制到多个节点

┌──────────┐ ┌──────────┐ ┌──────────┐
│ node1 │ │ node2 │ │ node3 │
│ │ │ │ │ │
│ [Queue A]│←─镜像──│ [Queue A]│←─镜像──│ [Queue A]│
│ (master) │ │ (slave) │ │ (slave) │
│ │ │ │ │ │
│ [Queue B]│ │ [Queue B]│ │ [Queue B]│
│ (slave) │←─镜像──│ (master) │←─镜像──│ (slave) │
└──────────┘ └──────────┘ └──────────┘

Master 宕机 → Slave 自动提升为 Master → 高可用

# 创建镜像策略(所有节点上执行,只需一次)
# 将所有以 "ha." 开头的队列镜像到所有节点
rabbitmqctl set_policy ha-all "^ha\\." '{
"ha-mode":"all",
"ha-params": {},
"ha-sync-mode":"automatic"
}'

# 参数说明:
# ha-mode = all → 镜像到所有节点
# ha-mode = exactly → 镜像到指定数量节点(ha-params=N)
# ha-mode = nodes → 镜像到指定节点(ha-params=["node1@host", …])
# ha-sync-mode = automatic → 自动同步

17.4 仲裁队列(Quorum Queue,推荐替代镜像队列)

RabbitMQ 3.8+ 引入 Quorum Queue,基于 Raft 协议:

比经典镜像队列更可靠、更简单:

┌──────────┐ ┌──────────┐ ┌──────────┐
│ node1 │ │ node2 │ │ node3 │
│ │ │ │ │ │
│ [Q.Quorum]│←Raft─│ [Q.Quorum]│←Raft─│ [Q.Quorum]│
│ (leader) │ │(follower)│ │(follower)│
└──────────┘ └──────────┘ └──────────┘

特点:
– 基于 Raft 一致性协议,自动选主
– 数据强一致,不丢消息
– 无需配置 policy
– 不支持非持久化(必须 durable)

# 声明 Quorum Queue(通过参数)
# 在管理界面创建队列时设置 arguments:
# x-queue-type = quorum

# 代码方式
Map<String, Object> args = new HashMap<>();
args.put("x-queue-type", "quorum");
channel.queueDeclare("my_quorum_queue", true, false, false, args);

对比经典镜像队列仲裁队列
一致性 最终一致 强一致(Raft)
配置 需设 policy 声明时指定
推荐 逐渐弃用 ✅ 推荐
持久化 可选 强制
性能

十八、常用设计模式

18.1 六大核心模式总览

┌─────────────────────────────────────────────────────────┐
│ RabbitMQ 六大模式 │
│ │
│ ① Simple(简单模式) P → Queue → C │
│ 点对点,一个生产者一个消费者 │
│ │
│ ② Work(工作模式) P → Queue → C1/C2/C3 │
│ 一个生产者多个消费者,争抢消费 │
│ │
│ ③ Fanout(广播模式) P → Exchange → Q1/Q2/Q3 → C │
│ 一条消息多个消费者都收到 │
│ │
│ ④ Routing(路由模式) P → Exchange → Q(匹配Key) → C │
│ 按 routingKey 精确路由 │
│ │
│ ⑤ Topics(主题模式) P → Exchange → Q(通配符) → C │
│ 按 routingKey 模式匹配路由 │
│ │
│ ⑥ RPC(远程调用模式) Client → Queue → Server → │
│ reply_to 回调队列 │
│ 同步等待响应 │
└─────────────────────────────────────────────────────────┘

18.2 RPC 模式

public class RPCServer {

private static final String RPC_QUEUE_NAME = "rpc_queue";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

channel.queueDeclare(RPC_QUEUE_NAME, false, false, false, null);
channel.basicQos(1);

System.out.println(" [x] Awaiting RPC requests");

channel.basicConsume(RPC_QUEUE_NAME, false, (tag, delivery) -> {
String request = new String(delivery.getBody(), "UTF-8");
System.out.println(" [.] Received: " + request);

// 处理请求
String response = processRequest(request);

// ★ 通过 replyTo 队列返回结果
AMQP.BasicProperties replyProps = new AMQP.BasicProperties
.Builder
()
.correlationId(delivery.getProperties().getCorrelationId())
.build();

channel.basicPublish("", delivery.getProperties().getReplyTo(),
replyProps, response.getBytes());

// ACK 原始请求
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
}, consumerTag -> {});
}
}

private static String processRequest(String request) {
// 模拟:求斐波那契数
int n = Integer.parseInt(request);
return String.valueOf(fib(n));
}

private static long fib(int n) {
if (n <= 1) return n;
return fib(n 1) + fib(n 2);
}
}

// 输出:
// [x] Awaiting RPC requests
// [.] Received: 10
// [.] Received: 20

public class RPCClient {

private static final String RPC_QUEUE_NAME = "rpc_queue";

public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {

// ★ 创建临时回调队列
String replyQueueName = channel.queueDeclare().getQueue();

// 生成唯一 correlationId
String corrId = UUID.randomUUID().toString();

AMQP.BasicProperties props = new AMQP.BasicProperties
.Builder
()
.correlationId(corrId)
.replyTo(replyQueueName)
.build();

// 发送请求
String request = "10"; // 请求 fib(10)
channel.basicPublish("", RPC_QUEUE_NAME, props, request.getBytes());
System.out.println(" [x] Requesting fib(" + request + ")");

// ★ 阻塞等待响应
BlockingQueue<String> responseHolder = new LinkedBlockingQueue<>();

channel.basicConsume(replyQueueName, true, (tag, delivery) -> {
if (delivery.getProperties().getCorrelationId().equals(corrId)) {
String response = new String(delivery.getBody(), "UTF-8");
responseHolder.offer(response);
}
}, tag -> {});

String response = responseHolder.take(); // 阻塞等待
System.out.println(" [.] Got response: fib(" + request + ") = " + response);
}
}
}

// 输出:
// [x] Requesting fib(10)
// [.] Got response: fib(10) = 55

18.3 RPC 流程图

Client Server
│ │
│ 1. 声明 reply_to 临时队列 │
│ │
│ 2. 发送请求到 rpc_queue │
│ (corrId + replyTo + body) │
│────────────────────────────────>│
│ │
│ │ 3. 消费 rpc_queue
│ │ 处理请求
│ │
│ 4. 返回结果到 reply_to 队列 │
│ (corrId + response) │
│<────────────────────────────────│
│ │
│ 5. 检查 corrId 匹配 │
│ 获取响应 │


十九、Spring Boot 整合实战

19.1 Maven 依赖

<dependencies>
<!– Spring Boot Starter AMQP –>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<!– Spring Boot Web –>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
</dependencies>

19.2 application.yml

spring:
rabbitmq:
host: localhost
port: 5672
username: admin
password: admin123
virtual-host: /
# 消费者配置
listener:
simple:
acknowledge-mode: manual # 手动 ACK
prefetch: 1 # 公平分发
concurrency: 3 # 消费线程数
retry:
enabled: true
max-attempts: 3 # 重试3次
initial-interval: 1000 # 初始间隔1s
# 生产者配置
publisher-confirm-type: correlated # 异步 Confirm
publisher-returns: true # 开启 Return

# 自定义队列名
app:
queue:
order: order.queue
order-dlx: order.dlx.queue
order-delay: order.delay.queue
exchange:
order: order.exchange
order-dlx: order.dlx.exchange

19.3 配置类

import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class RabbitMQConfig {

// ============ 队列名 ============
public static final String ORDER_EXCHANGE = "order.exchange";
public static final String ORDER_QUEUE = "order.queue";
public static final String ORDER_DLX_EXCHANGE = "order.dlx.exchange";
public static final String ORDER_DLX_QUEUE = "order.dlx.queue";

// ============ 死信交换机和队列 ============
@Bean
public DirectExchange orderDlxExchange() {
return new DirectExchange(ORDER_DLX_EXCHANGE, true, false);
}

@Bean
public Queue orderDlxQueue() {
return QueueBuilder.durable(ORDER_DLX_QUEUE).build();
}

@Bean
public Binding orderDlxBinding() {
return BindingBuilder.bind(orderDlxQueue())
.to(orderDlxExchange())
.with("order.dlx");
}

// ============ 正常交换机和队列(绑定死信) ============
@Bean
public DirectExchange orderExchange() {
return new DirectExchange(ORDER_EXCHANGE, true, false);
}

@Bean
public Queue orderQueue() {
return QueueBuilder.durable(ORDER_QUEUE)
// ★ 绑定死信队列
.withArgument("x-dead-letter-exchange", ORDER_DLX_EXCHANGE)
.withArgument("x-dead-letter-routing-key", "order.dlx")
.build();
}

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

19.4 生产者

import org.springframework.amqp.core.MessageProperties;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;

@Component
public class OrderMessageProducer {

@Autowired
private RabbitTemplate rabbitTemplate;

/**
* 发送订单消息
*/

public void sendOrder(String orderId, String orderInfo) {
rabbitTemplate.convertAndSend(
RabbitMQConfig.ORDER_EXCHANGE,
"order.create",
orderInfo,
message -> {
// 设置消息属性
message.getMessageProperties().setMessageId(orderId);
message.getMessageProperties().setDeliveryMode(MessageProperties.DEFAULT_DELIVERY_MODE);
return message;
}
);
System.out.println("[Producer] Sent order: " + orderId);
}

/**
* 发送延迟消息(TTL + DLX)
*/

public void sendDelayedOrder(String orderId, String orderInfo, int delaySeconds) {
rabbitTemplate.convertAndSend(
RabbitMQConfig.ORDER_EXCHANGE,
"order.create",
orderInfo,
message -> {
message.getMessageProperties().setMessageId(orderId);
message.getMessageProperties().setExpiration(String.valueOf(delaySeconds * 1000));
return message;
}
);
System.out.println("[Producer] Sent delayed order: " + orderId + " (delay=" + delaySeconds + "s)");
}
}

19.5 消费者

import com.rabbitmq.client.Channel;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

@Component
public class OrderMessageConsumer {

/**
* 消费订单消息
*/

@RabbitListener(queues = RabbitMQConfig.ORDER_QUEUE)
public void handleOrder(String orderInfo, Message message, Channel channel) throws Exception {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
String messageId = message.getMessageProperties().getMessageId();

try {
System.out.println("[Consumer] Received order: " + messageId + ", info: " + orderInfo);

// 模拟处理
processOrder(orderInfo);

// ★ 手动 ACK
channel.basicAck(deliveryTag, false);
System.out.println("[Consumer] ACK: " + messageId);

} catch (Exception e) {
System.err.println("[Consumer] Error processing: " + messageId + ", " + e.getMessage());

// 判断是否需要重试
// 获取重试次数
Long retryCount = getRetryCount(message);
if (retryCount < 3) {
// ★ NACK 并重新入队(重试)
channel.basicNack(deliveryTag, false, true);
System.out.println("[Consumer] NACK + requeue: " + messageId);
} else {
// ★ 超过重试次数,拒绝(进入死信队列)
channel.basicReject(deliveryTag, false);
System.err.println("[Consumer] Reject to DLX: " + messageId);
}
}
}

/**
* 消费死信队列(告警/存档/人工处理)
*/

@RabbitListener(queues = RabbitMQConfig.ORDER_DLX_QUEUE)
public void handleDeadLetter(String deadInfo, Message message, Channel channel) throws Exception {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
System.err.println("[DLX Consumer] Dead letter: " + deadInfo);
// 存入数据库 → 人工处理
channel.basicAck(deliveryTag, false);
}

private void processOrder(String orderInfo) throws Exception {
// 模拟处理
Thread.sleep(500);
}

private Long getRetryCount(Message message) {
// 从 header 中获取重试次数(Spring AMQP 默认不设)
Object retryCount = message.getMessageProperties().getHeader("x-retry-count");
return retryCount != null ? (Long) retryCount : 0L;
}
}

19.6 Confirm 与 Return 回调

import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import jakarta.annotation.PostConstruct;

@Component
public class RabbitMQConfirmCallback implements RabbitTemplate.ConfirmCallback, RabbitTemplate.ReturnsCallback {

@Autowired
private RabbitTemplate rabbitTemplate;

@PostConstruct
public void init() {
rabbitTemplate.setConfirmCallback(this);
rabbitTemplate.setReturnsCallback(this);
}

/**
* Confirm 回调:消息到达 Broker
*/

@Override
public void confirm(CorrelationData correlationData, boolean ack, String cause) {
if (ack) {
System.out.println("[Confirm] Message confirmed: " +
(correlationData != null ? correlationData.getId() : "null"));
} else {
System.err.println("[Confirm] Message NOT confirmed: " + cause);
// 重发…
}
}

/**
* Return 回调:消息不可路由
*/

@Override
public void returnedMessage(ReturnedMessage returned) {
System.err.println("[Return] Message returned: " +
"code=" + returned.getReplyCode() +
", text=" + returned.getReplyText() +
", exchange=" + returned.getExchange() +
", routingKey=" + returned.getRoutingKey() +
", body=" + new String(returned.getMessage().getBody()));
// 记录日志 / 存库…
}
}

19.7 测试 Controller

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;

@RestController
@RequestMapping("/order")
public class OrderController {

@Autowired
private OrderMessageProducer producer;

@PostMapping("/create")
public String createOrder(@RequestParam String orderId, @RequestParam String info) {
producer.sendOrder(orderId, info);
return "Order sent: " + orderId;
}

@PostMapping("/delayed")
public String createDelayedOrder(@RequestParam String orderId,
@RequestParam String info,
@RequestParam int delay) {
producer.sendDelayedOrder(orderId, info, delay);
return "Delayed order sent: " + orderId + " (delay=" + delay + "s)";
}
}

// 测试:
// curl -X POST "http://localhost:8080/order/create?orderId=001&info=Test"
// 返回:Order sent: 001

// curl -X POST "http://localhost:8080/order/delayed?orderId=002&info=Delayed&delay=30"
// 返回:Delayed order sent: 002 (delay=30s)


二十、性能调优

20.1 生产者优化

ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");

// ★ 连接参数优化
factory.setAutomaticRecoveryEnabled(true);
factory.setConnectionTimeout(10000);
factory.setRequestedHeartbeat(30);

// ★ Channel 优化:多 Channel 并发
Connection connection = factory.newConnection();
Channel[] channels = new Channel[10];
for (int i = 0; i < 10; i++) {
channels[i] = connection.createChannel();
channels[i].confirmSelect(); // 每个 Channel 开启 Confirm
}

// ★ 批量发送:减少网络往返
for (int i = 0; i < 1000; i++) {
channel.basicPublish(exchange, routingKey, props, body);
// 不每条都等 confirm
}
// 批量等 confirm
if (channel.waitForConfirms(5000)) {
System.out.println("All confirmed");
}

// ★ 异步 Confirm(最高吞吐量)
channel.addConfirmListener(ackCallback, nackCallback);
for (int i = 0; i < 10000; i++) {
channel.basicPublish(exchange, routingKey, props, body);
}

20.2 消费者优化

// ★ 预取值调优
// prefetch = 1 → 公平分发,吞吐低
// prefetch = N → 批量推送,吞吐高
channel.basicQos(10); // 每个消费者最多 10 条未确认

// ★ 多消费者并发
// Spring Boot: spring.rabbitmq.listener.simple.concurrency=10
// spring.rabbitmq.listener.simple.max-concurrency=50

// ★ 批量 ACK
channel.basicAck(deliveryTag, true); // multiple=true 批量确认

20.3 Broker 端优化

参数配置说明
vm_memory_high_watermark 0.4~0.6 内存水位线,超过触发阻塞
disk_free_limit 1GB+ 磁盘剩余空间警戒线
channel_max 2047 单连接最大 Channel 数
heartbeat 30~60 心跳间隔
hipe_compile false HiPE 编译(已弃用)

# /etc/rabbitmq/rabbitmq.conf
vm_memory_high_watermark.relative = 0.6
disk_free_limit.absolute = 2GB
channel_max = 2047
heartbeat = 30
collect_statistics_interval = 5000

# 日志
log.console = true
log.console.level = info
log.file = true
log.file.level = warning

20.4 性能对比参考

场景吞吐量(消息/秒)说明
简单生产消费(非持久化) 50000~100000 无 ACK,无持久化
持久化 + 手动 ACK 5000~15000 最常见
异步 Confirm + 批量 20000~50000 高吞吐
Fanout 广播 10000~30000 取决于队列数
跨网络 1000~5000 取决于网络延迟

二十一、监控与运维

21.1 管理界面监控

管理界面关键指标:

Overview:
├── Queued Messages(Ready + Unacknowledged)
├── Message Rates(Publish / Deliver / Ack)
└── Global Counts(Connections / Channels / Queues / Consumers)

Queues 页面:
├── Ready → 等待消费的消息数
├── Unacked → 已推送未确认的消息数
├── Total → Ready + Unacked
├── Consumers → 消费者数量
├── Incoming → 入队速率
├── Deliver/get → 投递速率
└── Ack → 确认速率

21.2 命令行监控

# 查看队列状态
rabbitmqctl list_queues name messages messages_ready messages_unacknowledged consumers

# 输出示例:
# Listing queues for vhost / …
# name messages ready unacked consumers
# task_queue 0 0 0 3
# dead_letter 5 5 0 0

# 查看 Exchange
rabbitmqctl list_exchanges name type

# 查看连接
rabbitmqctl list_connections name peer_host peer_port state

# 查看 Channel
rabbitmqctl list_channels number user consumer_count messages_unacknowledged

# 查看消费者
rabbitmqctl list_consumers queue_name consumer_tag ack_required

# 查看节点状态
rabbitmqctl node_health_check

21.3 HTTP API 监控

# 获取队列信息
curl -u admin:admin123 http://localhost:15672/api/queues/%2F/task_queue

# 获取 Overview
curl -u admin:admin123 http://localhost:15672/api/overview

# 获取所有队列
curl -u admin:admin123 http://localhost:15672/api/queues

21.4 告警指标

指标告警阈值说明
Queue Ready 消息数 > 10000 消费者跟不上
Queue Unacked 消息数 > 1000 消费者处理慢
Consumer 数量 = 0 消费者全部掉线
Connection 数量 > 1000 连接泄漏
内存使用率 > 80% 可能阻塞
磁盘剩余 < 2GB 可能阻塞
消息 Nack 率 > 5% 消费异常

二十二、常见问题与最佳实践

22.1 消息重复消费

问题:消费者处理成功后 ACK 未到达 Broker(网络问题)→ Broker 重发 → 重复消费

解决:幂等性设计
① 消息唯一 ID + 数据库去重表
② Redis 记录已处理的消息 ID
③ 业务层幂等(如:订单状态机)

// 方案1:数据库去重表
@Autowired
private ProcessedMessageRepository processedMessageRepository;

@RabbitListener(queues = "order.queue")
public void handleOrder(String orderInfo, Message message, Channel channel) throws Exception {
String messageId = message.getMessageProperties().getMessageId();
long deliveryTag = message.getMessageProperties().getDeliveryTag();

try {
// ★ 幂等检查
if (processedMessageRepository.existsById(messageId)) {
System.out.println("Duplicate message, skip: " + messageId);
channel.basicAck(deliveryTag, false);
return;
}

// 处理业务
processOrder(orderInfo);

// 记录已处理
processedMessageRepository.save(new ProcessedMessage(messageId));

channel.basicAck(deliveryTag, false);
} catch (Exception e) {
channel.basicNack(deliveryTag, false, true);
}
}

// 方案2:Redis 去重
@Autowired
private StringRedisTemplate redisTemplate;

@RabbitListener(queues = "order.queue")
public void handleOrderRedis(String orderInfo, Message message, Channel channel) throws Exception {
String messageId = message.getMessageProperties().getMessageId();
long deliveryTag = message.getMessageProperties().getDeliveryTag();

// ★ Redis SETNX 幂等
Boolean isNew = redisTemplate.opsForValue()
.setIfAbsent("msg:" + messageId, "1", 24, TimeUnit.HOURS);

if (Boolean.FALSE.equals(isNew)) {
channel.basicAck(deliveryTag, false); // 重复消息直接 ACK
return;
}

try {
processOrder(orderInfo);
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
redisTemplate.delete("msg:" + messageId); // 失败删除,允许重试
channel.basicNack(deliveryTag, false, true);
}
}

22.2 消息积压

原因:
① 消费者处理速度 < 生产速度
② 消费者全部宕机
③ 消费者 ACK 失败导致反复重试

应对:
① 紧急扩容消费者
② 检查消费者代码性能
③ 考虑惰性队列(防止内存溢出)
④ 积压严重时:临时多开消费者快速消费

22.3 消息顺序性

问题:多消费者并发消费 → 消息顺序无法保证

场景:
Producer → [Queue] → C1(msg1) / C2(msg2) / C3(msg3)
C2 先完成 → msg2 先处理完 → 顺序错乱

解决:
① 同一业务 Key 的消息发到同一队列
② 同一队列只有一个消费者
③ 或者:消费者内部按 Key 分区顺序处理

22.4 最佳实践清单

维度最佳实践
连接管理 长连接复用,多 Channel 并发,开启自动重连
队列声明 durable=true,合理设置 x-arguments
消息属性 设置 messageId,deliveryMode=2(持久化)
生产者 开启 Confirm(异步),Mandatory+Return
消费者 手动 ACK,合理 prefetch,幂等处理
异常处理 NACK+requeue 重试,超限进死信队列
死信队列 始终配置 DLX,防止消息丢失
监控 监控 Ready/Unacked/Consumer 数
集群 至少 3 节点,使用 Quorum Queue
安全 修改默认密码,限制 vhost 权限

二十三、总结与速查表

23.1 Exchange 类型速查

类型路由规则bindingKey 支持典型场景
Fanout 忽略 Key 广播通知
Direct 精确匹配 字符串 按类型路由
Topic 通配符匹配 * # 日志分类
Headers 头匹配 x-match=all/any 多维路由

23.2 消息可靠性保障速查

环节机制配置
生产者→Exchange Publisher Confirm confirmSelect() + 回调
Exchange→Queue Mandatory + Return basicPublish(…, mandatory=true)
Queue 持久化 Durable Queue queueDeclare(…, durable=true, …)
消息持久化 Delivery Mode 2 MessageProperties.PERSISTENT
消费者端 手动 ACK basicConsume(…, autoAck=false)
处理失败 NACK/Reject + DLX basicNack + 死信队列绑定
Broker 宕机 镜像/Quorum Queue 集群 + Quorum Queue

23.3 六大模式速查

模式Exchange路由消费者
Simple 默认(“”) queue名 1个
Work 默认(“”) queue名 多个
Fanout Fanout 多个
Routing Direct 精确Key 多个
Topics Topic 通配符 多个
RPC 默认(“”) replyTo 1个

23.4 常用参数速查

参数位置说明
durable=true Queue/Exchange 持久化
exclusive=true Queue 连接级独占
autoDelete=true Queue/Exchange 无消费者后自动删除
x-message-ttl Queue 参数 队列消息TTL
x-expires Queue 参数 队列空闲超时
x-max-length Queue 参数 队列最大长度
x-dead-letter-exchange Queue 参数 死信交换机
x-dead-letter-routing-key Queue 参数 死信路由键
x-max-priority Queue 参数 最大优先级
x-queue-type Queue 参数 quorum / classic
x-delay 消息 header 延迟时间(需插件)
x-match Binding 参数 all / any(Headers)
basicQos(prefetch) Channel 预取数量
deliveryMode=2 消息属性 持久化消息

23.5 知识体系总览

RabbitMQ 知识体系

├── 基础
│ ├── AMQP 协议
│ ├── 核心概念(Exchange/Queue/Binding/RoutingKey)
│ ├── Connection / Channel
│ └── 虚拟主机 vhost

├── 安装部署
│ ├── Docker / Linux / Windows
│ ├── 端口说明
│ ├── 常用命令
│ └── 管理界面

├── Exchange 类型
│ ├── Fanout(广播)
│ ├── Direct(精确路由)
│ ├── Topic(通配符路由)
│ └── Headers(头匹配)

├── 六大工作模式
│ ├── Simple(简单模式)
│ ├── Work(工作队列)
│ ├── Fanout(发布订阅)
│ ├── Routing(路由模式)
│ ├── Topics(主题模式)
│ └── RPC(远程调用)

├── 消息可靠性
│ ├── 生产者:Confirm + Mandatory/Return
│ ├── 持久化:Queue Durable + DeliveryMode 2
│ ├── 消费者:手动 ACK + QoS
│ └── 重试:NACK/Reject → 死信队列

├── 高级特性
│ ├── 死信队列(DLX)
│ ├── 延迟队列(TTL+DLX / 插件)
│ ├── 消息优先级
│ ├── 惰性队列(Lazy Queue)
│ └── 仲裁队列(Quorum Queue)

├── 高可用
│ ├── 集群搭建
│ ├── 镜像队列(经典)
│ ├── 仲裁队列(推荐)
│ └── 节点类型(disc / ram)

├── 实战整合
│ ├── Java 原生客户端
│ ├── Spring Boot 整合
│ ├── Confirm/Return 回调
│ └── 消费者幂等设计

├── 运维监控
│ ├── 管理界面
│ ├── 命令行工具
│ ├── HTTP API
│ └── 告警指标

└── 最佳实践
├── 消息重复消费 → 幂等
├── 消息积压 → 扩容/惰性队列
├── 消息顺序 → 同Key同队列
└── 死信处理 → DLX + 告警


本文到此结束。 涵盖了 RabbitMQ 从安装部署到生产实战的完整知识体系,每个知识点都配有可运行的 Java 代码实例和 ASCII 架构图。

如果本文对你有帮助,欢迎点赞、收藏、关注,你的支持是我持续输出的动力!

赞(0)
未经允许不得转载:171主机测评 » 消息队列篇——RabbitMQ
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址