You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RabbitMQ重启时Spring AMQP生产者消息丢失问题求助

RabbitMQ集群重启时Spring应用消息丢失问题

问题背景

使用Bitnami的RabbitMQ Chart部署3副本集群,通过LoadBalancer类型Service对外提供服务,应用配置spring.rabbitmq.addresses=<LOADBALANCER IP>。测试quorum queues(镜像经典队列也存在该问题)时发现,执行kubectl rollout restart statefulset/rabbitmq重启Broker期间,若应用正连接到重启的Pod发送消息,会出现消息丢失:百万消息最多丢失2万条。

需求

  • 临时连接错误时确保所有消息在连接恢复后送达;
  • 若无法实现,需为每条未送达消息生成日志。

测试过程与结果

1. 初始代码

@RestController
@RequestMapping("/")
@CommonsLog
public class RabbitController {

    private final RabbitTemplate rabbitTemplate;


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


    @GetMapping
    public ResponseEntity<String> index() {
        for (int i = 0; i < 1_000_000; i++) {
            try {
                this.sendMessage("myexchange", "myrouting", String.valueOf(i));
            } catch (Exception e) {
                log.error("Could not send message", e);
            }
        }
        return ResponseEntity.noContent().build();
    }


    private void sendMessage(String exchange, String routingKey, Object content) {
        MessageProperties properties = new MessageProperties();
        properties.setDeliveryMode(MessageDeliveryMode.PERSISTENT);
        Message message = new Message(content.toString().getBytes(), properties);
        this.rabbitTemplate.convertAndSend(exchange, routingKey, message);
    }
}
  • 重启Broker时仅生成1条错误日志,队列仅收到959,372条消息。

2. 配置重试机制

spring:
  rabbitmq:
    template:
      retry:
        enabled: true
        max-attempts: 100
        multiplier: 2
        max-interval: 120000
  • 丢失消息更多,仅909,813条入队。

3. 添加mandatory配置

spring:
  rabbitmq:
    template:
      mandatory: true
      retry:
        enabled: true
        max-attempts: 100
        multiplier: 2
        max-interval: 120000
  • 仍仅947,848条消息入队。

4. 启用publisher-confirm

spring:
  rabbitmq:
    publisher-confirm-type: correlated
    template:
      mandatory: true
      retry:
        enabled: true
        max-attempts: 100
        multiplier: 2
        max-interval: 120000
  • 发送15万条消息,仅丢失30条。

5. 启用publisher-returns

spring:
  rabbitmq:
    publisher-confirm-type: correlated
    publisher-returns: true
    template:
      mandatory: true
      retry:
        enabled: true
        max-attempts: 100
        multiplier: 2
        max-interval: 120000
  • 出现额外错误日志,但仍有28条消息丢失。

6. 移除异常捕获

移除sendMessage的try-catch后,程序在连接断开后直接停止,仅30,524条消息入队。


解决方案建议

1. 修正重试机制的错误使用

默认的RabbitTemplate重试是同步重试,在Broker重启期间,连接断开后重试会占用线程,导致后续消息无法发送甚至被丢弃。需改为异步重试+本地消息表的组合方案:

  • 发送前将消息持久化到本地数据库(本地消息表),标记为"待发送";
  • 异步发送消息,通过RabbitTemplate.confirmCallback接收确认结果:
    • 收到确认后更新本地消息表状态为"已发送";
    • 收到否定确认或超时,触发异步重试(使用Spring Retry的异步配置,避免阻塞主线程)。

2. 优化RabbitMQ连接配置

  • 配置连接工厂的超时与心跳参数,快速感知Broker状态:
spring:
  rabbitmq:
    connection-timeout: 5000
    requested-heartbeat: 10
    connection-factory:
      automatic-recovery-enabled: true
      network-recovery-interval: 5000

3. 避免LoadBalancer单点问题

LoadBalancer可能将请求路由到重启中的Pod,建议直接配置所有RabbitMQ节点地址到应用:

spring:
  rabbitmq:
    addresses: rabbitmq-0.rabbitmq-headless.default.svc.cluster.local:5672,rabbitmq-1.rabbitmq-headless.default.svc.cluster.local:5672,rabbitmq-2.rabbitmq-headless.default.svc.cluster.local:5672

Spring AMQP会自动尝试连接其他可用节点,减少连接到故障节点的概率。

4. 完善消息追踪与日志

  • 为每条消息生成唯一ID,存入MessageProperties:
properties.setMessageId(UUID.randomUUID().toString());
  • 在confirmCallback中记录未确认消息的详细信息:
rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
    if (!ack) {
        log.error("Message not confirmed, correlationId: {}, cause: {}", correlationData.getId(), cause);
        // 关联本地消息表标记为待重试
    }
});

5. 调整发送逻辑为异步批处理

将同步循环发送改为异步批处理,避免单线程阻塞:

@GetMapping
public ResponseEntity<String> index() {
    IntStream.range(0, 1_000_000)
        .mapToObj(String::valueOf)
        .forEach(content -> CompletableFuture.runAsync(() -> {
            try {
                sendMessage("myexchange", "myrouting", content);
            } catch (Exception e) {
                log.error("Failed to send message: {}", content, e);
            }
        }, executor));
    return ResponseEntity.noContent().build();
}

需配置合适的线程池,避免资源耗尽。


内容的提问来源于stack exchange,提问作者Linus

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 08:44:53