如何为Spring KafkaTemplate指定自定义ValueSerializer实例?
解决Spring KafkaTemplate配置自定义JsonSerializer实例的问题
在spring-kafka 2.1.0.RELEASE版本中,你可以通过直接向DefaultKafkaProducerFactory传入序列化器实例的方式,来配置带自定义ObjectMapper的JsonSerializer,而不是仅通过配置属性指定序列化器类。
具体实现步骤:
- 先创建你的自定义
JsonSerializer实例,传入配置好的customObjectMapper(比如设置了PropertyNamingStrategy.SNAKE_CASE) - 使用
DefaultKafkaProducerFactory的重载构造器,将配置属性、key序列化器、自定义value序列化器一起传入 - 基于这个ProducerFactory创建KafkaTemplate即可
修改后的代码示例:
@Bean public ProducerFactory<String, Object> producerFactory(KafkaProperties kafkaProperties, ObjectMapper customObjectMapper) { // 创建自定义的JsonSerializer实例 JsonSerializer<Object> valueSerializer = new JsonSerializer<>(customObjectMapper); valueSerializer.setAddTypeInfoHeaders(false); // 对应之前的JsonSerializer.ADD_TYPE_INFO_HEADERS配置 Map<String, Object> props = new HashMap<>(kafkaProperties.getDefaultSettings()); // 这里不再通过props指定VALUE_SERIALIZER_CLASS_CONFIG,而是直接传入实例 props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); // 使用带序列化器实例的构造器创建ProducerFactory return new DefaultKafkaProducerFactory<>(props, Serdes.String().serializer(), valueSerializer); } @Bean public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> producerFactory) { KafkaTemplate<String, Object> kafkaTemplate = new KafkaTemplate<>(producerFactory); kafkaTemplate.setDefaultTopic("topicName"); return kafkaTemplate; }
为什么之前的方法不生效?
- 通过
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class)配置时,Spring会通过反射创建JsonSerializer实例,无法传入你自定义的ObjectMapper,只能使用默认构造器生成的实例 kafkaTemplate.setMessageConverter(...)确实仅在调用send(Message<?>)方法(传入Spring的Message对象)时生效,直接调用send(key, value)这类方法会绕过MessageConverter,直接使用ProducerFactory中的序列化器
内容的提问来源于stack exchange,提问作者Vasyl Sarzhynskyi
相关产品推荐
相关产品推荐

