如何在reactor-kafka中对消费失败的ConsumerRecord进行重试
Reactor Kafka 单条消息消费重试配置方案
你当前配置的重试不生效,核心原因是使用了receiveAutoAck()方法:该方法会在消费者拉取到ConsumerRecord的同时自动提交偏移量,Kafka服务端会直接认为该条消息已消费成功,后续消费逻辑抛出异常时,无法再重新拉取该消息进行重试。
改造步骤
1. 切换为手动偏移量提交模式
将receiveAutoAck()替换为receive(),由业务逻辑控制偏移量提交时机:
reactiveKafkaConsumerTemplate .receive() .flatMap(record -> consumeWithRetry(record) // 消费成功后手动提交偏移量 .then(record.acknowledge()) // 单条消息所有重试都失败后的兜底处理 .onErrorResume(e -> { log.error("消息消费最终失败,topic: {}, 分区: {}, 偏移量: {}", record.topic(), record.partition(), record.offset(), e); // 可选:此处可添加投递死信队列的逻辑 // 如果希望后续还能重新消费该消息,可删除下一行ack代码 return record.acknowledge(); }) ) // 保留原有全局重试逻辑,仅处理重平衡、连接断开等消费链路级异常 .retryWhen(Retry.backoff(30, Duration.of(10, ChronoUnit.SECONDS))) .subscribe();
2. 调整重试逻辑,对齐Spring Kafka默认行为
retry(2)代表抛出异常后额外重试2次,加上首次执行总共消费3次,和Spring Kafka默认的3次重试规则完全一致:
public Mono<Void> consumeWithRetry(ConsumerRecord<String, MessageRecord> record) { MessageRecord message = record.value(); return consume(message) // 可选:可替换为带退避的重试策略,避免瞬时故障导致重试全部失败 // .retryWhen(Retry.backoff(2, Duration.ofSeconds(1))) .retry(2); } public Mono<Void> consume(MessageRecord message){ // 你的业务消费逻辑 return Mono.error(new RuntimeException("test retry")); }
注意事项
- 不要在消费流的全局位置使用
onErrorContinue处理单条消息错误,在flatMap内部通过onErrorResume处理单条消息的错误稳定性更高,不会意外跳过偏移量提交逻辑 - 如果不需要重试失败后立即提交偏移量,可删除兜底逻辑中的ack调用,消息会在消费者下次重平衡或者重启后重新拉取,需要注意调整Kafka消费者的
max.poll.interval.ms参数,避免重试耗时过长导致消费者被踢出消费组
内容的提问来源于stack exchange,提问作者Aniket Singla
相关产品推荐
相关产品推荐

