SpringBoot Kafka消费者返回结果与生产者输入不匹配问题排查
问题根源
你遇到的问题核心在于生产者对JSONObject的序列化方式错误,以及消费者端未匹配正确的反序列化逻辑:
- 生产者使用
JsonSerializer序列化org.json.JSONObject时,Jackson会将其视为普通Java对象,序列化它的内部属性(如mapType、empty),而非调用toString()输出你预期的地震数据JSON。你打印的message.toString()是正确的,但实际发送到Kafka的是JSONObject对象的属性JSON,而非业务数据。 - 消费者端的
JsonDeserializer收到这个错误的消息后,反序列化为JSONObject时,自然输出的是对象内部属性的字符串形式。
解决方案
方案一:使用JavaBean代替JSONObject(推荐)
用标准JavaBean封装业务数据,Jackson可以原生支持序列化/反序列化,避免JSONObject的兼容性问题。
1. 定义业务实体类
public class SeismicData { private String date; private int depth; private double latitude; private double longtitude; private double magnitude; private String time; // 必须提供无参构造方法 public SeismicData() {} // 生成所有字段的getter和setter方法 public String getDate() { return date; } public void setDate(String date) { this.date = date; } public int getDepth() { return depth; } public void setDepth(int depth) { this.depth = depth; } public double getLatitude() { return latitude; } public void setLatitude(double latitude) { this.latitude = latitude; } public double getLongtitude() { return longtitude; } public void setLongtitude(double longtitude) { this.longtitude = longtitude; } public double getMagnitude() { return magnitude; } public void setMagnitude(double magnitude) { this.magnitude = magnitude; } public String getTime() { return time; } public void setTime(String time) { this.time = time; } // 重写toString方便打印 @Override public String toString() { return "SeismicData{" + "date='" + date + '\'' + ", depth=" + depth + ", latitude=" + latitude + ", longtitude=" + longtitude + ", magnitude=" + magnitude + ", time='" + time + '\'' + '}'; } }
2. 修改生产者代码
@Service public class Producer { @Autowired KafkaTemplate<String, SeismicData> kafkaTemplate; // 直接发送JavaBean public void sendMessageToTopic(SeismicData message) { kafkaTemplate.send("seismic", message); } // 如果需要从JSONObject转换,添加该方法 public void sendFromJSONObject(JSONObject jsonObj) { SeismicData data = new SeismicData(); data.setDate(jsonObj.getString("date")); data.setDepth(jsonObj.getInt("depth")); data.setLatitude(jsonObj.getDouble("latitude")); data.setLongtitude(jsonObj.getDouble("longtitude")); data.setMagnitude(jsonObj.getDouble("magnitude")); data.setTime(jsonObj.getString("time")); kafkaTemplate.send("seismic", data); } }
3. 修改消费者代码
@Component public class KafkaListeners { @KafkaListener(topics="seismic", groupId = "groupId") void listener(SeismicData data) { System.out.println(data.toString()); } }
4. 修改消费者配置
@Configuration public class KafkaConsumerConfig { @Value("${spring.kafka.bootstrap-servers}") private String bootstrapServer; public Map<String, Object> consumerConfig() { Map<String, Object> properties = new HashMap<>(); properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer); properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); // 指定反序列化的目标类型 properties.put(JsonDeserializer.VALUE_DEFAULT_TYPE, SeismicData.class.getName()); // 信任所有包(生产环境可指定具体包) properties.put(JsonDeserializer.TRUSTED_PACKAGES, "*"); return properties; } @Bean public ConsumerFactory<String, SeismicData> consumerFactory() { return new DefaultKafkaConsumerFactory<>(consumerConfig()); } @Bean public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, SeismicData>> factory (ConsumerFactory<String, SeismicData> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, SeismicData> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); return factory; } }
5. 修改生产者配置
public Map<String, Object> producerConfig() { Map<String, Object> properties = new HashMap<>(); properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer); properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); return properties; }
方案二:直接传输JSON字符串(若必须保留JSONObject)
如果一定要使用JSONObject,可以将其转为字符串后发送,消费者接收字符串再转回JSONObject。
1. 修改生产者代码
@Service public class Producer { @Autowired KafkaTemplate<String, String> kafkaTemplate; public void sendMessageToTopic(JSONObject message) { // 将JSONObject转为字符串发送 kafkaTemplate.send("seismic", message.toString()); } }
2. 修改生产者配置
public Map<String, Object> producerConfig() { Map<String, Object> properties = new HashMap<>(); properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer); properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return properties; }
3. 修改消费者代码
@Component public class KafkaListeners { @KafkaListener(topics="seismic", groupId = "groupId") void listener(String data) { // 将字符串转回JSONObject JSONObject jsonObj = new JSONObject(data); System.out.println(jsonObj.toString()); } }
4. 修改消费者配置
@Configuration public class KafkaConsumerConfig { @Value("${spring.kafka.bootstrap-servers}") private String bootstrapServer; public Map<String, Object> consumerConfig() { Map<String, Object> properties = new HashMap<>(); properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer); properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); return properties; } @Bean public ConsumerFactory<String, String> consumerFactory() { return new DefaultKafkaConsumerFactory<>(consumerConfig()); } @Bean public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> factory (ConsumerFactory<String, String> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); return factory; } }
内容的提问来源于stack exchange,提问作者paul445
相关产品推荐
相关产品推荐

