RabbitMQ 消息队列入门与实战:从原理到 Spring Boot 整合

RabbitMQ 消息队列入门与实战:从原理到 Spring Boot 整合

admin
2026-07-24 / 0 评论 / 0 阅读

前言

RabbitMQ 是基于 AMQP 协议的开源消息中间件,广泛应用于微服务架构中的异步解耦、削峰填谷和可靠投递场景。本文将从核心概念讲起,最终完成一个 Spring Boot 整合 RabbitMQ 的实战项目。

一、核心概念

1.1 AMQP 模型

Producer → Exchange → (Binding) → Queue → Consumer
概念说明
Producer消息生产者
Exchange交换机,接收消息并路由到队列
Queue消息队列,存储消息
Binding交换机与队列之间的绑定关系
Routing Key路由键,Exchange 根据它决定消息去向
Consumer消息消费者

1.2 四种交换机类型

类型说明适用场景
Direct精确匹配 Routing Key点对点定向投递
Fanout广播到所有绑定队列广播通知
Topic模式匹配 Routing Key按主题订阅
Headers根据消息头匹配复杂路由条件

二、Docker 快速部署

docker run -d --name rabbitmq \
  -p 5672:5672 \
  -p 15672:15672 \
  -e RABBITMQ_DEFAULT_USER=admin \
  -e RABBITMQ_DEFAULT_PASS=admin123 \
  rabbitmq:3.13-management
  • 5672:AMQP 通信端口
  • 15672:管理界面端口,浏览器访问 http://localhost:15672

三、Spring Boot 整合

3.1 添加依赖

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

3.2 配置

spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: admin
    password: admin123
    virtual-host: /

3.3 队列与交换机配置

@Configuration
public class RabbitMQConfig {

    public static final String EXCHANGE = "order.exchange";
    public static final String QUEUE = "order.queue";
    public static final String ROUTING_KEY = "order.create";

    @Bean
    public TopicExchange orderExchange() {
        return new TopicExchange(EXCHANGE);
    }

    @Bean
    public Queue orderQueue() {
        return QueueBuilder.durable(QUEUE)
                .withArgument("x-message-ttl", 60000)
                .withArgument("x-dead-letter-exchange", "order.dlx")
                .build();
    }

    @Bean
    public Binding orderBinding() {
        return BindingBuilder.bind(orderQueue())
                .to(orderExchange())
                .with(ROUTING_KEY);
    }
}

3.4 生产者

@Service
public class OrderProducer {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    public void sendOrder(OrderDTO order) {
        rabbitTemplate.convertAndSend(
            RabbitMQConfig.EXCHANGE,
            RabbitMQConfig.ROUTING_KEY,
            order,
            message -> {
                message.getMessageProperties()
                       .setDeliveryMode(MessageDeliveryMode.PERSISTENT);
                return message;
            }
        );
    }
}

3.5 消费者

@Component
public class OrderConsumer {

    @RabbitListener(queues = RabbitMQConfig.QUEUE)
    public void handleOrder(OrderDTO order, Channel channel, Message message) throws IOException {
        try {
            processOrder(order);
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
        } catch (Exception e) {
            channel.basicNack(
                message.getMessageProperties().getDeliveryTag(),
                false,
                true
            );
        }
    }

    private void processOrder(OrderDTO order) {
        System.out.println("处理订单: " + order.getId());
    }
}

3.6 配置手动 ACK

spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: manual
        prefetch: 10
        retry:
          enabled: true
          max-attempts: 3

四、可靠性投递方案

4.1 生产者确认(Publisher Confirm)

@Configuration
public class RabbitConfig {
    @Bean
    public RabbitTemplate rabbitTemplate(ConnectionFactory cf) {
        RabbitTemplate template = new RabbitTemplate(cf);
        template.setConfirmCallback((correlationData, ack, cause) -> {
            if (!ack) {
                log.error("消息未到达Exchange: {}", cause);
            }
        });
        template.setReturnsCallback(returned -> {
            log.error("消息未路由到Queue: {}", returned.getMessage());
        });
        template.setMandatory(true);
        return template;
    }
}

4.2 死信队列(DLX)

消息变成死信的条件:

  1. 消息被消费者拒绝(basicNack + requeue=false)
  2. 消息 TTL 过期
  3. 队列达到最大长度
@Bean
public Queue deadLetterQueue() {
    return QueueBuilder.durable("order.dlq").build();
}

@Bean
public FanoutExchange deadLetterExchange() {
    return new FanoutExchange("order.dlx");
}

@Bean
public Binding dlxBinding() {
    return BindingBuilder.bind(deadLetterQueue()).to(deadLetterExchange());
}

五、常见问题

问题原因解决方案
消息丢失消费者宕机未 ACK手动 ACK + 持久化
消息重复网络抖动导致重复投递消费端幂等性(唯一ID去重)
消息堆积消费速度小于生产速度增加消费者 / 限流
队列阻塞单条消息处理过慢合理设置 prefetch

总结

RabbitMQ 在微服务架构中承担着异步解耦和流量削峰的关键角色。生产环境中务必做到:持久化 + 手动 ACK + 死信队列 + 消费幂等,才能保证消息的可靠投递。

0

评论 (0)

取消
0:00