Spring Boot Kafka手动ACK失败问题求助
问题描述
使用Spring Boot 2.7.2搭建Kafka消费者,尝试切换为手动确认模式,配置了spring.kafka.listener.ack-mode=MANUAL_IMMEDIATE,但消费者接收消息时持续抛出异常,核心错误提示:
java.lang.IllegalStateException: No Acknowledgment available as an argument, the listener container must have a MANUAL AckMode to populate the Acknowledgment.
用户提供的相关代码与配置
配置文件
kafka.consumer.groupId=mcs-ccp-event message.topic.name=mcs_ccp_test kafka.bootstrapAddress=kafka-dev-app-a1.com:9092 spring.kafka.consumer.auto-offset-reset=earliest spring.kafka.listener.ack-mode=MANUAL_IMMEDIATE
消费者配置类
@EnableKafka @Configuration @Slf4j public class KafkaConsumerConfig { @Value(value = "${kafka.bootstrapAddress}") private String bootstrapAddress; @Value(value = "${kafka.consumer.groupId}") private String groupId; public ConsumerFactory<String, Event> eventConsumerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress); props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); //props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(JsonDeserializer.VALUE_DEFAULT_TYPE, "com.xxx.mcsccpkafkaconsumer.vo.Event"); props.put(JsonDeserializer.USE_TYPE_INFO_HEADERS,false); return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), new JsonDeserializer<>(Event.class)); } @Bean public ConcurrentKafkaListenerContainerFactory<String, Event> eventKafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, Event> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(eventConsumerFactory()); return factory; } }
监听器代码
@KafkaListener(topics = "${message.topic.name}", containerFactory = "eventKafkaListenerContainerFactory", groupId = "${kafka.consumer.groupId}") public void eventListener(@Payload Event event, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition, Acknowledgment acknowledgment) { log.info("Received event message: {} from partition : {}", event, partition); persistEventToDB(event); acknowledgment.acknowledge(); this.eventLatch.countDown(); }
错误栈核心片段
org.springframework.kafka.listener.ListenerExecutionFailedException: invokeHandler Failed; nested exception is java.lang.IllegalStateException: No Acknowledgment available as an argument, the listener container must have a MANUAL AckMode to populate the Acknowledgment. Caused by: java.lang.IllegalStateException: No Acknowledgment available as an argument, the listener container must have a MANUAL AckMode to populate the Acknowledgment. Caused by: org.springframework.messaging.converter.MessageConversionException: Cannot convert from [com.xxx.mcsccpkafkaconsumer.vo.Event] to [org.springframework.kafka.support.Acknowledgment]
问题原因
自定义的ConcurrentKafkaListenerContainerFactory没有显式设置AckMode,配置文件中的spring.kafka.listener.ack-mode只会对默认的容器工厂生效,自定义工厂需要手动绑定AckMode参数,否则容器依然使用默认的自动提交模式,无法注入Acknowledgment对象,导致参数解析失败。
另外,手动确认模式下需要确保禁用自动提交,虽然Spring Kafka在MANUAL AckMode下会自动禁用ENABLE_AUTO_COMMIT_CONFIG,但显式配置更清晰。
解决方法
修改消费者配置类,在创建容器工厂时显式设置AckMode,并开启禁用自动提交:
@EnableKafka @Configuration @Slf4j public class KafkaConsumerConfig { @Value(value = "${kafka.bootstrapAddress}") private String bootstrapAddress; @Value(value = "${kafka.consumer.groupId}") private String groupId; public ConsumerFactory<String, Event> eventConsumerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress); props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); // 显式禁用自动提交 props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(JsonDeserializer.VALUE_DEFAULT_TYPE, "com.xxx.mcsccpkafkaconsumer.vo.Event"); props.put(JsonDeserializer.USE_TYPE_INFO_HEADERS,false); return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), new JsonDeserializer<>(Event.class)); } @Bean public ConcurrentKafkaListenerContainerFactory<String, Event> eventKafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, Event> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(eventConsumerFactory()); // 显式设置手动立即确认模式 factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); return factory; } }
验证说明
修改后,容器会使用MANUAL_IMMEDIATE模式,能够正确注入Acknowledgment对象,调用acknowledge()方法即可完成手动提交偏移量,异常会消失。
内容的提问来源于stack exchange,提问作者Ladu anand

