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

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时可能出现问题,咨询以下两个优化方案是否可行或更优:

  1. 扩展Spring原生JsonSerializer,仅为其ObjectMapper添加NullKeySerializer?
  2. 是否可通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 09:05:15