Spring Boot 整合 RabbitMQ 详细教程
1. 环境准备
1.1 Docker 安装 RabbitMQ(推荐)
# 拉取带管理界面的镜像
docker pull rabbitmq:4.3.2-management-alpine
# 运行容器
docker run -d \
--name rabbitmq \
-p 5672:5672 \
-p 15672:15672 \
-e RABBITMQ_DEFAULT_USER=admin \
-e RABBITMQ_DEFAULT_PASS=admin \
rabbitmq:3.12-managementdocker-compose.yml
services:
rabbitmq:
image: rabbitmq:4.3.2-management-alpine
container_name: rabbitmq
restart: always
mem_limit: 500m
memswap_limit: 500m
ports:
- "15672:15672"
- "5672:5672"
volumes:
- ./data:/var/lib/rabbitmq
- ./conf:/etc/rabbitmq
- ./log:/var/log/rabbitmq
environment:
- RABBITMQ_DEFAULT_USER=root
- RABBITMQ_DEFAULT_PASS=root
- TZ=Asia/Shanghai| 端口 | 用途 |
|---|---|
5672 | AMQP 协议端口(应用程序连接) |
15672 | 管理界面 HTTP 端口 |
访问管理界面:http://localhost:15672,用户名/密码:admin/admin
1.2 核心概念速览
| 概念 | 说明 |
|---|---|
| Producer | 消息生产者,发送消息到 Exchange |
| Consumer | 消息消费者,从 Queue 接收消息 |
| Exchange | 交换机,负责接收生产者消息并路由到队列 |
| Queue | 队列,存储消息的缓冲区 |
| Binding | 绑定,将 Exchange 和 Queue 关联的规则 |
| Routing Key | 路由键,Exchange 根据它决定消息发送到哪个队列 |
1.3 交换机类型
| 类型 | 路由规则 | 场景 |
|---|---|---|
direct | 精确匹配 Routing Key | 点对点消息 |
fanout | 广播到所有绑定队列 | 群发通知 |
topic | 模式匹配 Routing Key(支持 * 和 #) | 日志分类、新闻订阅 |
headers | 匹配消息 Headers | 复杂属性匹配(较少用) |
2. 项目搭建
2.1 创建 Spring Boot 项目
使用 Spring Initializr(https://start.spring.io)或 IDE 创建项目,添加以下依赖:
Maven (pom.xml):
<dependencies>
<!-- Spring Boot AMQP Starter -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
</dependencies>3. 基础配置
3.1 application.yml 配置
spring:
rabbitmq:
host: localhost
port: 5672
username: admin
password: admin
virtual-host: /
# 连接池配置(Spring Boot 3.x 默认使用 CachingConnectionFactory)
cache:
channel:
size: 25
checkout-timeout: 10s
# 生产者确认
# 开启发送确认
publisher-confirm-type: correlated
# 开启发送失败退回
publisher-returns: true
# 消费者配置
listener:
simple:
acknowledge-mode: manual # 手动确认
concurrency: 5 # 最小并发消费者数
max-concurrency: 10 # 最大并发消费者数
prefetch: 1 # 每次预取消息数
retry:
enabled: true
initial-interval: 1000ms
max-attempts: 3
max-interval: 10000ms
multiplier: 23.2 配置类:声明队列、交换机、绑定
package com.example.rabbitmq.config;
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 QUEUE_HELLO = "queue.hello";
public static final String QUEUE_WORK = "queue.work";
public static final String QUEUE_FANOUT_A = "queue.fanout.a";
public static final String QUEUE_FANOUT_B = "queue.fanout.b";
public static final String QUEUE_DIRECT_INFO = "queue.direct.info";
public static final String QUEUE_DIRECT_ERROR = "queue.direct.error";
public static final String QUEUE_TOPIC_NEWS = "queue.topic.news";
public static final String QUEUE_TOPIC_SPORT = "queue.topic.sport";
// 死信队列相关
public static final String QUEUE_NORMAL = "queue.normal";
public static final String QUEUE_DLX = "queue.dlx";
public static final String EXCHANGE_DLX = "exchange.dlx";
public static final String ROUTING_KEY_DLX = "routing.dlx";
// ========== 定义交换机 ==========
public static final String EXCHANGE_FANOUT = "exchange.fanout";
public static final String EXCHANGE_DIRECT = "exchange.direct";
public static final String EXCHANGE_TOPIC = "exchange.topic";
/**
* 声明 16 个分区队列及其绑定 Declarables
* 使用 Declarables 包装,RabbitAdmin 会自动识别并声明
*/
@Bean
public Declarables chatPartitionDeclarables() {
List<Declarable> declarables = new ArrayList<>();
Map<String, Object> args = new HashMap<>();
args.put("x-dead-letter-exchange", EXCHANGE_CHAT_DLX);
args.put("x-dead-letter-routing-key", ROUTING_KEY_DLX);
for (int i = 0; i < partitionCount; i++) {
String queueName = QUEUE_CHAT_MESSAGE_PREFIX + i;
Queue queue = new Queue(queueName, true, false, false, args);
Binding binding = BindingBuilder.bind(queue)
.to(chatExchange()) // 调用 @Bean 方法,CGLIB 代理会返回容器中的单例
.with(ROUTING_KEY_MESSAGE_PREFIX + i);
declarables.add(queue);
declarables.add(binding);
}
return new Declarables(declarables);
}
// ========== Hello World 队列 ==========
@Bean
public Queue helloQueue() {
// durable: 持久化队列,重启后仍在
return new Queue(QUEUE_HELLO, true); // durable = true
}
// ========== 工作队列 ==========
@Bean
public Queue workQueue() {
return QueueBuilder.durable(QUEUE_WORK).build();
}
// ========== Fanout 交换机及队列 ==========
@Bean
public FanoutExchange fanoutExchange() {
return new FanoutExchange(EXCHANGE_FANOUT);
}
@Bean
public Queue fanoutQueueA() {
return new Queue(QUEUE_FANOUT_A);
}
@Bean
public Queue fanoutQueueB() {
return new Queue(QUEUE_FANOUT_B);
}
@Bean
public Binding bindingFanoutA(FanoutExchange fanoutExchange, Queue fanoutQueueA) {
return BindingBuilder.bind(fanoutQueueA).to(fanoutExchange);
}
@Bean
public Binding bindingFanoutB(FanoutExchange fanoutExchange, Queue fanoutQueueB) {
return BindingBuilder.bind(fanoutQueueB).to(fanoutExchange);
}
// ========== Direct 交换机及队列 ==========
@Bean
public DirectExchange directExchange() {
return new DirectExchange(EXCHANGE_DIRECT);
}
@Bean
public Queue directQueueInfo() {
return new Queue(QUEUE_DIRECT_INFO);
}
@Bean
public Queue directQueueError() {
return new Queue(QUEUE_DIRECT_ERROR);
}
@Bean
public Binding bindingDirectInfo(DirectExchange directExchange, Queue directQueueInfo) {
return BindingBuilder.bind(directQueueInfo).to(directExchange).with("info");
}
@Bean
public Binding bindingDirectError(DirectExchange directExchange, Queue directQueueError) {
return BindingBuilder.bind(directQueueError).to(directExchange).with("error");
}
// ========== Topic 交换机及队列 ==========
@Bean
public TopicExchange topicExchange() {
return new TopicExchange(EXCHANGE_TOPIC);
}
@Bean
public Queue topicQueueNews() {
return new Queue(QUEUE_TOPIC_NEWS);
}
@Bean
public Queue topicQueueSport() {
return new Queue(QUEUE_TOPIC_SPORT);
}
// 匹配以 "news." 开头的路由键,如 news.china, news.world
@Bean
public Binding bindingTopicNews(TopicExchange topicExchange, Queue topicQueueNews) {
return BindingBuilder.bind(topicQueueNews).to(topicExchange).with("news.#");
}
// 匹配以 "sport." 开头的路由键
@Bean
public Binding bindingTopicSport(TopicExchange topicExchange, Queue topicQueueSport) {
return BindingBuilder.bind(topicQueueSport).to(topicExchange).with("sport.#");
}
// ========== 死信队列配置 ==========
@Bean
public DirectExchange dlxExchange() {
return new DirectExchange(EXCHANGE_DLX);
}
@Bean
public Queue dlxQueue() {
return new Queue(QUEUE_DLX);
}
@Bean
public Binding dlxBinding(DirectExchange dlxExchange, Queue dlxQueue) {
return BindingBuilder.bind(dlxQueue).to(dlxExchange).with(ROUTING_KEY_DLX);
}
// 正常队列,设置死信参数
@Bean
public Queue normalQueue() {
return QueueBuilder.durable(QUEUE_NORMAL)
.withArgument("x-dead-letter-exchange", EXCHANGE_DLX)
.withArgument("x-dead-letter-routing-key", ROUTING_KEY_DLX)
.withArgument("x-message-ttl", 10000) // 消息10秒过期
.build();
}
}
@RabbitListener(bindings = @QueueBinding(
value = @Queue(name = "store.addFavorite.success.queue", durable = "true"), // 队列 起名规则(服务名+业务名+成功+队列),durable持久化
exchange = @Exchange(name = "addFavorite.direct"), // 交换机名称,交换机默认类型就行direct,所以不用配置direct
key = "addFavorite.success" // 绑定的key4. 快速入门:Hello World
4.1 生产者
package com.example.rabbitmq.hello;
import com.example.rabbitmq.config.RabbitMQConfig;
import lombok.RequiredArgsConstructor;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Service;
@Service
@RequiredArgsConstructor
public class HelloProducer {
private final RabbitTemplate rabbitTemplate;
public void send(String message) {
rabbitTemplate.convertAndSend(RabbitMQConfig.QUEUE_HELLO, message);
System.out.println("[HelloProducer] 发送消息: " + message);
}
}4.2 消费者
package com.example.rabbitmq.hello;
import com.example.rabbitmq.config.RabbitMQConfig;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class HelloConsumer {
@RabbitListener(queues = RabbitMQConfig.QUEUE_HELLO)
public void receive(String message) {
log.info("[HelloConsumer] 收到消息: {}", message);
}
}4.3 测试接口
package com.example.rabbitmq.hello;
import lombok.RequiredArgsConstructor;
import org.springframework.web.bind.annotation.*;
@RestController
@RequestMapping("/hello")
@RequiredArgsConstructor
public class HelloController {
private final HelloProducer helloProducer;
@GetMapping("/send")
public String send(@RequestParam(defaultValue = "Hello RabbitMQ!") String msg) {
helloProducer.send(msg);
return "消息已发送: " + msg;
}
}5. 工作队列模式
多个消费者竞争消费同一个队列的消息,实现负载均衡。
5.1 生产者
package com.example.rabbitmq.work;
import com.example.rabbitmq.config.RabbitMQConfig;
import lombok.RequiredArgsConstructor;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Service;
@Service
@RequiredArgsConstructor
public class WorkProducer {
private final RabbitTemplate rabbitTemplate;
public void sendBatch(int count) {
for (int i = 0; i < count; i++) {
String message = "任务-" + i;
rabbitTemplate.convertAndSend(RabbitMQConfig.QUEUE_WORK, message);
}
System.out.println("[WorkProducer] 批量发送 " + count + " 条消息");
}
}5.2 两个消费者(竞争消费)
package com.example.rabbitmq.work;
import com.example.rabbitmq.config.RabbitMQConfig;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class WorkConsumer {
@RabbitListener(queues = RabbitMQConfig.QUEUE_WORK)
public void receiveA(String message) throws InterruptedException {
log.info("[WorkConsumer-A] 处理消息: {}", message);
Thread.sleep(100); // 模拟耗时处理
}
@RabbitListener(queues = RabbitMQConfig.QUEUE_WORK)
public void receiveB(String message) throws InterruptedException {
log.info("[WorkConsumer-B] 处理消息: {}", message);
Thread.sleep(200); // 模拟更耗时处理
}
}注意:默认轮询分发(Round-Robin),可通过
prefetch设置公平分发,让处理快的消费者消费更多。
6. 发布/订阅模式(Fanout)
消息广播到所有绑定队列,忽略 Routing Key。
6.1 生产者
package com.example.rabbitmq.fanout;
import com.example.rabbitmq.config.RabbitMQConfig;
import lombok.RequiredArgsConstructor;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Service;
@Service
@RequiredArgsConstructor
public class FanoutProducer {
private final RabbitTemplate rabbitTemplate;
public void broadcast(String message) {
// 第二个参数 routingKey 被忽略,设为 "" 即可
rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE_FANOUT, "", message);
System.out.println("[FanoutProducer] 广播消息: " + message);
}
}6.2 消费者
package com.example.rabbitmq.fanout;
import com.example.rabbitmq.config.RabbitMQConfig;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class FanoutConsumer {
@RabbitListener(queues = RabbitMQConfig.QUEUE_FANOUT_A)
public void receiveA(String message) {
log.info("[FanoutConsumer-A] 收到广播: {}", message);
}
@RabbitListener(queues = RabbitMQConfig.QUEUE_FANOUT_B)
public void receiveB(String message) {
log.info("[FanoutConsumer-B] 收到广播: {}", message);
}
}7. 路由模式(Direct)
精确匹配 Routing Key 进行消息路由。
7.1 生产者
package com.example.rabbitmq.direct;
import com.example.rabbitmq.config.RabbitMQConfig;
import lombok.RequiredArgsConstructor;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Service;
@Service
@RequiredArgsConstructor
public class DirectProducer {
private final RabbitTemplate rabbitTemplate;
public void sendInfo(String message) {
rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE_DIRECT, "info", "[INFO] " + message);
}
public void sendError(String message) {
rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE_DIRECT, "error", "[ERROR] " + message);
}
}7.2 消费者
package com.example.rabbitmq.direct;
import com.example.rabbitmq.config.RabbitMQConfig;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class DirectConsumer {
@RabbitListener(queues = RabbitMQConfig.QUEUE_DIRECT_INFO)
public void receiveInfo(String message) {
log.info("[DirectConsumer-Info] 收到: {}", message);
}
@RabbitListener(queues = RabbitMQConfig.QUEUE_DIRECT_ERROR)
public void receiveError(String message) {
log.info("[DirectConsumer-Error] 收到: {}", message);
}
}8. 主题模式(Topic)
支持通配符的模式匹配路由。
8.1 通配符规则
| 符号 | 含义 | 示例 |
|---|---|---|
* | 匹配一个单词 | sport.* 匹配 sport.football |
# | 匹配零个或多个单词 | news.# 匹配 news.china.beijing |
8.2 生产者
package com.example.rabbitmq.topic;
import com.example.rabbitmq.config.RabbitMQConfig;
import lombok.RequiredArgsConstructor;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Service;
@Service
@RequiredArgsConstructor
public class TopicProducer {
private final RabbitTemplate rabbitTemplate;
public void sendNews(String routingKey, String message) {
rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE_TOPIC, routingKey, message);
System.out.println("[TopicProducer] 发送 [" + routingKey + "]: " + message);
}
}8.3 测试用例
// 会被 news.# 匹配,进入 queue.topic.news
producer.sendNews("news.china", "中国新闻");
producer.sendNews("news.world.europe", "欧洲新闻");
// 会被 sport.# 匹配,进入 queue.topic.sport
producer.sendNews("sport.football", "足球新闻");
producer.sendNews("sport.basketball.nba", "NBA新闻");
// 不匹配任何规则,消息被丢弃(或进入死信,如果配置了)
producer.sendNews("weather.beijing", "天气预报");9. RPC 模式
使用 RabbitMQ 实现同步远程调用。
9.1 服务端
package com.example.rabbitmq.rpc;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.messaging.handler.annotation.SendTo;
import org.springframework.stereotype.Component;
@Component
public class RpcServer {
@RabbitListener(queues = "rpc_queue")
@SendTo("rpc_reply_queue") // 回复队列
public String handle(String request) {
System.out.println("[RpcServer] 收到请求: " + request);
// 模拟处理
return "处理结果: " + request.toUpperCase();
}
}9.2 客户端
package com.example.rabbitmq.rpc;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Service;
@Service
public class RpcClient {
private final RabbitTemplate rabbitTemplate;
public RpcClient(RabbitTemplate rabbitTemplate) {
this.rabbitTemplate = rabbitTemplate;
}
public String call(String message) {
// 使用 convertSendAndReceive 实现同步调用
Object response = rabbitTemplate.convertSendAndReceive("rpc_queue", message);
return (String) response;
}
}10. 消息确认机制
10.1 生产者确认(Publisher Confirm)
确保消息成功到达交换机/队列。
package com.example.rabbitmq.confirm;
import jakarta.annotation.PostConstruct;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.ReturnedMessage;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
@Slf4j
@Component
@RequiredArgsConstructor
public class ConfirmCallbackConfig {
private final RabbitTemplate rabbitTemplate;
@PostConstruct
public void init() {
// 消息到达 Exchange 的回调
rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
if (ack) {
log.info("[ConfirmCallback] 消息成功到达 Exchange");
} else {
log.error("[ConfirmCallback] 消息到达 Exchange 失败: {}", cause);
// 可在此进行补偿处理,如记录日志、存入数据库重试
}
});
// 消息无法路由到 Queue 的回调(如没有匹配的 Binding)
rabbitTemplate.setReturnsCallback((ReturnedMessage returned) -> {
log.error("[ReturnsCallback] 消息路由失败: exchange={}, routingKey={}, replyCode={}, replyText={}",
returned.getExchange(), returned.getRoutingKey(),
returned.getReplyCode(), returned.getReplyText());
});
}
}发送消息时带上 CorrelationData:
rabbitTemplate.convertAndSend(
RabbitMQConfig.EXCHANGE_NAME,
RabbitMQConfig.ROUTING_KEY,
message,
new CorrelationData(UUID.randomUUID().toString())
);10.2 消费者确认(Consumer Ack)
手动确认确保消息被正确处理后才会从队列删除。
package com.example.rabbitmq.ack;
import com.rabbitmq.client.Channel;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
import java.io.IOException;
@Slf4j
@Component
public class AckConsumer {
@RabbitListener(queues = "queue.ack")
public void receive(Message message, Channel channel) throws IOException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
String body = new String(message.getBody());
log.info("[AckConsumer] 处理消息: {}", body);
// 模拟业务处理
processBusiness(body);
// 手动确认(multiple=false 只确认当前消息)
channel.basicAck(deliveryTag, false);
log.info("[AckConsumer] 消息确认成功");
} catch (Exception e) {
log.error("[AckConsumer] 处理失败: {}", e.getMessage());
// 拒绝消息,requeue=false 进入死信队列(如果配置了 DLX)
// requeue=true 则重新入队(可能导致无限循环,慎用)
channel.basicNack(deliveryTag, false, false);
}
}
private void processBusiness(String body) {
// 业务逻辑
}
}11. 死信队列(DLX)
死信来源:
- 消息被拒绝(
basic.reject/basic.nack)且requeue=false - 消息 TTL 过期
- 队列达到最大长度
11.1 配置(已在上方 RabbitMQConfig 中声明)
// 正常队列绑定死信参数
@Bean
public Queue normalQueue() {
return QueueBuilder.durable(QUEUE_NORMAL)
.withArgument("x-dead-letter-exchange", EXCHANGE_DLX)
.withArgument("x-dead-letter-routing-key", ROUTING_KEY_DLX)
.withArgument("x-message-ttl", 10000) // 消息10秒过期
.build();
}11.2 死信消费者
package com.example.rabbitmq.dlx;
import com.example.rabbitmq.config.RabbitMQConfig;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class DlxConsumer {
@RabbitListener(queues = RabbitMQConfig.QUEUE_DLX)
public void receive(String message) {
log.info("[DlxConsumer] 收到死信消息: {}", message);
// 可进行告警、持久化到数据库、人工介入等处理
}
}12. 延迟队列
12.1 方案一:TTL + 死信队列(推荐,无需插件)
原理:消息在正常队列中等待 TTL 过期,自动进入死信队列,死信队列的消费者即延迟处理。
@Bean
public Queue delayQueue() {
return QueueBuilder.durable("queue.delay")
.withArgument("x-dead-letter-exchange", EXCHANGE_DLX)
.withArgument("x-dead-letter-routing-key", "routing.delay")
.withArgument("x-message-ttl", 30000) // 延迟30秒
.build();
}12.2 方案二:延迟消息插件(rabbitmq_delayed_message_exchange)
安装插件后使用自定义交换机:
@Bean
public CustomExchange delayedExchange() {
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct");
return new CustomExchange("exchange.delayed", "x-delayed-message", true, false, args);
}
// 发送时设置延迟时间
rabbitTemplate.convertAndSend("exchange.delayed", "routing.key", message, msg -> {
msg.getMessageProperties().setDelay(10000); // 延迟10秒
return msg;
});13. 消息幂等性
防止消息重复消费导致业务异常。
13.1 基于 Redis 实现幂等
package com.example.rabbitmq.idempotent;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import java.util.concurrent.TimeUnit;
@Slf4j
@Component
@RequiredArgsConstructor
public class IdempotentConsumer {
private final StringRedisTemplate redisTemplate;
@RabbitListener(queues = "queue.order")
public void handleOrder(String message) {
// 消息唯一标识(如订单号 + 消息ID)
String messageId = extractMessageId(message);
String key = "mq:idempotent:" + messageId;
// 使用 Redis SETNX 保证幂等
Boolean success = redisTemplate.opsForValue()
.setIfAbsent(key, "1", 24, TimeUnit.HOURS);
if (Boolean.FALSE.equals(success)) {
log.warn("[IdempotentConsumer] 消息已处理,跳过: {}", messageId);
return;
}
try {
// 执行业务逻辑
processOrder(message);
log.info("[IdempotentConsumer] 处理成功: {}", messageId);
} catch (Exception e) {
// 处理失败,删除 Redis 标记,允许重试
redisTemplate.delete(key);
throw new RuntimeException(e);
}
}
private String extractMessageId(String message) {
// 从消息中提取唯一ID
return message; // 简化示例
}
private void processOrder(String message) {
// 订单处理逻辑
}
}14. 监控与管理
14.1 RabbitMQ Management 插件
默认端口 15672,可查看:
- 队列消息堆积情况
- 消费者连接状态
- 交换机绑定关系
- 节点资源使用
14.2 Spring Boot Actuator 集成
添加依赖:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>访问 http://localhost:8080/actuator/health 查看 RabbitMQ 健康状态。
14.3 自定义队列监控
package com.example.rabbitmq.monitor;
import lombok.RequiredArgsConstructor;
import org.springframework.amqp.rabbit.core.RabbitAdmin;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.Properties;
@RestController
@RequiredArgsConstructor
public class QueueMonitorController {
private final RabbitAdmin rabbitAdmin;
@GetMapping("/queue/info")
public String getQueueInfo() {
Properties props = rabbitAdmin.getQueueProperties("queue.hello");
if (props != null) {
return String.format("队列: %s, 消息数: %s, 消费者数: %s",
props.getProperty("QUEUE_NAME"),
props.getProperty("QUEUE_MESSAGE_COUNT"),
props.getProperty("QUEUE_CONSUMER_COUNT"));
}
return "队列不存在";
}
}15. 命名规范
15.1 命名风格与分隔符
- 统一小写:避免大小写混用带来的问题,直接用全小写。
- 分隔符:推荐用点号
.或横杠-,首选点号,因为容易拼接层次结构。 - 不能有空格:AMQP 协议不允许。
- 长度:不要过长,建议控制在 255 字符以内(RabbitMQ 内部限制)。
好的示例: order.pay.queue、order.pay.exchange、order.pay.success
不好的示例: orderPayQueue(驼峰,风格不统一)、ORDER:PAY(不推荐特殊符号)
15.2 命名结构:业务模块.功能.类型
一种通用的分层结构如下:
<业务模块>.<功能/场景>.<资源类型>- 业务模块:如
order、user、inventory、notification - 功能/场景:如
create、cancel、timeout、dlx - 资源类型:
- 队列:
queue - 交换机:
exchange(或根据类型加前缀,见下) - 路由键:不写类型,直接就是
业务.功能或业务.功能.动作
- 队列:
示例
| 资源 | 命名 | 说明 |
|---|---|---|
| 队列 | order.pay.queue | 订单支付队列 |
| 交换机 | order.pay.exchange | 订单支付交换机 |
| 路由键 | order.pay.success | 支付成功路由键 |
15.3 交换机命名建议
交换机负责消息路由,其名称应该反映它服务的业务范围,而不是描述具体队列。
- 直连交换机 (Direct):
order.event.exchange、user.sync.exchange - 主题交换机 (Topic):
log.collector.exchange(用于按路由键模式匹配) - 扇出交换机 (Fanout):
broadcast.notice.exchange - 延迟交换机:
order.delay.exchange - 死信交换机:建议统一用后缀
.dlx,如order.dlx.exchange或order.pay.dlx
建议:不要直接在交换机名称中加上类型(如 direct),因为类型已经由代码声明,名字里加类型反而在改变类型时产生误导。但如果你希望名字自解释,可以在最后加 direct 字样,如 order.pay.direct,但不强制。
15.4 队列命名建议
队列是最终存储消息的地方,命名应明确消息的用途或消费者角色。
- 业务队列:
order.pay.queue、sms.send.queue - 死信队列:加上
.dlq后缀,如order.pay.dlq - 延迟队列(若不用插件而用死信模拟):
order.pay.delay.queue - 临时队列:带
tmp或anonymous,但 Spring 中多用自动删除队列,命名不重要。
原则:一个队列只服务于一个业务处理逻辑,不要多个消费者共用同一个队列处理完全不同的事情。
15.5 路由键命名建议
路由键决定了消息从交换机到队列的分发,必须与交换机的类型和绑定规则匹配。
- 对于 Direct 交换机:路由键通常等于队列绑定时指定的值,如
order.pay.success。 - 对于 Topic 交换机:路由键支持通配符,命名可用点号分段,如
order.pay.wechat、order.pay.alipay,绑定用order.pay.*。 - 对于 Fanout 交换机:路由键无意义,通常留空字符串
""。
推荐的命名方式:路由键尽量与业务动作关联,如 order.created、order.cancelled,同时与发送方定义的消息类型对等。
15.6 使用常量类管理
把所有名称定义在常量类中,避免在代码中散落字符串:
public final class RabbitMQConstants {
// 订单支付相关
public static final String ORDER_PAY_EXCHANGE = "order.pay.exchange";
public static final String ORDER_PAY_QUEUE = "order.pay.queue";
public static final String ORDER_PAY_ROUTING_KEY = "order.pay.success";
// 死信相关
public static final String ORDER_DLX_EXCHANGE = "order.dlx.exchange";
public static final String ORDER_DLX_QUEUE = "order.dlx.queue";
public static final String ORDER_DLX_ROUTING_KEY = "order.dlx"; // 死信路由键可以统一
// 延迟相关
public static final String ORDER_DELAY_EXCHANGE = "order.delay.exchange";
public static final String ORDER_DELAY_QUEUE = "order.delay.queue";
public static final String ORDER_DELAY_KEY = "order.delay";
}15.7 一个典型的命名示例
| 组件 | 名称 | 说明 |
|---|---|---|
| 业务交换机 | order.pay.exchange | 订单支付事件交换机 |
| 业务队列 | order.pay.queue | 订单支付处理队列 |
| 业务路由键 | order.pay.success | 支付成功消息的路由键 |
| 死信交换机 | order.pay.dlx | 订单支付死信交换机 |
| 死信队列 | order.pay.dlq | 订单支付死信队列 |
| 延迟交换机 | order.pay.delay.exchange | 支付延迟交换机 |
| 延迟路由键 | order.pay.delay | 延迟消息路由键 |
16. 常见问题
Q1:消息丢失怎么解决?
三层防护:
- 生产者端:开启
publisher-confirm和publisher-returns,失败时重试或持久化到数据库 - MQ 端:队列和消息都设置为持久化(
durable=true,deliveryMode=2) - 消费者端:开启手动 ACK,业务处理成功后再确认
Q2:消息重复消费怎么办?
- 保证消费者幂等性(Redis 去重、数据库唯一索引、业务状态机判断)
Q3:消息积压怎么处理?
- 增加消费者实例(横向扩容)
- 设置队列 TTL 和死信队列,避免无限堆积
- 紧急情况下,临时将消息转移到新队列,后续慢慢处理
Q4:如何确保消费顺序?
- 单个队列 + 单个消费者(牺牲吞吐量)
- 使用一致性 Hash 将相同 Key 的消息路由到同一队列
Q5:Spring Boot 连接不上 RabbitMQ?
检查清单:
- 网络连通:
telnet localhost 5672 - 用户名密码正确
- 用户具有对应 Virtual Host 的权限
- 防火墙/安全组放行端口
附录:完整项目结构
spring-boot-rabbitmq-demo/
├── src/main/java/com/example/rabbitmq/
│ ├── SpringBootRabbitmqDemoApplication.java
│ ├── config/
│ │ └── RabbitMQConfig.java
│ ├── hello/
│ │ ├── HelloProducer.java
│ │ ├── HelloConsumer.java
│ │ └── HelloController.java
│ ├── work/
│ ├── fanout/
│ ├── direct/
│ ├── topic/
│ ├── ack/
│ ├── dlx/
│ └── idempotent/
├── src/main/resources/
│ └── application.yml
└── pom.xml总结:RabbitMQ 是一个功能丰富、可靠性高的消息中间件。Spring Boot 通过
spring-boot-starter-amqp提供了开箱即用的支持。在实际项目中,建议始终开启生产者确认和消费者手动 ACK,并结合死信队列处理异常消息,确保消息系统的稳定性和可靠性。