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

SpringBoot Kafka消费者返回结果与生产者输入不匹配问题排查

问题根源

你遇到的问题核心在于生产者对JSONObject的序列化方式错误,以及消费者端未匹配正确的反序列化逻辑:

  1. 生产者使用JsonSerializer序列化org.json.JSONObject时,Jackson会将其视为普通Java对象,序列化它的内部属性(如mapType、empty),而非调用toString()输出你预期的地震数据JSON。你打印的message.toString()是正确的,但实际发送到Kafka的是JSONObject对象的属性JSON,而非业务数据。
  2. 消费者端的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 07:48:09