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

Kafka消费者无法接收数据问题排查求助

Kafka消费者无法接收消息且应用立即停止的解决方案

核心问题修复步骤

1. 解决应用启动后立即终止的问题

Spring Boot非Web应用默认启动完成后会关闭上下文,导致消费者线程无法持续运行。可以通过以下方式保持应用活跃:

在主启动类中添加阻塞逻辑:

@SpringBootApplication
public class KafkaConsumerApplication {
    public static void main(String[] args) {
        ConfigurableApplicationContext context = SpringApplication.run(KafkaConsumerApplication.class, args);
        // 注册关闭钩子,确保优雅停机
        context.registerShutdownHook();
        // 阻塞主线程,保持应用运行
        try {
            synchronized (context) {
                context.wait();
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

或者添加spring-boot-starter-web依赖,让应用以Web服务方式运行(上下文不会自动关闭),若不需要Web服务,可在配置文件中设置:

# application.properties
spring.main.web-application-type=none

2. 修正Kafka Bootstrap地址配置

Kafka使用TCP协议通信,BOOTSTRAP_SERVERS_CONFIG不需要http://前缀,修改配置类中的对应项:

props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");

3. 修复配置类的访问权限

你的factory配置类是默认包访问权限,Spring组件扫描无法识别该类,需改为public并使用有意义的类名:

@Configuration
public class KafkaConsumerConfig {

    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        System.out.println("consumer factory created");
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        // 显式配置消费者组ID,与@KafkaListener保持一致
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "id");
        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
        System.out.println("concurrent Kafka listener container factory created");
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setConcurrency(3);
        factory.getContainerProperties().setPollTimeout(3000);
        return factory;
    }
}

4. 优化监听类的泛型定义

监听类的泛型D extends EventData未在监听方法中使用,可能导致Spring组件扫描异常,建议移除不必要的泛型:

@Service
public class Listener {

    @KafkaListener(groupId = "id", topics = "quickstart-events")
    public void listen(String vRecord) {
        System.out.println("Received message: " + vRecord);
    }
}

额外验证项

  • 确认Kafka服务正常运行,可通过telnet localhost 9092或nc -zv localhost 9092测试端口连通性
  • 检查quickstart-events主题是否存在,执行命令:kafka-topics.sh --list --bootstrap-server localhost:9092
  • 查看应用启动日志,确认输出consumer factory created和concurrent Kafka listener container factory created,说明配置类已被加载
  • 确保主启动类的@SpringBootApplication注解能扫描到配置类和监听类(包结构需在主类所在包或子包下)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 11:42:05