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
相关产品推荐
相关产品推荐

