人人都会AI编程

14.4 消息队列整合:RabbitMQ / Kafka 基础使用与可靠投递

更新时间:2026-07-10

消息队列在现代分布式系统中承担着削峰填谷、异步解耦、事件驱动等重要职责。Spring 生态对 RabbitMQ 和 Kafka 都提供了深度的整合支持,开发者无需手动管理连接、通道或消费者线程,即可快速集成消息能力。然而,引入消息队列后,消息的可靠投递就成为必须严肃对待的课题——消息丢失、重复消费、顺序错乱等问题一旦出现,排查和修复成本极高。

本节将分别介绍 Spring Boot 如何整合 RabbitMQ 和 Kafka,涵盖基本收发、序列化配置以及可靠投递的关键实践。

14.4.1 RabbitMQ 基础整合

1. 依赖与配置

在 Spring Boot 项目中引入 RabbitMQ 支持,只需添加起步依赖:

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

然后在 application.yml 中配置连接信息:

spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
    virtual-host: /
    # 以下为生产环境建议配置
    publisher-confirm-type: correlated   # 开启发送方确认
    publisher-returns: true              # 开启路由失败回退
    template:
      mandatory: true                    # 消息无法路由时回调
    listener:
      simple:
        acknowledge-mode: manual         # 手动确认模式(默认 auto)
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 1000ms

2. 声明队列、交换器与绑定

推荐通过配置类显式声明消息系统的拓扑结构,避免依赖消费者自动创建带来的不可控风险:

@Configuration
public class RabbitConfig {

    public static final String ORDER_EXCHANGE = "order.exchange";
    public static final String ORDER_QUEUE = "order.queue";
    public static final String ORDER_ROUTING_KEY = "order.created";

    @Bean
    public DirectExchange orderExchange() {
        return new DirectExchange(ORDER_EXCHANGE, true, false);
    }

    @Bean
    public Queue orderQueue() {
        return QueueBuilder.durable(ORDER_QUEUE).build();
    }

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

3. 发送消息

使用 Spring 提供的 RabbitTemplate 即可完成消息发送:

@Component
public class OrderMessageSender {

    private final RabbitTemplate rabbitTemplate;

    public OrderMessageSender(RabbitTemplate rabbitTemplate) {
        this.rabbitTemplate = rabbitTemplate;
    }

    public void sendOrderCreated(Order order) {
        // 可将对象转为 JSON 后发送,需要配置 MessageConverter
        rabbitTemplate.convertAndSend(
            RabbitConfig.ORDER_EXCHANGE,
            RabbitConfig.ORDER_ROUTING_KEY,
            order
        );
    }
}

为了让对象自动序列化为 JSON,需配置 MessageConverter

@Bean
public MessageConverter jsonMessageConverter() {
    return new Jackson2JsonMessageConverter();
}

@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
    RabbitTemplate template = new RabbitTemplate(connectionFactory);
    template.setMessageConverter(jsonMessageConverter());
    return template;
}

4. 接收消息

消费者通过 @RabbitListener 注解绑定队列,消息体的反序列化由消息转换器自动完成:

@Component
public class OrderMessageListener {

    @RabbitListener(queues = RabbitConfig.ORDER_QUEUE)
    public void handleOrderCreated(Order order, Channel channel, Message message) 
            throws IOException {
        try {
            // 执行业务逻辑...
            System.out.println("收到订单创建消息:" + order.getId());
            
            // 手动确认
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
        } catch (Exception e) {
            // 根据业务决定是否重新入队
            channel.basicNack(message.getMessageProperties().getDeliveryTag(), 
                              false, true); // 重新入队
        }
    }
}

acknowledge-mode: manual 下,必须手动调用 basicAck 确认消息,否则消息会一直处于未确认状态,直到连接断开后重新入队。

14.4.2 RabbitMQ 可靠投递实践

1. 发送方确认(Publisher Confirm)

仅将消息写入 Socket 缓冲区并不代表 Broker 已成功接收。开启 publisher-confirm-type: correlated 后,可为每条消息设置确认回调:

@Component
public class ReliableOrderSender {

    private final RabbitTemplate rabbitTemplate;

    public ReliableOrderSender(RabbitTemplate rabbitTemplate) {
        this.rabbitTemplate = rabbitTemplate;
    }

    public void sendOrderWithConfirm(Order order) {
        CorrelationData correlationData = new CorrelationData(order.getId());
        // 设置确认回调
        correlationData.getFuture().whenComplete((confirm, throwable) -> {
            if (throwable != null) {
                // 发送异常,记录日志,启动补偿
                log.error("消息发送异常: {}", order.getId(), throwable);
                compensationService.retryOrAlert(order);
            } else if (confirm.isAck()) {
                log.info("消息已确认到达Broker: {}", order.getId());
            } else {
                log.warn("消息被Broker拒绝: {}", order.getId());
                // 可实现重发或告警
            }
        });
        rabbitTemplate.convertAndSend(
            RabbitConfig.ORDER_EXCHANGE,
            RabbitConfig.ORDER_ROUTING_KEY,
            order,
            correlationData
        );
    }
}

2. 路由失败处理

若消息发布到交换器后无法路由到任何队列(例如路由键写错),开启 mandatory: true 并注册 ReturnCallback 可捕获这类失败:

@PostConstruct
public void initReturnCallback() {
    rabbitTemplate.setReturnsCallback(returned -> {
        log.warn("消息无法路由: 交换器={} 路由键={} 消息内容={}",
                 returned.getExchange(), returned.getRoutingKey(),
                 new String(returned.getMessage().getBody()));
        // 可存入数据库或重发
    });
}

3. 消费者手动确认与重试

消费者侧最常见的问题是消费失败后直接丢失消息,或陷入无限重试的死循环。通过手动确认模式结合重试,可以实现精细化控制:

  • 业务异常(如库存不足)不应重试,应当在代码中 catch 后直接确认并记录或发送补偿事件。
  • 临时故障(如网络超时、数据库死锁)可以重试,但要限制次数并设置退避策略。
@RabbitListener(queues = RabbitConfig.ORDER_QUEUE)
public void handleWithRetry(Order order, Channel channel, Message message) throws IOException {
    long deliveryTag = message.getMessageProperties().getDeliveryTag();
    try {
        processOrder(order);
        channel.basicAck(deliveryTag, false);
    } catch (RetryableException e) {
        // 临时故障,重试有限次
        if (hasExceedMaxRetry(message)) {
            channel.basicNack(deliveryTag, false, false); // 拒绝且不重新入队
            compensationService.handleFailure(order, e);
        } else {
            channel.basicNack(deliveryTag, false, true); // 重新入队
        }
    } catch (Exception e) {
        // 非重试异常,直接拒绝并记录
        channel.basicNack(deliveryTag, false, false);
        compensationService.handleFailure(order, e);
    }
}

可以通过消息头记录重试次数,每次重新入队时次数递增。

4. 消息幂等性保证

即使消息确认机制再完善,网络抖动仍可能导致 Broker 重复投递。消费者必须实现幂等处理:

  • 为每条消息分配业务唯一标识(如订单号),并在消费前检查 Redis 或数据库是否已处理。
  • 利用数据库唯一约束,将“处理消息”与“记录处理状态”放在同一事务中。

14.4.3 Kafka 基础整合

1. 依赖与配置

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>
spring:
  kafka:
    bootstrap-servers: localhost:9092
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      acks: all                    # 等待所有副本确认
      retries: 3                   # 发送重试
      enable-idempotence: true     # 幂等发送,配合 acks=all
    consumer:
      group-id: order-group
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: "com.example.dto"
      enable-auto-commit: false    # 关闭自动提交,手动控制位移
    listener:
      ack-mode: manual             # 手动确认

2. 发送消息

使用 KafkaTemplate 发送消息:

@Component
public class OrderKafkaSender {

    private final KafkaTemplate<String, Order> kafkaTemplate;

    public OrderKafkaSender(KafkaTemplate<String, Order> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    public void sendOrder(String key, Order order) {
        ListenableFuture<SendResult<String, Order>> future = 
            kafkaTemplate.send("order-topic", key, order);
        future.addCallback(
            result -> log.info("消息发送成功: offset={}", 
                               result.getRecordMetadata().offset()),
            ex -> {
                log.error("消息发送失败", ex);
                // 启动补偿
            }
        );
    }
}

3. 接收消息

@Component
public class OrderKafkaListener {

    @KafkaListener(topics = "order-topic")
    public void onMessage(ConsumerRecord<String, Order> record, 
                          Acknowledgment ack) {
        try {
            Order order = record.value();
            processOrder(order);
            ack.acknowledge(); // 手动提交位移
        } catch (Exception e) {
            log.error("消费失败: offset={}", record.offset(), e);
            // 不提交位移,消息会再次消费;注意避免无限重试
        }
    }
}

14.4.4 Kafka 可靠投递与消费处理

1. 生产者可靠性

  • acks=all:确保所有 ISR 副本确认后才认为消息发送成功,避免 Leader 宕机丢失。
  • enable-idempotence=true:配合 acks=all,Broker 会对生产者消息去重,防止重试导致重复。
  • 发送结果异步回调:记录失败消息或重试,不能捕获异常就忽略。
  • 合理设置 retriesdelivery.timeout.ms:对于临时故障进行重试,但要设置超时防止长时间阻塞。

2. 消费者手动位移提交

Kafka 消费者的位移管理决定消息的投递语义。关闭自动提交后,应在处理完业务逻辑后再提交位移,实现至少一次语义。若要精确一次,需将业务结果与位移存放在同一个事务中(使用 Kafka 事务或外部协调)。

@KafkaListener(topics = "order-topic")
public void onMessage(ConsumerRecord<String, Order> record, 
                      Acknowledgment ack) {
    try {
        orderService.processAndMarkProcessed(record.value(), record.offset());
        ack.acknowledge();
    } catch (DuplicateException e) {
        // 已处理过,直接提交位移,避免重复处理
        ack.acknowledge();
    } catch (Exception e) {
        // 记录错误,不提交位移,线程挂起或暂停一段时间后继续
        throw e; // 配置中的 SeekToCurrentErrorHandler 会处理
    }
}

推荐配置 SeekToCurrentErrorHandler 配合 DeadLetterPublishingRecoverer,将消费多次失败的消息转移到死信主题,避免阻塞分区:

@Bean
public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory(
        ConsumerFactory<Object, Object> consumerFactory,
        KafkaTemplate<Object, Object> template) {
    ConcurrentKafkaListenerContainerFactory<Object, Object> factory = 
        new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);
    factory.setCommonErrorHandler(new DefaultErrorHandler(
        new DeadLetterPublishingRecoverer(template),
        new FixedBackOff(1000L, 3)  // 重试3次,间隔1秒
    ));
    return factory;
}

3. 幂等性

即使启用了生产者的幂等和消费者手动提交,业务层仍需考虑重复消费的情况。全局唯一标识 + 业务状态检查 是通用方案。例如在消息体中携带事件 ID,消费前在 Redis 中尝试设置一个具有过期时间的键(SETNX)或插入一张去重表。

void processOrder(Order order) {
    String eventId = order.getEventId();
    if (redisTemplate.opsForValue().setIfAbsent(eventId, "1", Duration.ofMinutes(30))) {
        // 未处理,执行核心业务
        businessLogic(order);
    } else {
        log.info("重复消息忽略: eventId={}", eventId);
    }
}

14.4.5 实践建议

  1. 选择合适的消息队列:RabbitMQ 擅长复杂的路由和消息管理,适合需要投递确认、死信队列等场景;Kafka 擅长高吞吐的日志和流式处理,适合事件溯源和实时数据管道。
  2. 消息体设计:尽量包含业务唯一 ID、事件产生时间戳和必要上下文,便于幂等、监控和补偿。
  3. 监控与告警:监控消息积压、消费延迟、确认失败率等指标,及时发现问题。
  4. 消费端做好幂等:不要完全依赖 Broker 的投递语义,业务层的幂等是最后的防线。
  5. 死信与补偿:对无法成功处理的消息,导入死信队列或数据库,并提供管理界面进行人工干预或定时补偿。

消息队列的整合看似简单,但可靠投递涉及生产端确认、消费端提交、重试策略、幂等设计等多个环节。Spring 的整合让我们能聚焦于这些可靠性考量本身,而不必陷入底层 API 的细节中。在实际项目中,根据业务对一致性的要求,选择合适级别的保障方案,是架构师和开发者必须掌握的技能。