Spring Kafka含Null键Map负载的Json序列化优化方案咨询
Kafka REST消息推送接口实现与优化咨询
实现代码详情
消息负载模型(SpecialData)
package com.learn.kafka.model; import com.fasterxml.jackson.annotation.JsonTypeInfo; import lombok.Data; import java.util.Map; @Data @JsonTypeInfo( use = JsonTypeInfo.Id.NAME, property = "type") public class SpecialData { Map<String, Object> messageInfo; }
Kafka消费者服务
package com.learn.kafka.service; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; import lombok.extern.slf4j.Slf4j; @Component @Slf4j public class ConsumerService { @KafkaListener(topics={"#{'${spring.kafka.topic}'}"},groupId="#{'${spring.kafka.consumer.group-id}'}") public void consumeMessage(String message){ log.info("Consumed message - {}",message); } }
Kafka生产者服务
package com.learn.kafka.service; import java.text.MessageFormat; import com.learn.kafka.model.SpecialData; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.kafka.core.KafkaTemplate; import lombok.extern.slf4j.Slf4j; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Service; @Service @Slf4j public class ProducerService{ @Value("${spring.kafka.topic:demo-topic}") String topicName; @Autowired KafkaTemplate<String,Object> kafkaTemplate; public String sendMessage(SpecialData messageModel){ log.info("Sending message from producer - {}",messageModel); Message message = constructMessage(messageModel); kafkaTemplate.send(message); return MessageFormat.format("Message Sent from Producer - {0}",message); } private Message constructMessage(SpecialData messageModel) { return MessageBuilder.withPayload(messageModel) .setHeader(KafkaHeaders.TOPIC,topicName) .setHeader("reason","for-Local-validation") .build(); } }
REST消息发送控制器
package com.learn.kafka.controller; import com.learn.kafka.model.SpecialData; import com.learn.kafka.service.ProducerService; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; import lombok.extern.slf4j.Slf4j; import java.util.HashMap; import java.util.Map; @RestController @RequestMapping("/api") @Slf4j public class MessageController { @Autowired private ProducerService producerService; @GetMapping("/send") public void sendMessage(){ SpecialData messageData = new SpecialData(); Map<String,Object> input = new HashMap<>(); input.put(null,"the key is null explicitly"); input.put("1","the key is one non-null"); messageData.setMessageInfo(input); producerService.sendMessage(messageData); } }
自定义Kafka序列化器
package com.learn.kafka; import com.fasterxml.jackson.core.JsonGenerator; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.SerializationFeature; import com.fasterxml.jackson.databind.SerializerProvider; import com.fasterxml.jackson.databind.ser.std.StdSerializer; import com.fasterxml.jackson.databind.util.StdDateFormat; import org.apache.kafka.common.errors.SerializationException; import org.apache.kafka.common.serialization.Serializer; import java.io.IOException; import java.util.Map; public class CustomSerializer implements Serializer<Object> { private static final ObjectMapper MAPPER = new ObjectMapper(); static { MAPPER.findAndRegisterModules(); MAPPER.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS); MAPPER.setDateFormat(new StdDateFormat().withColonInTimeZone(true)); MAPPER.getSerializerProvider().setNullKeySerializer(new NullKeySerializer()); } @Override public void configure(Map<String, ?> configs, boolean isKey) { } @Override public byte[] serialize(String topic, Object data) { try { if (data == null){ System.out.println("Null received at serializing"); return null; } System.out.println("Serializing..."); return MAPPER.writeValueAsBytes(data); } catch (Exception e) { e.printStackTrace(); throw new SerializationException("Error when serializing MessageDto to byte[]"); } } @Override public void close() { } static class NullKeySerializer extends StdSerializer<Object> { public NullKeySerializer() { this(null); } public NullKeySerializer(Class<Object> t) { super(t); } @Override public void serialize(Object obj, JsonGenerator gen, SerializerProvider provider) throws IOException { gen.writeFieldName("null"); } } }
application.yaml配置
spring: kafka: topic: input-topic consumer: bootstrap-servers: localhost:9092 group-id: input-group-id auto-offset-reset: earliest key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer producer: bootstrap-servers: localhost:9092 key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: com.learn.kafka.CustomSerializer properties: spring.json.add.type.headers: false
问题描述
当前代码可正常运行,能序列化包含Null键Map的SpecialData并发送至Kafka Broker,消费者使用StringDeserializer可正常接收打印消息。但考虑后续使用JsonDeserializer时可能出现问题,咨询以下两个优化方案是否可行或更优:
- 扩展Spring原生JsonSerializer,仅为其ObjectMapper添加NullKeySerializer?
- 是否可通过application.yaml简单配置实现Null键的序列化处理?
注:Null键序列化实现参考了Jackson处理Map空键的常规实现方式。
解答
1. 扩展Spring原生JsonSerializer是更优方案
Spring Kafka提供的JsonSerializer已经封装了Jackson ObjectMapper的基础配置,并且支持与Spring生态的其他配置(如全局ObjectMapper)协同工作。直接扩展它可以复用原有功能,仅添加Null键序列化的逻辑,避免重复造轮子。
实现示例:
package com.learn.kafka; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.SerializerProvider; import com.fasterxml.jackson.databind.ser.std.StdSerializer; import org.springframework.kafka.support.serializer.JsonSerializer; import java.io.IOException; public class CustomJsonSerializer extends JsonSerializer<Object> { public CustomJsonSerializer() { super(); configureNullKeySerializer(getObjectMapper()); } public CustomJsonSerializer(ObjectMapper objectMapper) { super(objectMapper); configureNullKeySerializer(objectMapper); } private void configureNullKeySerializer(ObjectMapper objectMapper) { SerializerProvider provider = objectMapper.getSerializerProvider(); provider.setNullKeySerializer(new NullKeySerializer()); } static class NullKeySerializer extends StdSerializer<Object> { public NullKeySerializer() { this(null); } public NullKeySerializer(Class<Object> t) { super(t); } @Override public void serialize(Object obj, com.fasterxml.jackson.core.JsonGenerator gen, SerializerProvider provider) throws IOException { gen.writeFieldName("null"); } } }
修改application.yaml中的生产者配置:
spring: kafka: producer: value-serializer: com.learn.kafka.CustomJsonSerializer
这种方式的优势:
- 复用Spring JsonSerializer的所有内置功能(如类型信息处理、日期格式化等)
- 可以直接使用Spring容器中配置的全局ObjectMapper(如果有)
- 代码更简洁,仅专注于添加Null键序列化逻辑
2. 无法仅通过application.yaml配置实现
Jackson的NullKeySerializer没有对应的配置属性,无法仅通过yaml配置直接启用。不过可以通过配置自定义ObjectMapper Bean,让Spring Kafka的JsonSerializer自动使用这个Bean,间接减少代码量:
定义全局ObjectMapper Bean:
package com.learn.kafka.config; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.SerializationFeature; import com.fasterxml.jackson.databind.util.StdDateFormat; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class JacksonConfig { @Bean public ObjectMapper objectMapper() { ObjectMapper mapper = new ObjectMapper(); mapper.findAndRegisterModules(); mapper.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS); mapper.setDateFormat(new StdDateFormat().withColonInTimeZone(true)); mapper.getSerializerProvider().setNullKeySerializer(new NullKeySerializer()); return mapper; } static class NullKeySerializer extends com.fasterxml.jackson.databind.ser.std.StdSerializer<Object> { public NullKeySerializer() { this(null); } public NullKeySerializer(Class<Object> t) { super(t); } @Override public void serialize(Object obj, com.fasterxml.jackson.core.JsonGenerator gen, com.fasterxml.jackson.databind.SerializerProvider provider) throws IOException { gen.writeFieldName("null"); } } }
修改application.yaml使用Spring原生JsonSerializer:
spring: kafka: producer: value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
这种方式不需要自定义序列化器,通过全局ObjectMapper配置实现Null键处理,也是一种简洁的方案。
内容的提问来源于stack exchange,提问作者Tim
相关产品推荐
相关产品推荐

