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

Kafka手动确认模式下消费失败消息无法立即重试问题求助

手动确认模式下Kafka Consumer无法立即重试失败消息的问题

我正在编写一个Kafka Consumer,已将确认属性设置为手动模式。当消费处理消息失败时我不会执行ack操作,现在希望消费者能立即重新处理这条失败消息,但未能实现。

当前消费者配置类

public ConsumerFactory<String, String> consumerFactory() {
    Map<String, Object> props = new HashMap<>();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAPSERVERS);
    props.put(ConsumerConfig.GROUP_ID_CONFIG, "${crmdsforecast.judjement-consumer-groupId}");
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class);
    props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,"earliest" );
    props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
    
    return new DefaultKafkaConsumerFactory<>(props);
}

@Bean
@ConditionalOnMissingBean(name = "kafkaListenerContainerFactory")
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {

    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
    factory.setConsumerFactory(consumerFactory());
    factory.getContainerProperties().setAckOnError(false);
    SeekToCurrentErrorHandler errorHandler = new SeekToCurrentErrorHandler((record, exception) -> {
        System.out.println("Error while processing the record {}"+  exception.getCause().getMessage());
    }, new FixedBackOff(3000L, 2L));
    factory.setErrorHandler(errorHandler);
    
    return factory;
}

当前消费者方法

@KafkaListener( containerFactory = "kafkaListenerContainerFactory",id = "${id}", topics = "${topicname}")
public void consume(String message,Acknowledgment acknowledgment){
    Dto payload = new Dto();
    try{
        payload = payloadDeserializer.convertIntoDtoObject(message);
        
        if(payload != null)
            //Do Something;
        
        acknowledgment.acknowledge();
    }
    catch(JsonProcessingException e){
        log.error("Error occured while Deserializing the String input {}",message);
        acknowledgment.acknowledge();
    }
    catch(Exception e){
        log.error("Some error occured while updating revenueLines {}",e.getMessage());
        //Here I expect that if error comes consumer should reread the message.
    }
}

已尝试的解决方案

方案1:使用SeekToCurrentBatchErrorHandler

修改后的containerFactory方法:

@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {

    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
    SeekToCurrentBatchErrorHandler errorHandler = new SeekToCurrentBatchErrorHandler();
    errorHandler.setBackOff( new ExponentialBackOff(500L, 2L));
    factory.setBatchErrorHandler(errorHandler);
    factory.setBatchListener(true);
    factory.setConsumerFactory(consumerFactory());
    
    return factory;
}

方案2:创建DefaultErrorHandler Bean并设置commonErrorHandler

@Bean
public DefaultErrorHandler errorHandler() {
    BackOff fixedBackOff = new FixedBackOff(1000, 3);
    DefaultErrorHandler errorHandler = new DefaultErrorHandler((consumerRecord, e) -> {
        System.out.println(String.format("consumed record %s because this exception was thrown",consumerRecord.toString(),e.getClass().getName()));
    }, fixedBackOff);
    errorHandler.addNotRetryableExceptions(NullPointerException.class);
    return errorHandler;
}

请问有人能帮忙解决这个问题吗?


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 13:07:48