微服务架构下服务间通信离不开消息中间件。RabbitMQ基于AMQP协议,支持多种交换机路由、死信队列、延迟消息,配合Spring Boot开发效率高。本文从依赖配置到延迟队列,完整实现一套可靠消息链路,覆盖生产者确认、消费者幂等、死信与延迟场景。
RabbitMQ核心概念与Spring Boot依赖
核心模型:Producer(生产者)- Exchange(交换机)- Queue(队列)- Consumer(消费者)。交换机负责路由,四种类型:direct精确匹配、topic通配、fanout广播、headers头匹配。业务上最常用topic。引入依赖:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
配置文件:
spring:
rabbitmq:
host: 10.0.0.12
port: 5672
username: admin
password: ${RABBIT_PASSWORD}
publisher-confirm-type: correlated # 生产者确认
publisher-returns: true
listener:
simple:
acknowledge-mode: manual # 手动确认
交换机、队列与死信配置
生产代码里用@Bean声明,运行时自动创建:
@Configuration
public class RabbitConfig {
public static final String EXCHANGE = "order.exchange";
public static final String QUEUE = "order.queue";
public static final String DLQ = "order.dlq";
@Bean
TopicExchange exchange() {
return new TopicExchange(EXCHANGE, true, false);
}
@Bean
Queue queue() {
return QueueBuilder.durable(QUEUE)
.withArgument("x-dead-letter-exchange", EXCHANGE)
.withArgument("x-dead-letter-routing-key", "order.dlq")
.build();
}
@Bean
Queue dlq() { return new Queue(DLQ_QUEUE, true); }
@Bean
Binding binding() {
return BindingBuilder.bind(queue()).to(exchange()).with("order.#");
}
}
可靠消息投递:确认回调与手动ack
生产者端开启确认,发布失败可感知并重试:
@Service
public class OrderPublisher {
@Autowired private RabbitTemplate rabbitTemplate;
public void publish(String routingKey, Object payload) {
CorrelationData cd = new CorrelationData(UUID.randomUUID().toString());
rabbitTemplate.setConfirmCallback((corr, ack, cause) -> {
if (!ack) log.error("消息确认失败: {} {}", corr, cause);
});
rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE, routingKey, payload, cd);
}
}
消费者端自动确认模式下,消息在回调执行后自动ack;处理失败可重试,配合死信队列丢弃/转人工。重试次数有限时,把消息发到DLQ再走补偿。
延迟队列实现订单超时关闭等定时逻辑
RabbitMQ原生延迟队列:消息先进一个不消费的队列,设置TTL,到期自动进入死信交换机,实现定时触发:
@Bean
Queue delayQueue() {
return QueueBuilder.durable("order.delay")
.ttl(30000) // 30秒
.deadLetterExchange(RabbitConfig.EXCHANGE)
.deadLetterRoutingKey("order.timeout")
.build();
}
@RabbitListener(queues = "order.timeout")
public void handleTimeout(Order order) {
orderService.autoCancel(order.getId());
}
这样订单超时未支付自动关闭、优惠券到期提醒等业务无需定时轮询数据库。
幂等消费与重复消息防护
MQ的at-least-once语义下,消费者必须幂等:用消息内业务ID做去重(Redis SETNX或数据库唯一索引)。幂等实现:
if (redis.setIfAbsent("msg:" + msgId, "1", 24h)) {
handleBiz(msg); // 首次处理
}
消息堆积排查与性能调优
线上问题三类:消费者处理慢导致堆积(加大并发、检查外部依赖);消息丢失(开启持久化+确认);重复消费(幂等兜底)。运维命令 rabbitmqctl list_queues,看消息量与消费者数。性能调优:预取prefetch控制在合理范围(默认250过高,调小到20-50),多消费者提高吞吐。RabbitMQ选型适合中小消息量、路由复杂的业务;超大吞吐换Kafka/Pulsar。链路加上追踪ID便于排错。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-ji-cheng-rabbitmq-shi-zhan-ke-kao-xiao-xi-yu-yan/