KafkaListener 批量消费异常时回退偏移量

1.实现BatchErrorHandler

@Component
public  class KafkaBatchExceptionHandler implements BatchErrorHandler {

    @Override
    public void handle(Exception e, ConsumerRecords<?, ?> consumerRecords) {
    }

    @Override
    public void handle(Exception thrownException, ConsumerRecords<?, ?> consumerRecords, Consumer<?, ?> consumer) {
        Map<TopicPartition, Long> offsetsToReset = new HashMap<>();
        consumerRecords.forEach((record) -> {
            offsetsToReset.compute(new TopicPartition(record.topic(), record.partition()),
                    (k, v) -> v == null ? record.offset() : Math.min(v, record.offset()));

        });
        offsetsToReset.forEach((k, v) -> consumer.seek(k, v));

    }


}

2.注册进KafkaListenerContainerFactory中

   @Bean
    KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, GenericRecord>> kafkaListenerContainerFactory(@Autowired KafkaBatchExceptionHandler batchErrorHandler) {
        ConcurrentKafkaListenerContainerFactory<String, GenericRecord> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setConcurrency(null);
        factory.setBatchListener(true);
        factory.getContainerProperties()
                .setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
        factory.getContainerProperties().setPollTimeout(30000);
        factory.setBatchErrorHandler(batchErrorHandler);
        return factory;
    }
©著作权归作者所有,转载或内容合作请联系作者
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

推荐阅读更多精彩内容