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
相关产品推荐
相关产品推荐

