本文目录导读:

我来详细介绍SpringBoot集成RabbitMQ实现消息收发。
添加依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
配置RabbitMQ
application.yml配置
spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
virtual-host: /
# 发送确认
publisher-confirm-type: correlated
# 发送回退
publisher-returns: true
# 手动ACK
listener:
simple:
acknowledge-mode: manual
prefetch: 1
消息队列配置
交换机、队列、绑定配置
@Configuration
public class RabbitConfig {
// 1. 直连交换机
public static final String DIRECT_EXCHANGE = "direct.exchange";
public static final String DIRECT_QUEUE = "direct.queue";
public static final String DIRECT_ROUTING_KEY = "direct.routing.key";
// 2. 主题交换机
public static final String TOPIC_EXCHANGE = "topic.exchange";
public static final String TOPIC_QUEUE = "topic.queue";
public static final String TOPIC_ROUTING_KEY = "topic.#";
// 3. 广播交换机
public static final String FANOUT_EXCHANGE = "fanout.exchange";
public static final String FANOUT_QUEUE_A = "fanout.queue.a";
public static final String FANOUT_QUEUE_B = "fanout.queue.b";
// 创建直连交换机
@Bean
public DirectExchange directExchange() {
return new DirectExchange(DIRECT_EXCHANGE, true, false);
}
// 创建主题交换机
@Bean
public TopicExchange topicExchange() {
return new TopicExchange(TOPIC_EXCHANGE, true, false);
}
// 创建广播交换机
@Bean
public FanoutExchange fanoutExchange() {
return new FanoutExchange(FANOUT_EXCHANGE, true, false);
}
// 创建队列
@Bean
public Queue directQueue() {
return QueueBuilder.durable(DIRECT_QUEUE).build();
}
@Bean
public Queue topicQueue() {
return QueueBuilder.durable(TOPIC_QUEUE).build();
}
@Bean
public Queue fanoutQueueA() {
return QueueBuilder.durable(FANOUT_QUEUE_A).build();
}
@Bean
public Queue fanoutQueueB() {
return QueueBuilder.durable(FANOUT_QUEUE_B).build();
}
// 绑定关系
@Bean
public Binding directBinding() {
return BindingBuilder.bind(directQueue())
.to(directExchange())
.with(DIRECT_ROUTING_KEY);
}
@Bean
public Binding topicBinding() {
return BindingBuilder.bind(topicQueue())
.to(topicExchange())
.with(TOPIC_ROUTING_KEY);
}
@Bean
public Binding fanoutBindingA() {
return BindingBuilder.bind(fanoutQueueA())
.to(fanoutExchange());
}
@Bean
public Binding fanoutBindingB() {
return BindingBuilder.bind(fanoutQueueB())
.to(fanoutExchange());
}
}
消息发送
消息实体类
@Data
@NoArgsConstructor
@AllArgsConstructor
public class OrderMessage implements Serializable {
private static final long serialVersionUID = 1L;
private String orderId;
private String userId;
private BigDecimal amount;
private LocalDateTime createTime;
}
消息发送者
@Component
@Slf4j
public class MessageSender {
@Autowired
private RabbitTemplate rabbitTemplate;
/**
* 发送直连交换机消息
*/
public void sendDirectMessage(OrderMessage orderMessage) {
CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString());
rabbitTemplate.convertAndSend(
RabbitConfig.DIRECT_EXCHANGE,
RabbitConfig.DIRECT_ROUTING_KEY,
orderMessage,
correlationData
);
log.info("直接交换机消息发送成功: {}", orderMessage);
}
/**
* 发送主题交换机消息
*/
public void sendTopicMessage(String routingKey, Object message) {
rabbitTemplate.convertAndSend(
RabbitConfig.TOPIC_EXCHANGE,
routingKey,
message
);
log.info("主题交换机消息发送成功: routingKey={}, message={}", routingKey, message);
}
/**
* 发送广播交换机消息
*/
public void sendFanoutMessage(Object message) {
rabbitTemplate.convertAndSend(
RabbitConfig.FANOUT_EXCHANGE,
"",
message
);
log.info("广播交换机消息发送成功: {}", message);
}
/**
* 延迟消息发送(需要安装延迟插件)
*/
public void sendDelayMessage(Object message, long delayTime) {
rabbitTemplate.convertAndSend(
"delay.exchange",
"delay.routing.key",
message,
msg -> {
msg.getMessageProperties().setDelay(Long.valueOf(delayTime).intValue());
return msg;
}
);
log.info("延迟消息发送成功: {}", message);
}
}
消息接收
手动ACK消息消费者
@Component
@Slf4j
public class MessageConsumer {
/**
* 直连队列消费者
*/
@RabbitListener(queues = RabbitConfig.DIRECT_QUEUE)
public void handleDirectMessage(OrderMessage message,
Message mqMessage,
Channel channel) throws IOException {
long deliveryTag = mqMessage.getMessageProperties().getDeliveryTag();
try {
log.info("接收到直连消息: {}", message);
// 业务处理
processOrder(message);
// 手动ACK
channel.basicAck(deliveryTag, false);
log.info("消息处理成功, deliveryTag: {}", deliveryTag);
} catch (Exception e) {
log.error("消息处理失败: {}", e.getMessage());
// 判断是否重试
boolean isRetry = mqMessage.getMessageProperties()
.getHeader("x-retry-count") != null &&
(Integer) mqMessage.getMessageProperties().getHeader("x-retry-count") < 3;
if (isRetry) {
// 重试,重新入队
channel.basicNack(deliveryTag, false, true);
} else {
// 丢弃或进入死信队列
channel.basicNack(deliveryTag, false, false);
log.warn("消息超过最大重试次数,已丢弃: {}", message);
}
}
}
/**
* 主题队列消费者
*/
@RabbitListener(queues = RabbitConfig.TOPIC_QUEUE)
public void handleTopicMessage(Object message,
Message mqMessage,
Channel channel) throws IOException {
long deliveryTag = mqMessage.getMessageProperties().getDeliveryTag();
try {
log.info("接收到主题消息: {}, routingKey={}",
message, mqMessage.getMessageProperties().getReceivedRoutingKey());
// 业务处理
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
channel.basicNack(deliveryTag, false, true);
}
}
/**
* 广播队列A消费者
*/
@RabbitListener(queues = RabbitConfig.FANOUT_QUEUE_A)
public void handleFanoutMessageA(Object message,
Message mqMessage,
Channel channel) throws IOException {
long deliveryTag = mqMessage.getMessageProperties().getDeliveryTag();
try {
log.info("广播队列A接收到消息: {}", message);
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
channel.basicNack(deliveryTag, false, true);
}
}
/**
* 广播队列B消费者
*/
@RabbitListener(queues = RabbitConfig.FANOUT_QUEUE_B)
public void handleFanoutMessageB(Object message,
Message mqMessage,
Channel channel) throws IOException {
long deliveryTag = mqMessage.getMessageProperties().getDeliveryTag();
try {
log.info("广播队列B接收到消息: {}", message);
channel.basicAck(deliveryTag, false);
} catch (Exception e) {
channel.basicNack(deliveryTag, false, true);
}
}
private void processOrder(OrderMessage orderMessage) {
// 业务处理逻辑
log.info("处理订单: {}", orderMessage.getOrderId());
}
}
回调配置
@Component
@Slf4j
public class RabbitTemplateConfig implements RabbitTemplate.ConfirmCallback,
RabbitTemplate.ReturnsCallback {
@PostConstruct
public void init() {
rabbitTemplate.setConfirmCallback(this);
rabbitTemplate.setReturnsCallback(this);
rabbitTemplate.setMandatory(true);
}
@Autowired
private RabbitTemplate rabbitTemplate;
/**
* 消息发送确认回调
*/
@Override
public void confirm(CorrelationData correlationData,
boolean ack,
String cause) {
if (ack) {
log.info("消息发送成功: correlationId={}", correlationData.getId());
} else {
log.error("消息发送失败: correlationId={}, cause={}",
correlationData.getId(), cause);
// 可以进行重试处理
}
}
/**
* 消息回退回调(消息无法路由到队列时触发)
*/
@Override
public void returnedMessage(ReturnedMessage returnedMessage) {
log.error("消息无法路由: exchange={}, routingKey={}, replyCode={}, replyText={}",
returnedMessage.getExchange(),
returnedMessage.getRoutingKey(),
returnedMessage.getReplyCode(),
returnedMessage.getReplyText());
// 可以将消息保存到数据库,进行后续处理
}
}
使用示例
@RestController
@RequestMapping("/mq")
@Slf4j
public class RabbitController {
@Autowired
private MessageSender messageSender;
@PostMapping("/send/direct")
public String sendDirect(@RequestBody OrderMessage orderMessage) {
messageSender.sendDirectMessage(orderMessage);
return "发送成功";
}
@PostMapping("/send/topic")
public String sendTopic(@RequestParam String routingKey,
@RequestBody Object message) {
messageSender.sendTopicMessage(routingKey, message);
return "发送成功";
}
@PostMapping("/send/fanout")
public String sendFanout(@RequestBody Object message) {
messageSender.sendFanoutMessage(message);
return "发送成功";
}
}
重要配置说明
- 手动ACK:确保消息被正确处理后才确认
- 消息持久化:队列和消息需要设置durable=true
- 重试机制:设置合理重试次数,避免死循环
- 回调机制:实现ConfirmCallback和ReturnsCallback监控消息状态
- 死信队列:处理失败消息,避免消息丢失
这样配置,就可以在SpringBoot中完整地使用RabbitMQ进行消息收发。