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

Kafka消费者监听无记录接收问题排查求助

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

我开发了一个简单的Kafka应用,用于消费JSON格式的用户数据。本地已运行端口为9092的Kafka服务及ZooKeeper,拥有符合User.json schema的data.json模拟数据,且User.java实体类与该schema结构匹配。配置了KafkaConfig、KafkaListeners等相关组件后,应用可正常启动,但KafkaMessageListenerContainer始终接收0条记录,无法消费数据。相关代码如下,恳请提供排查及解决建议:

KafkaConfig代码

import com.microservices.kafka.model.User;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.context.annotation.Bean;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.*;
import org.springframework.kafka.support.converter.StringJsonMessageConverter;
import org.springframework.kafka.support.serializer.JsonDeserializer;
import org.springframework.kafka.support.serializer.JsonSerializer;
import org.springframework.stereotype.Component;

import java.util.HashMap;
import java.util.Map;

@Component
@EnableKafka
public class KafkaConfig {

    private static final String BOOTSTRAP_SERVERS = "localhost:9092";

    public static ProducerFactory<String, User> producerFactory() {
        Map<String, Object> config = new HashMap<>();
        config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS);
        config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
        return new DefaultKafkaProducerFactory<>(config);
    }

    public static ConsumerFactory<String, User> consumerFactory(String groupId) {
        Map<String, Object> config = new HashMap<>();
        config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS);
        config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
        config.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        return new DefaultKafkaConsumerFactory<>(config, new StringDeserializer(),
                new JsonDeserializer<>(User.class));
    }

    @Bean
    public KafkaTemplate<String, User> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, User> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, User> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory("employeeGroup"));
        factory.setMessageConverter(new StringJsonMessageConverter());
        return factory;
    }
}

User实体类代码

package com.microservices.kafka.model;

import lombok.AllArgsConstructor;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import lombok.Setter;

import javax.persistence.Entity;
import javax.persistence.Id;

@Entity
@Getter
@Setter
@AllArgsConstructor
@RequiredArgsConstructor
public class User {

    @Id
    private int id;

    private String firstName;

    private String lastName;

    private String email;

    private String status;
}

KafkaApplication启动类代码

package com.microservices.kafka;

import com.microservices.kafka.config.KafkaConfig;
import com.microservices.kafka.model.User;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.core.io.ClassPathResource;
import org.springframework.core.io.Resource;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.converter.JsonMessageConverter;
import org.springframework.kafka.support.converter.RecordMessageConverter;
import org.springframework.kafka.support.converter.StringJsonMessageConverter;
import org.springframework.kafka.support.serializer.JsonDeserializer;

import java.io.IOException;
import java.io.InputStreamReader;
import java.io.Reader;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;

@SpringBootApplication(exclude = {DataSourceAutoConfiguration.class})
public class KafkaApplication {

    private static final String TOPIC_NAME = "employeeTopic";
    private static final String GROUP_ID = "employeeGroup";

    private static final ConsumerRecord<String, User> record = null;

    private static final Logger logger = org.slf4j.LoggerFactory.getLogger(KafkaApplication.class);

    public static void main(String[] args) {
        SpringApplication.run(KafkaApplication.class, args);
    }

    @Bean
    public RecordMessageConverter converter() {
        return new StringJsonMessageConverter();
    }

    @Bean
    public JsonMessageConverter jsonConverter() {
        return new JsonMessageConverter();
    }

    @Bean
    public JsonDeserializer<User> deserializer() {
        return new JsonDeserializer<>(User.class);
    }

    @Bean
    public String loadData() throws IOException {
        Resource resource = new ClassPathResource("/data/data.json");
        byte[] bytes = Files.readAllBytes(Paths.get(resource.getURI()));
        logger.info("Loading data from file " + resource.getFilename());
        return new String(bytes);
    }
}

KafkaListeners监听类代码

import com.microservices.kafka.model.User;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
import java.io.IOException;
import com.fasterxml.jackson.databind.ObjectMapper;

/**
 * The @Autowired annotation is used to inject an instance of ObjectMapper into the class.
 * This is an instance of Jackson's ObjectMapper class, which is used to deserialize JSON data into Java objects.
 *
 * The @KafkaListener annotation is used to define a Kafka listener that will listen for messages
 * on the employeeTopic topic. The groupId parameter specifies the ID of the consumer group that this
 * listener is part of. The containerFactory parameter specifies the name of the Kafka listener container factory to use.
 *
 * The listen() method is the actual message listener method that will be called when a message is
 * received. It takes a ConsumerRecord<String, String> parameter that contains the Kafka message.
 * In this case, we assume that the message contains a JSON string representing a User object.
 *
 * The readValue() method of the ObjectMapper class is used to deserialize the JSON string into a User object.
 * We assume that the JSON string is stored in the value() property of the ConsumerRecord.
 *
 * Finally, we print out the received User object to the console.
 */
@Component
public class KafkaListeners {

    @Autowired
    private ObjectMapper objectMapper;

    @KafkaListener(
            topics = "employeeTopic",
            groupId = "group_id",
            containerFactory = "kafkaListenerContainerFactory"
    )
    public void listen(ConsumerRecord<String, String> record) throws IOException {
        String fileContent = record.value();
        User user = objectMapper.readValue(fileContent, User.class);
        System.out.println("Received user: " + user);
    }
}

User.json Schema

{
  "type": "object",
  "javaType": "com.microservices.kafka.model.User",
  "properties": {
    "id": {
      "type": "integer"
    },
    "firstName": {
      "type": "string"
    },
    "lastName": {
      "type": "string"
    },
    "email": {
      "type": "string"
    },
    "status": {
      "type": "string"
    }
  }
}

排查及解决建议

1. 修正消费组ID不匹配问题

KafkaConfig中消费者工厂使用的组ID是employeeGroup,但KafkaListeners里的@KafkaListener指定的是group_id,两者不一致会导致消费者无法加入正确的消费组,无法读取对应偏移量。
修复:将监听类的groupId改为employeeGroup:

@KafkaListener(
        topics = "employeeTopic",
        groupId = "employeeGroup",
        containerFactory = "kafkaListenerContainerFactory"
)

2. 统一消息类型与泛型配置

配置的ConcurrentKafkaListenerContainerFactory泛型是<String, User>,说明期望直接接收User类型消息,但监听方法参数是ConsumerRecord<String, String>,类型不匹配会导致消息转换失败被丢弃。
修复:修改监听方法参数,移除手动反序列化逻辑:

public void listen(ConsumerRecord<String, User> record) {
    User user = record.value();
    System.out.println("Received user: " + user);
}

同时可以删除监听类中的ObjectMapper注入,容器已通过配置的JsonDeserializer完成反序列化。

3. 添加消息发送逻辑

当前代码仅加载了data.json,但没有将数据发送到employeeTopic,消费者自然无法接收。
修复:在启动类中添加应用就绪后的消息发送逻辑:

@Autowired
private KafkaTemplate<String, User> kafkaTemplate;

@Autowired
private ObjectMapper objectMapper;

@EventListener(ApplicationReadyEvent.class)
public void sendData() throws IOException {
    String jsonData = loadData();
    // 若data.json是单个User对象
    User user = objectMapper.readValue(jsonData, User.class);
    kafkaTemplate.send("employeeTopic", String.valueOf(user.getId()), user);
    // 若data.json是User数组,遍历发送
    // User[] users = objectMapper.readValue(jsonData, User[].class);
    // for (User u : users) {
    //     kafkaTemplate.send("employeeTopic", String.valueOf(u.getId()), u);
    // }
    logger.info("数据已发送到employeeTopic");
}

4. 补充User类无参构造函数

JsonDeserializer需要类有无参构造函数才能完成反序列化,当前User类仅包含带参构造函数,会导致反序列化失败。
修复:在User类上添加@NoArgsConstructor注解:

@Entity
@Getter
@Setter
@AllArgsConstructor
@RequiredArgsConstructor
@NoArgsConstructor
public class User {
    // 原有字段
}

5. 验证Topic存在性与权限

  • 用Kafka命令行工具检查employeeTopic是否存在:
    kafka-topics.sh --list --bootstrap-server localhost:9092
    
  • 若不存在则创建:
    kafka-topics.sh --create --topic employeeTopic --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1
    
  • 本地环境一般无需权限配置,若有ACL需确认消费者读取权限。

6. 开启调试日志定位问题

在application.yml中添加日志配置,查看消费者连接、分区分配、偏移量等细节:

logging:
  level:
    org.springframework.kafka: DEBUG
    org.apache.kafka.clients.consumer: DEBUG

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 15:22:01