为何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在序列化响应时抛出异常。
解决方式(三选一即可)
- 返回自定义响应对象:只封装业务需要的发送结果信息,避免返回包含
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() )); }
- 不返回发送结果:如果前端不需要发送结果,直接返回
CompletableFuture<Void>
@PostMapping("/send/transitECS") public CompletableFuture<Void> sendTransit(@Valid @RequestBody TransitECS transitECS) { return this.transitECSKafkaTemplate.sendDefault(transitECS) .thenAccept(sendResult -> {}); }
- 配置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
相关产品推荐
相关产品推荐

