Spring Kafka批量监听器asyncAcks模式配置问题咨询
问题
使用Spring Kafka消费Kafka消息,希望采用批量监听器的asyncAcks模式处理主题,但配置后发现KafkaMessageListenerContainer的asyncAcks仍为false。
已编写的批量监听器代码:
@KafkaListener(groupId = "${app.kafka.shipment-historic-topic.group-id}", topics = "${app.kafka.shipment-historic-topic.topic-name}", autoStartup = "${app.kafka.shipment-historic-topic.enabled}", batch = "true", properties = {"spring.json.value.default.type=com.xxx.events.model.shipment.ShipmentEvent"} ) public void dispatchEvent(@Payload List<ConsumerRecord<String, ShipmentEvent>> shipmentEventRecords ) { //consumer logic }
当前的application.yaml配置:
spring: kafka: listener: ack-mode: BATCH async-acks: true
请问如何配置才能让批量监听器容器启用asyncAcks模式?
解决方案
确认Spring Kafka版本
async-acks对批量(BATCH) ack模式的支持是从Spring Kafka 2.7.0版本开始引入的,如果你的版本低于这个,该配置不会生效,请检查项目依赖中的Spring Kafka版本是否达标。检查配置优先级
@KafkaListener注解中的properties属性优先级高于全局yaml配置,如果注解内设置了与ack相关的属性,会覆盖全局配置,确保注解里没有覆盖async-acks的设置。验证容器工厂配置
如果项目自定义了ConcurrentKafkaListenerContainerFactory,需要确保工厂代码中没有手动设置asyncAcks=false。如果自定义了工厂,需在配置中显式开启:
@Bean public ConcurrentKafkaListenerContainerFactory<String, ShipmentEvent> kafkaListenerContainerFactory( ConsumerFactory<String, ShipmentEvent> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, ShipmentEvent> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.setBatchListener(true); // 显式开启asyncAcks factory.getContainerProperties().setAsyncAcks(true); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.BATCH); return factory; }
- 显式在注解中指定配置
如果全局配置仍不生效,可以直接在@KafkaListener的properties中添加async-acks配置,强制开启:
@KafkaListener( groupId = "${app.kafka.shipment-historic-topic.group-id}", topics = "${app.kafka.shipment-historic-topic.topic-name}", autoStartup = "${app.kafka.shipment-historic-topic.enabled}", batch = "true", properties = { "spring.json.value.default.type=com.xxx.events.model.shipment.ShipmentEvent", "spring.kafka.listener.async-acks=true" } )
内容的提问来源于stack exchange,提问作者Sanchelz
相关产品推荐
相关产品推荐

