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

