Spring Kafka手动确认消息后仍重试异常消息问题排查
Kafka消费者手动确认后仍重试的问题解决
问题描述
使用@KafkaListener实现的Kafka消费者,需求是先确认Kafka消息,再执行REST API调用。但执行REST API抛出预期异常后,即便已在调用前手动确认消息,消费者仍会重试这条消息。
相关代码
消费者代码
@KafkaListener(topics = "${spring.kafka.topic.name}", groupId = "${consumer.topicGroupId}") public void listenEvent(ConsumerRecord<String, Event> consumerRecord, Acknowledgment acknowledgment) throws IOException { acknowledgment.acknowledge(); patchService.patch(arguments...); }
PatchService代码
public class PatchService { public String patch(arguments....) { try { ResponseEntity<String> response = restTemplate.exchange(uri, HttpMethod.PATCH, request, String.class); return response.getBody(); } catch (HttpClientErrorException ex) { log.error("Error updating API having error response {}", ex.getResponseBodyAsString()); throw new APIException(ex.getStatusCode(), ex.getResponseBodyAsString()); } } }
消费者配置代码
@Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() throws FileNotFoundException { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<String, String>(); factory.setConsumerFactory(consumerFactory()); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); return factory; }
consumerFactory中已配置:
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
问题原因
Spring Kafka的监听器容器默认逻辑是:只要监听器方法抛出未捕获的异常,就会判定消息处理失败,进而触发重试机制,哪怕已经手动调用了acknowledge()确认消息。容器的重试触发逻辑和手动确认是独立的,异常抛出会覆盖已完成的确认动作,导致消息被重新投递。
解决方案
要避免重试,核心是不让监听器方法抛出未捕获的异常,具体可按以下方式修改:
方式一:在消费者方法内捕获目标异常
在listenEvent方法内部捕获PatchService抛出的APIException,仅记录日志而不向上抛出:
@KafkaListener(topics = "${spring.kafka.topic.name}", groupId = "${consumer.topicGroupId}") public void listenEvent(ConsumerRecord<String, Event> consumerRecord, Acknowledgment acknowledgment) throws IOException { acknowledgment.acknowledge(); try { patchService.patch(arguments...); } catch (APIException e) { log.error("处理REST API调用失败,消息已确认无需重试", e); // 仅记录日志,不抛出异常 } }
方式二:全局关闭重试机制
如果希望全局禁用所有消息的重试逻辑,可在容器工厂配置中关闭重试:
@Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() throws FileNotFoundException { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<String, String>(); factory.setConsumerFactory(consumerFactory()); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); // 关闭重试 factory.setRetryTemplate(null); return factory; }
优先推荐方式一,因为可以针对特定异常做精准处理,不会影响其他正常异常的重试逻辑。
内容的提问来源于stack exchange,提问作者Ankur Goel
相关产品推荐
相关产品推荐

