Spring Boot集成RabbitMQ实战:可靠消息与延迟队列实现

微服务架构下服务间通信离不开消息中间件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/

(0)
小编小编
上一篇 6小时前
下一篇 6小时前

相关推荐