RabbitMQ 从入门到精通:Spring Boot 实战三部曲(一)—— 基础核心与快速上手
专题导读:本系列共三篇,从基础到高级,带你系统掌握 RabbitMQ 在 Spring Boot 项目中的实战应用。
- 第一篇:基础核心与快速上手(本文)
- 第二篇:进阶特性与可靠性保障
- 第三篇:高级应用与性能优化
📖 前言
在当今的分布式系统中,消息队列已成为不可或缺的基础设施。RabbitMQ 作为最流行的消息中间件之一,以其可靠性、灵活性和易用性著称。
本文将从 RabbitMQ 的基础概念出发,结合 Spring Boot 项目,带你快速上手 RabbitMQ 开发,掌握核心概念和实际应用。
学完本文你将掌握:
- ✅ RabbitMQ 的核心概念与应用场景
- ✅ 五种工作模式详解
- ✅ Spring Boot 集成 RabbitMQ 的完整配置
- ✅ 实际业务场景中的消息发送与接收
- ✅ 常见问题的解决方案
一、RabbitMQ 是什么?为什么需要它?
1.1 RabbitMQ 简介
RabbitMQ 是一个开源的消息代理和队列服务器,基于 AMQP(Advanced Message Queuing Protocol)协议实现。
核心特点:
- 🚀 可靠性:支持消息持久化、事务、确认机制
- 🔄 灵活性:支持多种消息路由模式
- 🛡️ 高可用:支持集群、镜像队列
- 🌐 多语言支持:提供多种语言的客户端
- 📊 管理界面:内置 Web 管理控制台
1.2 典型应用场景
1┌─────────────────────────────────────────────┐ 2│ RabbitMQ 应用场景 │ 3├──────────────┬──────────────────────────────┤ 4│ 异步处理 │ 注册发送邮件、订单处理 │ 5│ 应用解耦 │ 微服务间通信、系统拆分 │ 6│ 流量削峰 │ 秒杀活动、突发流量 │ 7│ 日志收集 │ 分布式日志聚合 │ 8│ 任务队列 │ 定时任务、批量处理 │ 9│ 事件驱动 │ 状态变更通知、数据同步 │ 10└──────────────┴──────────────────────────────┘ 11
1.3 与其他消息队列对比
| 特性 | RabbitMQ | Kafka | RocketMQ |
|---|---|---|---|
| 吞吐量 | 万级 | 十万级 | 十万级 |
| 时效性 | 微秒级 | 毫秒级 | 毫秒级 |
| 可用性 | 高 | 非常高 | 非常高 |
| 可靠性 | 高 | 中 | 高 |
| 功能丰富度 | ⭐⭐⭐⭐⭐ | ⭐⭐⭐ | ⭐⭐⭐⭐ |
| 适用场景 | 中小规模、复杂路由 | 大数据、日志 | 大规模、金融 |
二、环境搭建与快速开始
2.1 RabbitMQ 安装
Docker 安装(推荐)
1# 拉取镜像 2docker pull rabbitmq:3.12-management 3 4# 启动容器 5docker run -d --name rabbitmq \ 6 -p 5672:5672 \ 7 -p 15672:15672 \ 8 -e RABBITMQ_DEFAULT_USER=admin \ 9 -e RABBITMQ_DEFAULT_PASS=admin123 \ 10 rabbitmq:3.12-management 11 12# 查看日志 13docker logs -f rabbitmq 14
访问管理界面:
- URL: http://localhost:15672
- 用户名: admin
- 密码: admin123
Linux 安装
1# Ubuntu/Debian 2sudo apt-get install rabbitmq-server 3 4# CentOS/RHEL 5sudo yum install rabbitmq-server 6 7# 启动服务 8sudo systemctl start rabbitmq-server 9sudo systemctl enable rabbitmq-server 10 11# 启用管理插件 12sudo rabbitmq-plugins enable rabbitmq_management 13
2.2 核心概念
1┌──────────────────────────────────────────────┐ 2│ RabbitMQ 核心概念 │ 3├──────────┬───────────────────────────────────┤ 4│ Producer │ 消息生产者,发送消息 │ 5│ Consumer │ 消息消费者,接收并处理消息 │ 6│ Queue │ 消息队列,存储消息 │ 7│ Exchange │ 交换机,接收并路由消息 │ 8│ Binding │ 绑定关系,连接 Exchange 和 Queue │ 9│ Routing │ 路由键,决定消息路由规则 │ 10│ Key │ │ 11└──────────┴───────────────────────────────────┘ 12
消息流转过程:
1Producer → Exchange → (Binding + Routing Key) → Queue → Consumer 2
三、五种工作模式详解
3.1 简单队列模式(Simple Queue)
最简单的模式:一个生产者,一个消费者,一个队列。
1Producer → [Queue] → Consumer 2
Spring Boot 实现
1. Maven 依赖
1<dependencies> 2 <!-- Spring Boot Starter AMQP --> 3 <dependency> 4 <groupId>org.springframework.boot</groupId> 5 <artifactId>spring-boot-starter-amqp</artifactId> 6 </dependency> 7</dependencies> 8
2. application.yml 配置
1spring: 2 rabbitmq: 3 host: localhost 4 port: 5672 5 username: admin 6 password: admin123 7 virtual-host: / 8 # 连接池配置 9 connection-timeout: 15000 10 # 发布者确认 11 publisher-confirm-type: correlated 12 publisher-returns: true 13 # 消费者配置 14 listener: 15 simple: 16 acknowledge-mode: manual # 手动确认 17 concurrency: 5 # 最小消费者数量 18 max-concurrency: 10 # 最大消费者数量 19 prefetch: 1 # 每次预取消息数 20
3. 配置类
1@Configuration 2public class RabbitMQConfig { 3 4 /** 5 * 声明简单队列 6 */ 7 @Bean 8 public Queue simpleQueue() { 9 // durable: 是否持久化 10 // exclusive: 是否排他 11 // autoDelete: 是否自动删除 12 return new Queue("simple.queue", true, false, false); 13 } 14} 15
4. 生产者
1@Component 2@Slf4j 3public class SimpleProducer { 4 5 @Autowired 6 private RabbitTemplate rabbitTemplate; 7 8 /** 9 * 发送消息 10 */ 11 public void sendMessage(String message) { 12 rabbitTemplate.convertAndSend("simple.queue", message); 13 log.info("发送消息: {}", message); 14 } 15 16 /** 17 * 发送对象消息 18 */ 19 public void sendObject(Object obj) { 20 rabbitTemplate.convertAndSend("simple.queue", obj); 21 log.info("发送对象消息: {}", obj); 22 } 23} 24
5. 消费者
1@Component 2@Slf4j 3public class SimpleConsumer { 4 5 /** 6 * 监听简单队列 7 */ 8 @RabbitListener(queues = "simple.queue") 9 public void receiveMessage(String message) { 10 log.info("收到消息: {}", message); 11 // 处理业务逻辑 12 processMessage(message); 13 } 14 15 /** 16 * 监听对象消息 17 */ 18 @RabbitListener(queues = "simple.queue") 19 public void receiveObject(Message message, Channel channel) throws Exception { 20 try { 21 // 反序列化消息 22 String body = new String(message.getBody(), "UTF-8"); 23 log.info("收到对象消息: {}", body); 24 25 // 处理业务逻辑 26 processMessage(body); 27 28 // 手动确认 29 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); 30 31 } catch (Exception e) { 32 log.error("消息处理失败", e); 33 // 拒绝消息,重新入队 34 channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true); 35 } 36 } 37 38 private void processMessage(String message) { 39 // 业务处理逻辑 40 System.out.println("Processing: " + message); 41 } 42} 43
6. 测试代码
1@SpringBootTest 2@RunWith(SpringRunner.class) 3public class SimpleQueueTest { 4 5 @Autowired 6 private SimpleProducer simpleProducer; 7 8 @Test 9 public void testSendMessage() { 10 // 发送字符串消息 11 simpleProducer.sendMessage("Hello RabbitMQ!"); 12 13 // 发送对象消息 14 User user = new User(); 15 user.setId(1L); 16 user.setName("张三"); 17 user.setEmail("zhangsan@example.com"); 18 simpleProducer.sendObject(user); 19 } 20} 21
3.2 工作队列模式(Work Queue)
多个消费者竞争消费同一个队列的消息,实现负载均衡。
1Producer → [Queue] → Consumer1 2 → Consumer2 3 → Consumer3 4
实战案例:订单处理
1. 配置类
1@Configuration 2public class WorkQueueConfig { 3 4 /** 5 * 声明工作队列 6 */ 7 @Bean 8 public Queue workQueue() { 9 return new Queue("work.queue", true); 10 } 11} 12
2. 生产者
1@Component 2@Slf4j 3public class OrderProducer { 4 5 @Autowired 6 private RabbitTemplate rabbitTemplate; 7 8 /** 9 * 发送订单消息 10 */ 11 public void sendOrder(Order order) { 12 rabbitTemplate.convertAndSend("work.queue", order); 13 log.info("发送订单消息: orderId={}", order.getOrderId()); 14 } 15 16 /** 17 * 批量发送订单 18 */ 19 public void batchSendOrders(List<Order> orders) { 20 for (Order order : orders) { 21 rabbitTemplate.convertAndSend("work.queue", order); 22 } 23 log.info("批量发送订单,数量: {}", orders.size()); 24 } 25} 26
3. 多个消费者
1@Component 2@Slf4j 3public class OrderConsumer1 { 4 5 @RabbitListener(queues = "work.queue") 6 public void processOrder(Order order, Channel channel, Message message) throws Exception { 7 try { 8 log.info("消费者1处理订单: orderId={}", order.getOrderId()); 9 10 // 模拟处理耗时 11 Thread.sleep(1000); 12 13 // 处理订单逻辑 14 processOrderLogic(order); 15 16 // 手动确认 17 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); 18 19 } catch (Exception e) { 20 log.error("订单处理失败", e); 21 channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true); 22 } 23 } 24 25 private void processOrderLogic(Order order) { 26 // 订单处理逻辑 27 System.out.println("Processing order: " + order.getOrderId()); 28 } 29} 30 31@Component 32@Slf4j 33public class OrderConsumer2 { 34 35 @RabbitListener(queues = "work.queue") 36 public void processOrder(Order order, Channel channel, Message message) throws Exception { 37 try { 38 log.info("消费者2处理订单: orderId={}", order.getOrderId()); 39 40 Thread.sleep(1000); 41 42 processOrderLogic(order); 43 44 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); 45 46 } catch (Exception e) { 47 log.error("订单处理失败", e); 48 channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true); 49 } 50 } 51 52 private void processOrderLogic(Order order) { 53 System.out.println("Processing order: " + order.getOrderId()); 54 } 55} 56
4. 测试代码
1@SpringBootTest 2public class WorkQueueTest { 3 4 @Autowired 5 private OrderProducer orderProducer; 6 7 @Test 8 public void testWorkQueue() throws InterruptedException { 9 // 发送100个订单 10 for (int i = 1; i <= 100; i++) { 11 Order order = new Order(); 12 order.setOrderId((long) i); 13 order.setAmount(new BigDecimal(i * 100)); 14 order.setCreateTime(new Date()); 15 16 orderProducer.sendOrder(order); 17 } 18 19 // 等待处理完成 20 Thread.sleep(30000); 21 } 22} 23
结果:
- 两个消费者平均分配100个订单
- 每个消费者处理约50个订单
- 实现负载均衡
3.3 发布订阅模式(Publish/Subscribe)
一个生产者发送消息,多个消费者都能收到相同的消息。
1Producer → [Fanout Exchange] → Queue1 → Consumer1 2 → Queue2 → Consumer2 3 → Queue3 → Consumer3 4
实战案例:用户注册通知
1. 配置类
1@Configuration 2public class FanoutConfig { 3 4 /** 5 * 声明 Fanout 交换机 6 */ 7 @Bean 8 public FanoutExchange fanoutExchange() { 9 return new FanoutExchange("user.fanout.exchange"); 10 } 11 12 /** 13 * 声明邮件队列 14 */ 15 @Bean 16 public Queue emailQueue() { 17 return new Queue("email.queue", true); 18 } 19 20 /** 21 * 声明短信队列 22 */ 23 @Bean 24 public Queue smsQueue() { 25 return new Queue("sms.queue", true); 26 } 27 28 /** 29 * 声明站内信队列 30 */ 31 @Bean 32 public Queue notificationQueue() { 33 return new Queue("notification.queue", true); 34 } 35 36 /** 37 * 绑定邮件队列到交换机 38 */ 39 @Bean 40 public Binding emailBinding(Queue emailQueue, FanoutExchange fanoutExchange) { 41 return BindingBuilder.bind(emailQueue).to(fanoutExchange); 42 } 43 44 /** 45 * 绑定短信队列到交换机 46 */ 47 @Bean 48 public Binding smsBinding(Queue smsQueue, FanoutExchange fanoutExchange) { 49 return BindingBuilder.bind(smsQueue).to(fanoutExchange); 50 } 51 52 /** 53 * 绑定站内信队列到交换机 54 */ 55 @Bean 56 public Binding notificationBinding(Queue notificationQueue, FanoutExchange fanoutExchange) { 57 return BindingBuilder.bind(notificationQueue).to(fanoutExchange); 58 } 59} 60
2. 生产者
1@Component 2@Slf4j 3public class UserRegisterProducer { 4 5 @Autowired 6 private RabbitTemplate rabbitTemplate; 7 8 /** 9 * 发送用户注册消息 10 */ 11 public void sendUserRegister(User user) { 12 rabbitTemplate.convertAndSend("user.fanout.exchange", "", user); 13 log.info("发送用户注册消息: userId={}, username={}", 14 user.getId(), user.getUsername()); 15 } 16} 17
3. 消费者
1@Component 2@Slf4j 3public class EmailConsumer { 4 5 @RabbitListener(queues = "email.queue") 6 public void sendEmail(User user) { 7 log.info("发送邮件通知: userId={}, email={}", 8 user.getId(), user.getEmail()); 9 10 // 调用邮件服务 11 sendWelcomeEmail(user.getEmail(), user.getUsername()); 12 } 13 14 private void sendWelcomeEmail(String email, String username) { 15 // 发送邮件逻辑 16 System.out.println("Sending welcome email to: " + email); 17 } 18} 19 20@Component 21@Slf4j 22public class SmsConsumer { 23 24 @RabbitListener(queues = "sms.queue") 25 public void sendSms(User user) { 26 log.info("发送短信通知: userId={}, phone={}", 27 user.getId(), user.getPhone()); 28 29 // 调用短信服务 30 sendWelcomeSms(user.getPhone(), user.getUsername()); 31 } 32 33 private void sendWelcomeSms(String phone, String username) { 34 // 发送短信逻辑 35 System.out.println("Sending welcome SMS to: " + phone); 36 } 37} 38 39@Component 40@Slf4j 41public class NotificationConsumer { 42 43 @RabbitListener(queues = "notification.queue") 44 public void sendNotification(User user) { 45 log.info("发送站内信通知: userId={}", user.getId()); 46 47 // 创建站内信 48 createNotification(user.getId(), "欢迎注册"); 49 } 50 51 private void createNotification(Long userId, String content) { 52 // 创建站内信逻辑 53 System.out.println("Creating notification for user: " + userId); 54 } 55} 56
4. 测试代码
1@SpringBootTest 2public class FanoutTest { 3 4 @Autowired 5 private UserRegisterProducer producer; 6 7 @Test 8 public void testFanout() { 9 User user = new User(); 10 user.setId(1L); 11 user.setUsername("张三"); 12 user.setEmail("zhangsan@example.com"); 13 user.setPhone("13800138000"); 14 15 producer.sendUserRegister(user); 16 17 // 三个消费者都会收到消息 18 } 19} 20
3.4 路由模式(Routing)
根据 routing key 将消息路由到不同的队列。
1Producer → [Direct Exchange] → Routing Key="error" → Error Queue 2 → Routing Key="warning" → Warning Queue 3 → Routing Key="info" → Info Queue 4
实战案例:日志分级处理
1. 配置类
1@Configuration 2public class DirectConfig { 3 4 /** 5 * 声明 Direct 交换机 6 */ 7 @Bean 8 public DirectExchange directExchange() { 9 return new DirectExchange("log.direct.exchange"); 10 } 11 12 /** 13 * 声明错误日志队列 14 */ 15 @Bean 16 public Queue errorQueue() { 17 return new Queue("log.error.queue", true); 18 } 19 20 /** 21 * 声明警告日志队列 22 */ 23 @Bean 24 public Queue warningQueue() { 25 return new Queue("log.warning.queue", true); 26 } 27 28 /** 29 * 声明信息日志队列 30 */ 31 @Bean 32 public Queue infoQueue() { 33 return new Queue("log.info.queue", true); 34 } 35 36 /** 37 * 绑定错误日志队列 38 */ 39 @Bean 40 public Binding errorBinding(Queue errorQueue, DirectExchange directExchange) { 41 return BindingBuilder.bind(errorQueue) 42 .to(directExchange) 43 .with("error"); 44 } 45 46 /** 47 * 绑定警告日志队列 48 */ 49 @Bean 50 public Binding warningBinding(Queue warningQueue, DirectExchange directExchange) { 51 return BindingBuilder.bind(warningQueue) 52 .to(directExchange) 53 .with("warning"); 54 } 55 56 /** 57 * 绑定信息日志队列 58 */ 59 @Bean 60 public Binding infoBinding(Queue infoQueue, DirectExchange directExchange) { 61 return BindingBuilder.bind(infoQueue) 62 .to(directExchange) 63 .with("info"); 64 } 65} 66
2. 消息实体
1@Data 2@AllArgsConstructor 3@NoArgsConstructor 4public class LogMessage implements Serializable { 5 6 private static final long serialVersionUID = 1L; 7 8 private String level; // 日志级别 9 private String message; // 日志内容 10 private String serviceName;// 服务名称 11 private Long timestamp; // 时间戳 12} 13
3. 生产者
1@Component 2@Slf4j 3public class LogProducer { 4 5 @Autowired 6 private RabbitTemplate rabbitTemplate; 7 8 /** 9 * 发送错误日志 10 */ 11 public void sendErrorLog(String message, String serviceName) { 12 LogMessage logMessage = new LogMessage("error", message, serviceName, System.currentTimeMillis()); 13 rabbitTemplate.convertAndSend("log.direct.exchange", "error", logMessage); 14 log.info("发送错误日志: {}", message); 15 } 16 17 /** 18 * 发送警告日志 19 */ 20 public void sendWarningLog(String message, String serviceName) { 21 LogMessage logMessage = new LogMessage("warning", message, serviceName, System.currentTimeMillis()); 22 rabbitTemplate.convertAndSend("log.direct.exchange", "warning", logMessage); 23 log.info("发送警告日志: {}", message); 24 } 25 26 /** 27 * 发送信息日志 28 */ 29 public void sendInfoLog(String message, String serviceName) { 30 LogMessage logMessage = new LogMessage("info", message, serviceName, System.currentTimeMillis()); 31 rabbitTemplate.convertAndSend("log.direct.exchange", "info", logMessage); 32 log.info("发送信息日志: {}", message); 33 } 34} 35
4. 消费者
1@Component 2@Slf4j 3public class ErrorLogConsumer { 4 5 @RabbitListener(queues = "log.error.queue") 6 public void consumeErrorLog(LogMessage logMessage) { 7 log.error("【错误日志】服务: {}, 内容: {}", 8 logMessage.getServiceName(), logMessage.getMessage()); 9 10 // 保存错误日志到数据库 11 saveErrorLog(logMessage); 12 13 // 发送告警通知 14 sendAlert(logMessage); 15 } 16 17 private void saveErrorLog(LogMessage logMessage) { 18 // 保存逻辑 19 } 20 21 private void sendAlert(LogMessage logMessage) { 22 // 告警逻辑 23 } 24} 25 26@Component 27@Slf4j 28public class WarningLogConsumer { 29 30 @RabbitListener(queues = "log.warning.queue") 31 public void consumeWarningLog(LogMessage logMessage) { 32 log.warn("【警告日志】服务: {}, 内容: {}", 33 logMessage.getServiceName(), logMessage.getMessage()); 34 35 // 保存警告日志 36 saveWarningLog(logMessage); 37 } 38 39 private void saveWarningLog(LogMessage logMessage) { 40 // 保存逻辑 41 } 42} 43 44@Component 45@Slf4j 46public class InfoLogConsumer { 47 48 @RabbitListener(queues = "log.info.queue") 49 public void consumeInfoLog(LogMessage logMessage) { 50 log.info("【信息日志】服务: {}, 内容: {}", 51 logMessage.getServiceName(), logMessage.getMessage()); 52 53 // 保存信息日志 54 saveInfoLog(logMessage); 55 } 56 57 private void saveInfoLog(LogMessage logMessage) { 58 // 保存逻辑 59 } 60} 61
5. 测试代码
1@SpringBootTest 2public class DirectTest { 3 4 @Autowired 5 private LogProducer logProducer; 6 7 @Test 8 public void testDirect() { 9 // 发送不同级别的日志 10 logProducer.sendErrorLog("数据库连接失败", "order-service"); 11 logProducer.sendWarningLog("内存使用率超过80%", "user-service"); 12 logProducer.sendInfoLog("用户登录成功", "auth-service"); 13 14 // 每条日志会被路由到对应的队列 15 } 16} 17
3.5 主题模式(Topic)
根据 routing key 的模式匹配进行路由,更灵活的路由方式。
1Producer → [Topic Exchange] → *.error.* → Error Queue 2 → order.# → Order Queue 3 → user.*.create → UserCreate Queue 4
匹配规则:
*:匹配一个单词#:匹配零个或多个单词
实战案例:电商消息路由
1. 配置类
1@Configuration 2public class TopicConfig { 3 4 /** 5 * 声明 Topic 交换机 6 */ 7 @Bean 8 public TopicExchange topicExchange() { 9 return new TopicExchange("ecommerce.topic.exchange"); 10 } 11 12 /** 13 * 声明订单队列 14 */ 15 @Bean 16 public Queue orderQueue() { 17 return new Queue("order.queue", true); 18 } 19 20 /** 21 * 声明用户队列 22 */ 23 @Bean 24 public Queue userQueue() { 25 return new Queue("user.queue", true); 26 } 27 28 /** 29 * 声明所有消息队列 30 */ 31 @Bean 32 public Queue allQueue() { 33 return new Queue("all.queue", true); 34 } 35 36 /** 37 * 绑定订单队列(匹配 order.*) 38 */ 39 @Bean 40 public Binding orderBinding(Queue orderQueue, TopicExchange topicExchange) { 41 return BindingBuilder.bind(orderQueue) 42 .to(topicExchange) 43 .with("order.*"); 44 } 45 46 /** 47 * 绑定用户队列(匹配 user.#) 48 */ 49 @Bean 50 public Binding userBinding(Queue userQueue, TopicExchange topicExchange) { 51 return BindingBuilder.bind(userQueue) 52 .to(topicExchange) 53 .with("user.#"); 54 } 55 56 /** 57 * 绑定所有消息队列(匹配 #) 58 */ 59 @Bean 60 public Binding allBinding(Queue allQueue, TopicExchange topicExchange) { 61 return BindingBuilder.bind(allQueue) 62 .to(topicExchange) 63 .with("#"); 64 } 65} 66
2. 生产者
1@Component 2@Slf4j 3public class EcommerceProducer { 4 5 @Autowired 6 private RabbitTemplate rabbitTemplate; 7 8 /** 9 * 发送订单创建消息 10 */ 11 public void sendOrderCreate(Order order) { 12 rabbitTemplate.convertAndSend("ecommerce.topic.exchange", 13 "order.create", order); 14 log.info("发送订单创建消息: orderId={}", order.getOrderId()); 15 } 16 17 /** 18 * 发送订单支付消息 19 */ 20 public void sendOrderPay(Order order) { 21 rabbitTemplate.convertAndSend("ecommerce.topic.exchange", 22 "order.pay", order); 23 log.info("发送订单支付消息: orderId={}", order.getOrderId()); 24 } 25 26 /** 27 * 发送用户注册消息 28 */ 29 public void sendUserRegister(User user) { 30 rabbitTemplate.convertAndSend("ecommerce.topic.exchange", 31 "user.register", user); 32 log.info("发送用户注册消息: userId={}", user.getId()); 33 } 34 35 /** 36 * 发送用户修改消息 37 */ 38 public void sendUserUpdate(User user) { 39 rabbitTemplate.convertAndSend("ecommerce.topic.exchange", 40 "user.profile.update", user); 41 log.info("发送用户修改消息: userId={}", user.getId()); 42 } 43} 44
3. 消费者
1@Component 2@Slf4j 3public class OrderConsumer { 4 5 @RabbitListener(queues = "order.queue") 6 public void processOrder(Message message, String routingKey) { 7 log.info("订单消费者收到消息: routingKey={}", routingKey); 8 // 处理订单相关逻辑 9 } 10} 11 12@Component 13@Slf4j 14public class UserConsumer { 15 16 @RabbitListener(queues = "user.queue") 17 public void processUser(Message message, String routingKey) { 18 log.info("用户消费者收到消息: routingKey={}", routingKey); 19 // 处理用户相关逻辑 20 } 21} 22 23@Component 24@Slf4j 25public class AllConsumer { 26 27 @RabbitListener(queues = "all.queue") 28 public void processAll(Message message, String routingKey) { 29 log.info("全局消费者收到消息: routingKey={}", routingKey); 30 // 记录所有消息 31 } 32} 33
4. 测试代码
1@SpringBootTest 2public class TopicTest { 3 4 @Autowired 5 private EcommerceProducer producer; 6 7 @Test 8 public void testTopic() { 9 // 订单消息 10 Order order1 = new Order(); 11 order1.setOrderId(1L); 12 producer.sendOrderCreate(order1); 13 14 Order order2 = new Order(); 15 order2.setOrderId(2L); 16 producer.sendOrderPay(order2); 17 18 // 用户消息 19 User user1 = new User(); 20 user1.setId(1L); 21 producer.sendUserRegister(user1); 22 23 User user2 = new User(); 24 user2.setId(2L); 25 producer.sendUserUpdate(user2); 26 27 /* 28 * 路由结果: 29 * order.create → order.queue + all.queue 30 * order.pay → order.queue + all.queue 31 * user.register → user.queue + all.queue 32 * user.profile.update → user.queue + all.queue 33 */ 34 } 35} 36
四、Spring Boot 集成最佳实践
4.1 完整配置类
1@Configuration 2@Slf4j 3public class RabbitMQCompleteConfig { 4 5 @Value("${spring.rabbitmq.host}") 6 private String host; 7 8 @Value("${spring.rabbitmq.port}") 9 private int port; 10 11 @Value("${spring.rabbitmq.username}") 12 private String username; 13 14 @Value("${spring.rabbitmq.password}") 15 private String password; 16 17 /** 18 * 连接工厂 19 */ 20 @Bean 21 public ConnectionFactory connectionFactory() { 22 CachingConnectionFactory factory = new CachingConnectionFactory(); 23 factory.setHost(host); 24 factory.setPort(port); 25 factory.setUsername(username); 26 factory.setPassword(password); 27 factory.setVirtualHost("/"); 28 29 // 开启发布者确认 30 factory.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.CORRELATED); 31 factory.setPublisherReturns(true); 32 33 return factory; 34 } 35 36 /** 37 * RabbitTemplate 38 */ 39 @Bean 40 public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) { 41 RabbitTemplate template = new RabbitTemplate(connectionFactory); 42 43 // 消息序列化 44 template.setMessageConverter(new Jackson2JsonMessageConverter()); 45 46 // 强制返回 47 template.setMandatory(true); 48 49 // 确认回调 50 template.setConfirmCallback((correlationData, ack, cause) -> { 51 if (ack) { 52 log.info("消息发送成功: {}", correlationData); 53 } else { 54 log.error("消息发送失败: {}, 原因: {}", correlationData, cause); 55 } 56 }); 57 58 // 返回回调 59 template.setReturnsCallback(returned -> { 60 log.error("消息返回: {}, 路由键: {}, 交换机: {}", 61 new String(returned.getMessage().getBody()), 62 returned.getRoutingKey(), 63 returned.getExchange()); 64 }); 65 66 return template; 67 } 68 69 /** 70 * 消息转换器 71 */ 72 @Bean 73 public MessageConverter messageConverter() { 74 return new Jackson2JsonMessageConverter(); 75 } 76} 77
4.2 通用工具类
1@Component 2@Slf4j 3public class RabbitMQUtil { 4 5 @Autowired 6 private RabbitTemplate rabbitTemplate; 7 8 /** 9 * 发送消息到指定队列 10 */ 11 public void sendToQueue(String queueName, Object message) { 12 rabbitTemplate.convertAndSend(queueName, message); 13 log.info("发送消息到队列: {}, 消息: {}", queueName, message); 14 } 15 16 /** 17 * 发送消息到交换机 18 */ 19 public void sendToExchange(String exchange, String routingKey, Object message) { 20 rabbitTemplate.convertAndSend(exchange, routingKey, message); 21 log.info("发送消息到交换机: {}, 路由键: {}", exchange, routingKey); 22 } 23 24 /** 25 * 延迟发送消息 26 */ 27 public void sendWithDelay(String queueName, Object message, long delayMillis) { 28 MessagePostProcessor processor = msg -> { 29 msg.getMessageProperties().setDelay((int) delayMillis); 30 return msg; 31 }; 32 rabbitTemplate.convertAndSend(queueName, message, processor); 33 log.info("延迟发送消息: {}ms", delayMillis); 34 } 35} 36
4.3 消息实体规范
1@Data 2@AllArgsConstructor 3@NoArgsConstructor 4@Builder 5public class BaseMessage implements Serializable { 6 7 private static final long serialVersionUID = 1L; 8 9 /** 10 * 消息ID 11 */ 12 private String messageId; 13 14 /** 15 * 消息类型 16 */ 17 private String messageType; 18 19 /** 20 * 消息内容 21 */ 22 private Object data; 23 24 /** 25 * 创建时间 26 */ 27 private Long timestamp; 28 29 /** 30 * 版本号 31 */ 32 private String version; 33 34 @PrePersist 35 public void prePersist() { 36 if (this.messageId == null) { 37 this.messageId = UUID.randomUUID().toString(); 38 } 39 if (this.timestamp == null) { 40 this.timestamp = System.currentTimeMillis(); 41 } 42 if (this.version == null) { 43 this.version = "1.0"; 44 } 45 } 46} 47
五、实战项目:电商订单系统
5.1 项目结构
1ecommerce-order/ 2├── config/ 3│ └── RabbitMQConfig.java 4├── producer/ 5│ ├── OrderProducer.java 6│ └── PaymentProducer.java 7├── consumer/ 8│ ├── OrderConsumer.java 9│ ├── PaymentConsumer.java 10│ └── InventoryConsumer.java 11├── service/ 12│ ├── OrderService.java 13│ └── PaymentService.java 14└── model/ 15 ├── Order.java 16 └── Payment.java 17
5.2 核心代码实现
订单创建流程:
1@Service 2@Slf4j 3public class OrderService { 4 5 @Autowired 6 private OrderProducer orderProducer; 7 8 @Autowired 9 private OrderMapper orderMapper; 10 11 /** 12 * 创建订单 13 */ 14 @Transactional 15 public Long createOrder(OrderRequest request) { 16 // 1. 创建订单 17 Order order = new Order(); 18 order.setUserId(request.getUserId()); 19 order.setProductId(request.getProductId()); 20 order.setQuantity(request.getQuantity()); 21 order.setAmount(request.getAmount()); 22 order.setStatus(OrderStatus.CREATED); 23 order.setCreateTime(new Date()); 24 25 orderMapper.insert(order); 26 27 // 2. 发送订单创建消息 28 orderProducer.sendOrderCreated(order); 29 30 log.info("订单创建成功: orderId={}", order.getId()); 31 return order.getId(); 32 } 33} 34
订单消费者:
1@Component 2@Slf4j 3public class OrderConsumer { 4 5 @Autowired 6 private InventoryService inventoryService; 7 8 @Autowired 9 private NotificationService notificationService; 10 11 @RabbitListener(queues = "order.created.queue") 12 public void handleOrderCreated(Order order, Channel channel, Message message) throws Exception { 13 try { 14 log.info("处理订单创建: orderId={}", order.getId()); 15 16 // 1. 扣减库存 17 boolean success = inventoryService.deductStock( 18 order.getProductId(), 19 order.getQuantity() 20 ); 21 22 if (!success) { 23 throw new RuntimeException("库存不足"); 24 } 25 26 // 2. 更新订单状态 27 updateOrderStatus(order.getId(), OrderStatus.PAID); 28 29 // 3. 发送通知 30 notificationService.sendOrderNotification(order); 31 32 // 手动确认 33 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); 34 35 } catch (Exception e) { 36 log.error("订单处理失败", e); 37 // 拒绝消息,不重新入队(进入死信队列) 38 channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, false); 39 } 40 } 41 42 private void updateOrderStatus(Long orderId, OrderStatus status) { 43 // 更新订单状态 44 } 45} 46
六、常见问题与解决方案
6.1 消息丢失问题
解决方案:开启持久化和确认机制
1spring: 2 rabbitmq: 3 # 发布者确认 4 publisher-confirm-type: correlated 5 publisher-returns: true 6 # 消费者手动确认 7 listener: 8 simple: 9 acknowledge-mode: manual 10
1// 队列持久化 2@Bean 3public Queue queue() { 4 return new Queue("queue.name", true); // durable=true 5} 6 7// 交换机持久化 8@Bean 9public DirectExchange exchange() { 10 return new DirectExchange("exchange.name", true, false); 11} 12 13// 消息持久化 14rabbitTemplate.convertAndSend("exchange", "routingKey", message, 15 msg -> { 16 msg.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT); 17 return msg; 18 }); 19
6.2 消息重复消费
解决方案:消息幂等性处理
1@Component 2public class IdempotentConsumer { 3 4 @Autowired 5 private RedisTemplate<String, Object> redisTemplate; 6 7 @RabbitListener(queues = "order.queue") 8 public void consume(Order order, Channel channel, Message message) throws Exception { 9 String messageId = order.getMessageId(); 10 String key = "message:processed:" + messageId; 11 12 try { 13 // 检查是否已处理 14 Boolean processed = redisTemplate.opsForValue() 15 .setIfAbsent(key, "1", 24, TimeUnit.HOURS); 16 17 if (Boolean.FALSE.equals(processed)) { 18 log.warn("消息已处理,跳过: {}", messageId); 19 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); 20 return; 21 } 22 23 // 处理业务逻辑 24 processOrder(order); 25 26 // 确认消息 27 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); 28 29 } catch (Exception e) { 30 log.error("消息处理失败", e); 31 channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true); 32 } 33 } 34} 35
6.3 消息积压问题
解决方案:增加消费者数量
1spring: 2 rabbitmq: 3 listener: 4 simple: 5 concurrency: 10 # 增加并发消费者数量 6 max-concurrency: 50 7 prefetch: 10 # 增加预取数量 8
七、总结与展望
7.1 本文要点回顾
✅ RabbitMQ 基础概念:理解核心组件和工作原理
✅ 五种工作模式:掌握 Simple、Work、Fanout、Direct、Topic
✅ Spring Boot 集成:完成配置类和工具类封装
✅ 实战案例:电商订单系统完整实现
✅ 常见问题:消息丢失、重复消费、积压的解决方案
7.2 下篇预告
在下一篇文章《RabbitMQ 从入门到精通:Spring Boot 实战三部曲(二)—— 进阶特性与可靠性保障》中,我们将深入探讨:
- 🔒 消息可靠性保障:确认机制、持久化、事务
- 💀 死信队列:处理失败消息的完整方案
- ⏰ 延迟队列:实现定时任务的多种方式
- 🔄 消息幂等性:保证消息不被重复消费
- 📊 监控与管理:RabbitMQ Management 深度使用
- 🛡️ 高可用架构:镜像队列与集群
7.3 学习建议
- 动手实践:亲自搭建环境,运行示例代码
- 理解原理:不仅要会用,更要理解为什么
- 关注可靠性:生产环境必须考虑消息的可靠传递
- 监控先行:及时发现问题,避免消息积压
📚 参考资料
- RabbitMQ 官方文档:https://www.rabbitmq.com/documentation.html
- Spring AMQP 文档:https://spring.io/projects/spring-amqp
- 《RabbitMQ 实战指南》
觉得有用?欢迎点赞、收藏、转发!
下一篇更精彩,敬请期待! 🚀
系列文章:
- [第一篇] RabbitMQ 从入门到精通:Spring Boot 实战三部曲(一)—— 基础核心与快速上手
- [第二篇] RabbitMQ 从入门到精通:Spring Boot 实战三部曲(二)—— 进阶特性与可靠性保障
- [第三篇] RabbitMQ 从入门到精通:Spring Boot 实战三部曲(三)—— 高级应用与性能优化
《RabbitMQ 从入门到精通:Spring Boot 实战三部曲(一)—— 基础核心与快速上手》 是转载文章,点击查看原文。