You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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消息,配置参数值均正确。请问我的配置哪里存在问题?


问题排查与修复

核心问题分析

  1. 反序列化器配置冲突:
    在taskDtoConsumerFactory中,你既在配置Map里设置了VALUE_DESERIALIZER_CLASS_CONFIG = TaskDtoDeserializer.class,又在构造DefaultKafkaConsumerFactory时手动传入了ErrorHandlingDeserializer<TaskDto>。DefaultKafkaConsumerFactory会优先使用构造函数传入的反序列化器,重复配置会导致内部逻辑混乱,可能引发反序列化异常或消费者连接问题。

  2. 容器工厂未明确绑定:
    如果应用中存在多个ConcurrentKafkaListenerContainerFactory Bean,@KafkaListener默认会使用名为kafkaListenerContainerFactory的Bean,但为了避免歧义,最好显式指定。

  3. 连接断开日志关联:
    日志中的Node -1 disconnected通常表示消费者与Kafka集群的连接异常,结合配置冲突来看,大概率是反序列化失败导致消费者无法正常处理消息,进而触发连接中断、AdminClient重启的连锁反应。

修复步骤

  1. 移除重复的反序列化器配置:
    删除配置Map中多余的VALUE_DESERIALIZER_CLASS_CONFIG设置,因为已经通过构造函数传入了带错误处理的反序列化器:

    // 删掉此行
    // prop.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, TaskDtoDeserializer.class);
    
  2. 显式绑定容器工厂:
    在@KafkaListener注解中指定自定义的容器工厂,确保配置生效:

    @KafkaListener(
        topics = "${spring.kafka.consumer.from-hotfolder-transferer.topic}",
        containerFactory = "kafkaListenerContainerFactory"
    )
    
  3. 检查反序列化器实现:
    确认TaskDtoDeserializer能正确处理Topic中的消息格式,比如JSON结构是否匹配、字段类型是否兼容等。如果反序列化频繁失败,即使有ErrorHandlingDeserializer,也可能导致消费者无法正常工作。

  4. 验证最大轮询间隔配置:
    检查MAX_POLL_INTERVAL_MS_CONFIG的值,确保它大于单条消息的最大处理耗时。如果该值过小,Kafka会判定消费者超时,触发重平衡,也会出现类似的连接断开日志。


内容的提问来源于stack exchange,提问作者Marrakesh

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 15:34:53