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

Spring Kafka消费JSON消息时如何将可变结构的data字段解析为字符串

实现方案

步骤1:自定义JSON反序列化器,实现任意JSON节点转字符串

该反序列化器会直接将读取到的JSON节点(不管是对象、数组还是基础类型)转换为原始JSON字符串返回:

import com.fasterxml.jackson.core.JsonParser;
import com.fasterxml.jackson.databind.DeserializationContext;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.deser.std.StdDeserializer;
import java.io.IOException;

public class RawJsonToStringDeserializer extends StdDeserializer<String> {
    public RawJsonToStringDeserializer() {
        super(String.class);
    }

    @Override
    public String deserialize(JsonParser parser, DeserializationContext context) throws IOException {
        JsonNode node = parser.getCodec().readTree(parser);
        return node.toString();
    }
}

步骤2:定义消息接收实体类

在data字段上指定使用上述自定义反序列化器即可:

import com.fasterxml.jackson.databind.annotation.JsonDeserialize;
import lombok.Data;

@Data // 无Lombok则手动添加getter、setter方法
public class KafkaConsumeMessage {
    private String identifier;
    @JsonDeserialize(using = RawJsonToStringDeserializer.class)
    private String data;
}

步骤3:修改Kafka消费者配置

将值反序列化器替换为Spring自带的JsonDeserializer,指定反序列化目标类为上面定义的实体类:

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.support.serializer.JsonDeserializer;
import java.util.HashMap;
import java.util.Map;

@Bean
public ConsumerFactory<String, KafkaConsumeMessage> consumerFactory() {
    Map<String, Object> config = new HashMap<>();
    config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaConfig.getConsumerBootstrapServers());
    config.put(ConsumerConfig.GROUP_ID_CONFIG, GROUP_ID_CONFIG);
    config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
    // 替换为你自己的KafkaConsumeMessage全类名
    config.put(JsonDeserializer.VALUE_DEFAULT_TYPE, "com.yourpackage.KafkaConsumeMessage");
    config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, AUTO_OFFSET_RESET_CONFIG);
    // 可选配置:JSON存在未定义字段时不报错
    config.put(JsonDeserializer.FAIL_ON_UNKNOWN_PROPERTIES, false);

    return new DefaultKafkaConsumerFactory<>(config, new StringDeserializer(),
            new JsonDeserializer<>(KafkaConsumeMessage.class));
}

@Bean
public ConcurrentKafkaListenerContainerFactory<String, KafkaConsumeMessage> kafkaListenerFactory(ConsumerFactory<String, KafkaConsumeMessage> consumerFactory) {
    ConcurrentKafkaListenerContainerFactory<String, KafkaConsumeMessage> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);
    return factory;
}

步骤4:直接用实体类接收消息

@KafkaListener(topics = "你的Topic名称", containerFactory = "kafkaListenerFactory")
public void consumeMessage(KafkaConsumeMessage message) {
    // 直接获取固定结构的identifier字段
    String identifier = message.getIdentifier();
    // 获取data字段的原始JSON字符串,可按需后续再单独解析
    String dataRawJson = message.getData();
    // 业务逻辑处理
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 01:45:04