KafkaListener无法读取Topic消息:配置问题排查求助
我编写了带有@KafkaListener注解的DocumentListener类,代码如下:
@Slf4j @Component @RequiredArgsConstructor public class DocumentListener { @KafkaListener(topics = "${spring.kafka.consumer.from-hotfolder-transferer.topic}") public void listen(@Payload TaskDto taskDto, @Header(KafkaConst.HEADER_KEY) String serviceName) { ... some logic } }
消息发送至Topic后,日志输出如下:
2023-11-08T15:30:05.783+03:00 INFO [worker-ocr ] 28396 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer : gpn-1: partitions assigned: [task-to-ocr-0] 2023-11-08T15:39:03.175+03:00 INFO [worker-ocr ] 28396 --- [ntainer#0-0-C-1] org.apache.kafka.clients.NetworkClient : [Consumer clientId=gpn-1-0, groupId=gpn-1] Node -1 disconnected. 2023-11-08T15:58:27.194+03:00 INFO [worker-ocr ] 28396 --- [ntainer#0-0-C-1] o.a.k.clients.admin.AdminClientConfig : AdminClientConfig values: ______ here some values kafka admin ______ 2023-11-08T15:58:27.219+03:00 INFO [worker-ocr ] 28396 --- [ntainer#0-0-C-1] o.a.kafka.common.utils.AppInfoParser : Kafka version: 3.4.1 2023-11-08T15:58:27.219+03:00 INFO [worker-ocr ] 28396 --- [ntainer#0-0-C-1] o.a.kafka.common.utils.AppInfoParser : Kafka commitId: 8a516edc2755df89 2023-11-08T15:58:27.219+03:00 INFO [worker-ocr ] 28396 --- [ntainer#0-0-C-1] o.a.kafka.common.utils.AppInfoParser : Kafka startTimeMs: 1699448307219 2023-11-08T15:58:27.235+03:00 INFO [worker-ocr ] 28396 --- [| adminclient-1] o.a.kafka.common.utils.AppInfoParser : App info kafka.admin.client for adminclient-1 unregistered 2023-11-08T15:58:27.242+03:00 INFO [worker-ocr ] 28396 --- [| adminclient-1] o.apache.kafka.common.metrics.Metrics : Metrics scheduler closed 2023-11-08T15:58:27.243+03:00 INFO [worker-ocr ] 28396 --- [| adminclient-1] o.apache.kafka.common.metrics.Metrics : Closing reporter org.apache.kafka.common.metrics.JmxReporter 2023-11-08T15:58:27.243+03:00 INFO [worker-ocr ] 28396 --- [| adminclient-1] o.apache.kafka.common.metrics.Metrics : Metrics reporters closed
我的消费者配置类代码如下:
@Configuration @EnableKafka @RequiredArgsConstructor public class KafkaConfig { private final KafkaProducerProperties kafkaProducerProperties; @Bean public ConsumerFactory<String, TaskDto> taskDtoConsumerFactory(KafkaProperties kafkaProperties, KafkaCustomProperties kafkaCustomProperties) { KafkaCustomProperties.KafkaConsumerProperties consumer = kafkaCustomProperties.getConsumer().get(KafkaConst.CONSUMER_CONFIG); Map<String, Object> prop = kafkaProperties.buildConsumerProperties(); prop.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); prop.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, TaskDtoDeserializer.class); prop.put(ConsumerConfig.GROUP_ID_CONFIG, consumer.getClientId()); prop.put(ConsumerConfig.CLIENT_ID_CONFIG, consumer.getClientId()); prop.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, consumer.getPollInterval()); return new DefaultKafkaConsumerFactory<>( prop, new StringDeserializer(), new ErrorHandlingDeserializer<>(new TaskDtoDeserializer()) ); } @Bean public ConcurrentKafkaListenerContainerFactory<String, TaskDto> kafkaListenerContainerFactory( ConsumerFactory<String, TaskDto> taskDtoConsumerFactory, KafkaCustomProperties kafkaCustomProperties) { ConcurrentKafkaListenerContainerFactory<String, TaskDto> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(taskDtoConsumerFactory); factory.setCommonErrorHandler(errorHandler(kafkaCustomProperties)); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.RECORD); factory.getContainerProperties().setObservationEnabled(true); factory.afterPropertiesSet(); return factory; } }
已确认Kafka服务正常,用大数据工具可正常读取该Topic消息,配置参数值均正确。请问我的配置哪里存在问题?
核心问题分析
反序列化器配置冲突:
在taskDtoConsumerFactory中,你既在配置Map里设置了VALUE_DESERIALIZER_CLASS_CONFIG = TaskDtoDeserializer.class,又在构造DefaultKafkaConsumerFactory时手动传入了ErrorHandlingDeserializer<TaskDto>。DefaultKafkaConsumerFactory会优先使用构造函数传入的反序列化器,重复配置会导致内部逻辑混乱,可能引发反序列化异常或消费者连接问题。容器工厂未明确绑定:
如果应用中存在多个ConcurrentKafkaListenerContainerFactoryBean,@KafkaListener默认会使用名为kafkaListenerContainerFactory的Bean,但为了避免歧义,最好显式指定。连接断开日志关联:
日志中的Node -1 disconnected通常表示消费者与Kafka集群的连接异常,结合配置冲突来看,大概率是反序列化失败导致消费者无法正常处理消息,进而触发连接中断、AdminClient重启的连锁反应。
修复步骤
移除重复的反序列化器配置:
删除配置Map中多余的VALUE_DESERIALIZER_CLASS_CONFIG设置,因为已经通过构造函数传入了带错误处理的反序列化器:// 删掉此行 // prop.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, TaskDtoDeserializer.class);显式绑定容器工厂:
在@KafkaListener注解中指定自定义的容器工厂,确保配置生效:@KafkaListener( topics = "${spring.kafka.consumer.from-hotfolder-transferer.topic}", containerFactory = "kafkaListenerContainerFactory" )检查反序列化器实现:
确认TaskDtoDeserializer能正确处理Topic中的消息格式,比如JSON结构是否匹配、字段类型是否兼容等。如果反序列化频繁失败,即使有ErrorHandlingDeserializer,也可能导致消费者无法正常工作。验证最大轮询间隔配置:
检查MAX_POLL_INTERVAL_MS_CONFIG的值,确保它大于单条消息的最大处理耗时。如果该值过小,Kafka会判定消费者超时,触发重平衡,也会出现类似的连接断开日志。
内容的提问来源于stack exchange,提问作者Marrakesh

