MQ 从 0 到 1(三):Spring Boot 集成 RabbitMQ
前两篇讲了 MQ 的作用和 RabbitMQ 的基础模型。
这一篇开始写 Spring Boot 里的使用方式。
目标不是把所有配置背下来,而是先搭出一条完整链路:
Controller -> Service -> RabbitTemplate 发送消息 -> RabbitMQ -> @RabbitListener 消费消息 -> 执行业务逻辑引入依赖
Spring Boot 项目里通常使用 Spring AMQP。
Maven 依赖:
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId></dependency>版本一般交给 Spring Boot 管理,不需要自己手动写。
基础配置
application.yml 可以这样写:
spring: rabbitmq: host: localhost port: 5672 username: guest password: guest virtual-host: /本地学习环境可以先用默认账号。
真实项目里要单独创建用户、vhost 和权限,不要直接用默认账号。
定义消息对象
比如订单创建消息:
public class OrderCreatedMessage {
private String messageId;
private Long orderId;
private Long userId;
private String occurredAt;
public String getMessageId() { return messageId; }
public void setMessageId(String messageId) { this.messageId = messageId; }
public Long getOrderId() { return orderId; }
public void setOrderId(Long orderId) { this.orderId = orderId; }
public Long getUserId() { return userId; }
public void setUserId(Long userId) { this.userId = userId; }
public String getOccurredAt() { return occurredAt; }
public void setOccurredAt(String occurredAt) { this.occurredAt = occurredAt; }}消息对象建议包含 messageId。
后面做幂等、排查日志、追踪链路都会用到。
定义交换机、队列和绑定
可以用配置类声明 RabbitMQ 资源。
@Configurationpublic class RabbitOrderConfig {
public static final String ORDER_EXCHANGE = "order.topic.exchange";
public static final String ORDER_CREATED_ROUTING_KEY = "order.created";
public static final String ORDER_SMS_QUEUE = "order.sms.queue";
@Bean public TopicExchange orderExchange() { return ExchangeBuilder .topicExchange(ORDER_EXCHANGE) .durable(true) .build(); }
@Bean public Queue orderSmsQueue() { return QueueBuilder .durable(ORDER_SMS_QUEUE) .build(); }
@Bean public Binding orderSmsBinding() { return BindingBuilder .bind(orderSmsQueue()) .to(orderExchange()) .with(ORDER_CREATED_ROUTING_KEY); }}这段代码声明了:
交换机:order.topic.exchange队列:order.sms.queue绑定规则:order.created当生产者发送消息到 order.topic.exchange,并且 routing key 是 order.created 时,消息就会进入 order.sms.queue。
配置 JSON 转换
默认情况下,消息体可能会按 Java 序列化方式处理。
业务项目里更推荐使用 JSON。
@Configurationpublic class RabbitMessageConfig {
@Bean public MessageConverter messageConverter() { return new Jackson2JsonMessageConverter(); }}这样发送对象时,会转换成 JSON。
消费者接收时,也可以自动转换成对应对象。
发送消息
发送消息使用 RabbitTemplate。
@Servicepublic class OrderMessageProducer {
private final RabbitTemplate rabbitTemplate;
public OrderMessageProducer(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; }
public void sendOrderCreatedMessage(OrderCreatedMessage message) { rabbitTemplate.convertAndSend( RabbitOrderConfig.ORDER_EXCHANGE, RabbitOrderConfig.ORDER_CREATED_ROUTING_KEY, message ); }}业务代码里可以这样调用:
@Servicepublic class OrderService {
private final OrderMessageProducer orderMessageProducer;
public OrderService(OrderMessageProducer orderMessageProducer) { this.orderMessageProducer = orderMessageProducer; }
public void createOrder(CreateOrderCommand command) { Long orderId = saveOrder(command);
OrderCreatedMessage message = new OrderCreatedMessage(); message.setMessageId(UUID.randomUUID().toString()); message.setOrderId(orderId); message.setUserId(command.getUserId()); message.setOccurredAt(LocalDateTime.now().toString());
orderMessageProducer.sendOrderCreatedMessage(message); }
private Long saveOrder(CreateOrderCommand command) { return 10001L; }}这里先写简单版本。
后面讲可靠性时会补充一个关键问题:
数据库事务提交成功,但消息发送失败怎么办?
真实项目不能只写 saveOrder() 后直接 sendMessage() 就结束。
消费消息
消费者可以用 @RabbitListener。
@Componentpublic class OrderSmsConsumer {
private final SmsService smsService;
public OrderSmsConsumer(SmsService smsService) { this.smsService = smsService; }
@RabbitListener(queues = RabbitOrderConfig.ORDER_SMS_QUEUE) public void handle(OrderCreatedMessage message) { smsService.sendOrderCreatedMessage(message.getOrderId()); }}当 order.sms.queue 里有消息时,这个方法就会被触发。
消费者里不要写太多东西
消费者方法最好保持薄一点。
不推荐这样:
@RabbitListener(queues = RabbitOrderConfig.ORDER_SMS_QUEUE)public void handle(OrderCreatedMessage message) { // 校验参数 // 查订单 // 查用户 // 拼短信模板 // 调短信接口 // 写发送记录 // 处理异常 // 更新状态}更推荐:
@RabbitListener(queues = RabbitOrderConfig.ORDER_SMS_QUEUE)public void handle(OrderCreatedMessage message) { orderSmsApplicationService.sendOrderCreatedSms(message);}消费者只负责接消息和转发给应用服务。
具体业务逻辑放在服务层里,方便测试,也方便以后被定时补偿任务复用。
手动确认的思路
入门阶段可以先用默认确认机制。
但真实项目里,经常需要手动 ack。
大致流程是:
收到消息执行业务业务成功 -> ack业务失败 -> nack 或 reject示意代码:
@RabbitListener(queues = RabbitOrderConfig.ORDER_SMS_QUEUE)public void handle(OrderCreatedMessage message, Channel channel, Message rawMessage) throws IOException { long deliveryTag = rawMessage.getMessageProperties().getDeliveryTag();
try { orderSmsApplicationService.sendOrderCreatedSms(message); channel.basicAck(deliveryTag, false); } catch (Exception ex) { channel.basicNack(deliveryTag, false, false); }}这里的重点不是背 API,而是理解:
- 成功后确认;
- 失败时不要假装成功;
- 是否重新入队要谨慎;
- 长期失败的消息应该进入死信队列或失败表。
如果一直 requeue=true,可能会导致同一条失败消息不断重试,拖垮消费者。
一个完整调用过程
把前面的内容串起来:
1. 用户创建订单2. OrderService 保存订单3. OrderMessageProducer 发送 OrderCreatedMessage4. 消息进入 order.topic.exchange5. routing key = order.created6. 交换机根据绑定规则投递到 order.sms.queue7. OrderSmsConsumer 监听队列8. 消费者调用短信服务9. 业务成功后确认消息这个链路就是 RabbitMQ 在业务系统里的最小闭环。
常见命名方式
命名建议保持统一。
交换机:
业务域.exchange类型.exchangeorder.topic.exchangeuser.topic.exchangecoupon.topic.exchange队列:
业务域.消费目的.queueorder.sms.queueorder.point.queueorder.search.queuerouting key:
业务域.事件动作order.createdorder.paidorder.cancelleduser.registeredcoupon.received统一命名能让排查问题轻松很多。
小结
这一篇完成了 Spring Boot 集成 RabbitMQ 的基本链路:
- 引入
spring-boot-starter-amqp; - 配置 RabbitMQ 连接;
- 定义消息对象;
- 声明交换机、队列和绑定;
- 使用
RabbitTemplate发送消息; - 使用
@RabbitListener消费消息; - 推荐使用 JSON 消息转换;
- 真实项目要继续补可靠性、幂等和失败处理。
If this article helped you, please share it with others!
Some information may be outdated






