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

Spring Kafka异步处理手动ACK问题:解决@RetryableTopic偏移过期与重试丢失

问题分析与解决方案

问题根源

  1. 手动ACK与@RetryableTopic逻辑冲突:你手动调用acknowledgment.acknowledge()会直接提交偏移,完全绕过了@RetryableTopic的重试/DLT转发逻辑。框架需要根据消息处理的成功/失败状态来决定是否提交偏移、是否触发重试,手动ACK会导致失败消息还没进入重试流程就被确认,直接丢失。
  2. 错误信号未传递:doOnError仅记录日志但未重新抛出错误,导致Mono最终以成功状态结束,@RetryableTopic无法感知处理失败,因此不会触发重试,直接提交偏移。
  3. 异步乱序ACK:异步线程切换导致ACK顺序与消息消费顺序不一致,出现stale record或offsets list is empty错误,本质是偏移提交时机被手动操作打乱。

解决方案

1. 移除手动ACK,交由框架自动管理偏移

@RetryableTopic结合异步返回类型(Mono/Flux)时,Spring Kafka会自动根据异步流的状态处理偏移提交和重试逻辑,无需手动调用acknowledge()。

2. 确保错误信号正确传递

在doOnError后重新抛出错误,让@RetryableTopic感知到处理失败,触发重试流程。

3. 调整容器工厂的ACK配置

保持enable-auto-commit: false,将AckMode设置为MANUAL_IMMEDIATE,配合异步ACK确保框架能及时、正确地处理偏移提交。

修正后的完整代码

容器工厂配置

@Bean
public <K, V extends Message> ConcurrentKafkaListenerContainerFactory<K, V> kafkaListenerContainerFactory(ConsumerFactory<K, V> consumerFactory) {
    final var factory = new ConcurrentKafkaListenerContainerFactory<K,V>();
    factory.setConsumerFactory(consumerFactory);
    factory.getContainerProperties().setCommitLogLevel(Level.INFO);
    factory.setConcurrency(8);
    factory.getContainerProperties().setAsyncAcks(true);
    factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE); // 调整为MANUAL_IMMEDIATE
    return factory;
}

消息处理方法

@Component
class KafkaComponent {
    
    @RetryableTopic(attempts = "4", 
                    backoff = @Backoff(delay = 900000 /*15 min*/, multiplier = 3.0, maxDelay = 10800000), 
                    autoCreateTopics = "false",
                    topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE)
    @KafkaListener(topics = "topic", autoStartup = "true", containerFactory = "kafkaListenerContainerFactory")
    public Mono<Void> processMessage(PrecalcMessage message, 
                                     ConsumerRecordMetadata metadata, 
                                     @Header(name = RetryTopicHeaders.DEFAULT_HEADER_ATTEMPTS, defaultValue = "1") int attempt, 
                                     @Header(KafkaHeaders.RECEIVED_PARTITION) int partition, 
                                     Acknowledgement acknowledgement) {
        
        log.info("Received message {} from partition {} with metadata {}", message, partition, metadata);
        AtomicReference<String> resultRef = new AtomicReference<>("failed");

        ThreadContext.put("id", message.id);
        ThreadContext.put("date", message.date);
        ThreadContext.put("attempt", String.valueOf(attempt));

        return Mono.defer(() -> readAllMessageFromHbase(message).collectList())
            .flatMap(precalcInfos -> 
                aggregate(precalcInfos)
                    .flatMap(aggregates ->
                        Flux.fromIterable(List.of("hbase-1", "hbase-2"))
                            .parallel()
                            .flatMap(cluster -> save(cluster, aggregates))
                            .reduce(Result::merge)
                            .then()
                    )
            )
            .doOnSuccess(ctx -> resultRef.set("success")) // 移除手动ACK调用
            .doOnError(error -> log.error("Failed message processing", error))
            .onErrorResume(error -> Mono.error(error)) // 重新抛出错误,触发重试逻辑
            .doFinally(signal -> {
                metrics.publish("messageProcessing", resultRef.get());
                ThreadContext.clearAll();
            })
            .contextCapture()
            .subscribeOn(Schedulers.boundedElastic());                  
    }
}

额外注意事项

  • @RetryableTopic的attempts参数:设置为4表示总处理次数为4次(原始处理+3次重试),超出次数后消息会转发到DLT主题,确保不会丢失。
  • 异步流错误处理:确保所有业务异常都能被捕获并重新抛出,避免框架误判处理成功。
  • 线程上下文:contextCapture()已确保异步线程切换时上下文的传递,无需额外处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 23:30:01