Spring Kafka中kafka_deliveryAttempt头缺失问题及配置疑问
Spring Kafka中启用KafkaHeaders.DELIVERY_ATTEMPT头的正确配置方式
问题场景
我在Spring Kafka中实现了带有KafkaHeaders.DELIVERY_ATTEMPT头的监听器,代码如下:
@KafkaListener(topics = "${my_topic}", autoStartup = "true") public void listenForMessage(String message, @Header(KafkaHeaders.DELIVERY_ATTEMPT) int delivery) { log.info("Message from kafka topic {}", message); // 业务逻辑省略 }
根据官方文档说明,已添加@Header参数,但未启用容器的deliveryAttemptHeader属性(默认关闭以避免性能开销)时,出现以下错误:
org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method could not be invoked with the incoming message Endpoint handler details: Method [public void ****.listenForMessage(java.lang.String,int)] Bean [****MessageConsumer$$EnhancerBySpringCGLIB$$c2e2a4bd@3cc101bf]; nested exception is org.springframework.messaging.MessageHandlingException: Missing header 'kafka_deliveryAttempt' for method parameter type [int], failedMessage=GenericMessage [payload=.....
我尝试通过ContainerCustomizer配置启用该属性,但问题仍未解决(调试时已确认属性已设置):
@Bean public ContainerCustomizer<String, Message, ConcurrentMessageListenerContainer<String, Message>> containerCustomizer( ConcurrentKafkaListenerContainerFactory<String, Message> factory) { ContainerCustomizer<String, Message, ConcurrentMessageListenerContainer<String, Message>> custCustomizer = container -> { container.getContainerProperties().setDeliveryAttemptHeader(true); }; factory.setContainerCustomizer(custCustomizer); return custCustomizer; }
解决方案
问题出在ContainerCustomizer的配置时机或关联方式上,以下两种正确配置方式可解决问题:
方式一:直接配置容器工厂
无需单独定义ContainerCustomizer Bean,直接在ConcurrentKafkaListenerContainerFactory中启用属性:
@Bean public ConcurrentKafkaListenerContainerFactory<String, Message> kafkaListenerContainerFactory( ConsumerFactory<String, Message> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, Message> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); // 启用deliveryAttemptHeader属性 factory.getContainerProperties().setDeliveryAttemptHeader(true); return factory; }
方式二:正确关联ContainerCustomizer
若坚持使用ContainerCustomizer,需确保它被正确注入到容器工厂中,避免在Bean定义中直接操作工厂导致时机问题:
@Bean public ContainerCustomizer<String, Message, ConcurrentMessageListenerContainer<String, Message>> containerCustomizer() { return container -> { container.getContainerProperties().setDeliveryAttemptHeader(true); }; } @Bean public ConcurrentKafkaListenerContainerFactory<String, Message> kafkaListenerContainerFactory( ConsumerFactory<String, Message> consumerFactory, ContainerCustomizer<String, Message, ConcurrentMessageListenerContainer<String, Message>> customizer) { ConcurrentKafkaListenerContainerFactory<String, Message> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.setContainerCustomizer(customizer); return factory; }
关键说明
deliveryAttemptHeader必须设置为true,Spring Kafka才会在消息中添加kafka_deliveryAttempt头- 若自定义了
ConcurrentKafkaListenerContainerFactory,需确保@KafkaListener使用的是该自定义工厂实例(Spring默认优先使用自定义Bean) - 调试时需确认最终生成的容器实例中
deliveryAttemptHeader属性确实为true,避免因工厂未被正确应用导致配置无效
内容的提问来源于stack exchange,提问作者gstackoverflow
相关产品推荐
相关产品推荐

