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

为何Kafka JsonSerializer无法序列化ProducerRecord?

Spring Kafka使用JsonSerializer发送JSON对象时出现ProducerRecord序列化异常问题

问题描述

尝试使用Spring Kafka的JsonSerializer从生产者发送JSON对象,但发送消息时出现与ProducerRecord序列化相关的异常。原本期望JsonSerializer能处理消息负载的序列化,却发现Jackson在尝试序列化包含负载对象及Kafka元数据的ProducerRecord时失败。

生产者配置

import com.app.integration.impl.dto.TransitECS;
import com.app.integration.impl.util.KafkaConstants;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.kafka.support.serializer.JsonSerializer;

import java.util.Map;

@Configuration
public class KafkaProducerConfig {

    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Bean
    public Map<String, Object> producerConfigs() {
        return Map.of(
                ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers,
                ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class,
                ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class
        );
    }

    @Bean
    public ProducerFactory<String, TransitECS> transitECSProducerFactory() {
        return new DefaultKafkaProducerFactory<>(producerConfigs());
    }

    @Bean
    public KafkaTemplate<String, TransitECS> transitECSKafkaTemplate() {
        KafkaTemplate<String, TransitECS> kafkaTemplate = new KafkaTemplate<>(transitECSProducerFactory());
        kafkaTemplate.setDefaultTopic(KafkaConstants.TRANSIT_ECS_TOPIC);
        return kafkaTemplate;
    }
}

发送方法

@RestController
@RequestMapping("/api/v1/topic")
@RequiredArgsConstructor
public class TopicController {
    private final KafkaTemplate<String, TransitECS> transitECSKafkaTemplate;

    /**
     * 该方法用于向主题发送Transit ECS
     *
     * @param transitECS 要发送的Transit ECS
     */
    @PostMapping("/send/transitECS")
    public CompletableFuture<SendResult<String, TransitECS>> sendTransit(@Valid @RequestBody TransitECS transitECS) {
        return this.transitECSKafkaTemplate.sendDefault(transitECS);
    }
}

异常信息

com.fasterxml.jackson.databind.exc.InvalidDefinitionException: No serializer found for class org.apache.kafka.clients.producer.ProducerRecord and no properties discovered to create BeanSerializer (to avoid exception, disable SerializationFeature.FAIL_ON_EMPTY_BEANS) (through reference chain: org.springframework.kafka.support.SendResult["producerRecord"])

引发异常的代码片段

package com.fasterxml.jackson.databind.ser.impl;

public class UnknownSerializer
    extends ToEmptyObjectSerializer // since 2.13
{
   ...

    @Override
    public void serialize(Object value, JsonGenerator gen, SerializerProvider ctxt) throws IOException
    {
        // 27-Nov-2009, tatu: As per [JACKSON-201] may or may not fail...
        if (ctxt.isEnabled(SerializationFeature.FAIL_ON_EMPTY_BEANS)) {
            failForEmpty(ctxt, value); // <------ 进入此处
        }
        super.serialize(value, gen, ctxt);
    }

    protected void failForEmpty(SerializerProvider prov, Object value)
            throws JsonMappingException {
        Class<?> cl = value.getClass();
        if (NativeImageUtil.needsReflectionConfiguration(cl)) {
            prov.reportBadDefinition(handledType(), String.format(
                    "No serializer found for class %s and no properties discovered to create BeanSerializer (to avoid exception, disable SerializationFeature.FAIL_ON_EMPTY_BEANS). This appears to be a native image, in which case you may need to configure reflection for the class that is to be serialized",
                    cl.getName()));
        } else {
            prov.reportBadDefinition(handledType(), String.format(
                    "No serializer found for class %s and no properties discovered to create BeanSerializer (to avoid exception, disable SerializationFeature.FAIL_ON_EMPTY_BEANS)",
                    cl.getName()));
        }
    }
}

已尝试的操作

  • 禁用SerializationFeature.FAIL_ON_EMPTY_BEANS但未解决问题,修改后的ProducerFactory配置如下:
@Bean
public ProducerFactory<String, TransitECS> transitECSProducerFactory() {
    ObjectMapper mapper = new ObjectMapper();
    mapper.setVisibility(PropertyAccessor.FIELD, JsonAutoDetect.Visibility.ANY);
    mapper.configure(SerializationFeature.FAIL_ON_EMPTY_BEANS, false);
    return new DefaultKafkaProducerFactory<>(
            producerConfigs(),
            new StringSerializer(),
            new JsonSerializer<>(mapper)
    );
}

解决方案

问题的核心不是Kafka生产者的序列化配置错误,而是:
控制器返回的CompletableFuture<SendResult>会被Spring MVC自动序列化为响应内容,而SendResult内部包含ProducerRecord对象,这个Kafka内部类没有对应的Jackson序列化器,导致Spring MVC在序列化响应时抛出异常。

解决方式(三选一即可)

  1. 返回自定义响应对象:只封装业务需要的发送结果信息,避免返回包含ProducerRecord的SendResult
@PostMapping("/send/transitECS")
public CompletableFuture<Map<String, Object>> sendTransit(@Valid @RequestBody TransitECS transitECS) {
    return this.transitECSKafkaTemplate.sendDefault(transitECS)
            .thenApply(sendResult -> Map.of(
                    "success", true,
                    "topic", sendResult.getRecordMetadata().topic(),
                    "partition", sendResult.getRecordMetadata().partition(),
                    "offset", sendResult.getRecordMetadata().offset()
            ));
}
  1. 不返回发送结果:如果前端不需要发送结果,直接返回CompletableFuture<Void>
@PostMapping("/send/transitECS")
public CompletableFuture<Void> sendTransit(@Valid @RequestBody TransitECS transitECS) {
    return this.transitECSKafkaTemplate.sendDefault(transitECS)
            .thenAccept(sendResult -> {});
}
  1. 配置Spring MVC的ObjectMapper忽略ProducerRecord:如果必须返回SendResult,可以通过配置Spring MVC的全局ObjectMapper,对ProducerRecord进行忽略处理
@Configuration
public class WebConfig implements WebMvcConfigurer {
    @Override
    public void extendMessageConverters(List<HttpMessageConverter<?>> converters) {
        for (HttpMessageConverter<?> converter : converters) {
            if (converter instanceof MappingJackson2HttpMessageConverter jacksonConverter) {
                ObjectMapper mapper = jacksonConverter.getObjectMapper();
                mapper.addMixIn(ProducerRecord.class, ProducerRecordMixIn.class);
            }
        }
    }

    abstract static class ProducerRecordMixIn {
        @JsonIgnore
        abstract ProducerRecord<?, ?> getProducerRecord();
    }
}

为什么禁用FAIL_ON_EMPTY_BEANS无效?

因为你配置的ObjectMapper是给Kafka生产者的JsonSerializer使用的,而抛出异常的是Spring MVC序列化响应时使用的全局ObjectMapper,两者不是同一个实例,所以该配置无法影响Spring MVC的序列化行为。


内容的提问来源于stack exchange,提问作者Samuele Domenico Ruffino

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 12:11:00