Kafka消费者监听无记录接收问题排查求助
我开发了一个简单的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

