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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 17:36:03