如何在Spring Cloud Stream中对单个Kafka主题的不同数据类型消息进行反序列化?
如何在Spring Cloud Stream中对单个Kafka主题的不同数据类型消息进行反序列化?
我之前在做Spring Boot + Spring Cloud Stream对接Kafka的项目时,刚好遇到过和你一模一样的需求——单个主题里混着不同类型的JSON消息,靠header里的type字段区分,要自动反序列化成对应的Java实体类。折腾了一番后,总结出了几种可行的方案,下面分享给你:
一、先准备基础实体类
首先得把你要处理的消息实体类定义好,比如Customer和Address:
public class Customer { private String id; private String name; // 记得加上无参构造、全参构造,以及getter/setter方法 // 这里为了简洁省略,实际项目里可不能忘哦 } public class Address { private String street; private String city; // 同样,构造方法、getter/setter不能少 }
二、方案一:自定义消息转换器(优雅统一处理)
这种方式是全局配置一个转换器,所有进入的消息都会先经过它,根据type header自动选择对应的实体类反序列化,后续消费端直接拿实体类用就行,非常省心。
1. 实现自定义消息转换器
继承Spring的AbstractMessageConverter,核心逻辑就是读取header里的type,然后用Jackson反序列化:
import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.messaging.Message; import org.springframework.messaging.converter.AbstractMessageConverter; import org.springframework.messaging.converter.MessageConversionException; import org.springframework.util.MimeType; public class DynamicJsonMessageConverter extends AbstractMessageConverter { private final ObjectMapper objectMapper; private static final String TYPE_HEADER = "type"; public DynamicJsonMessageConverter(ObjectMapper objectMapper) { super(MimeType.valueOf("application/json")); this.objectMapper = objectMapper; } @Override protected boolean supports(Class<?> clazz) { // 支持我们定义的所有实体类,这里直接把Customer和Address加进去 return Customer.class.isAssignableFrom(clazz) || Address.class.isAssignableFrom(clazz); } @Override protected Object convertFromInternal(Message<?> message, Class<?> targetClass, Object conversionHint) { // 从header里拿type字段 String typeHeader = message.getHeaders().get(TYPE_HEADER, String.class); if (typeHeader == null) { throw new MessageConversionException("消息里缺少必填的'type'头信息"); } // 根据type匹配对应的实体类 Class<?> actualClass; switch (typeHeader) { case "customer": actualClass = Customer.class; break; case "address": actualClass = Address.class; break; default: throw new MessageConversionException("不支持的消息类型: " + typeHeader); } // 反序列化payload try { byte[] payload = (byte[]) message.getPayload(); return objectMapper.readValue(payload, actualClass); } catch (Exception e) { throw new MessageConversionException("反序列化消息失败", e); } } }
2. 注册自定义转换器到Spring容器
写个配置类,把上面的转换器注册进去,让Spring Cloud Stream使用它:
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.messaging.converter.MessageConverter; import org.springframework.cloud.stream.config.StreamMessageConverterCustomizer; import com.fasterxml.jackson.databind.ObjectMapper; @Configuration public class StreamConfig { @Bean public DynamicJsonMessageConverter dynamicJsonMessageConverter(ObjectMapper objectMapper) { // 用Spring容器里的ObjectMapper,避免重复创建 return new DynamicJsonMessageConverter(objectMapper); } @Bean public StreamMessageConverterCustomizer messageConverterCustomizer(DynamicJsonMessageConverter converter) { return (messageConverters) -> { // 把自定义转换器放到最前面,确保优先使用 messageConverters.add(0, converter); }; } }
三、方案二:结合@StreamListener的条件路由
配置好自定义转换器后,就可以用@StreamListener配合条件表达式,把不同类型的消息路由到对应的处理方法,方法参数直接写实体类就行:
import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.messaging.Sink; import org.springframework.stereotype.Component; @Component public class MessageConsumer { // 只处理type为customer的消息 @StreamListener(target = Sink.INPUT, condition = "headers['type']=='customer'") public void handleCustomer(Customer customer) { // 这里直接拿到反序列化后的Customer对象,做业务处理 System.out.println("收到Customer消息: " + customer); } // 只处理type为address的消息 @StreamListener(target = Sink.INPUT, condition = "headers['type']=='address'") public void handleAddress(Address address) { System.out.println("收到Address消息: " + address); } }
四、方案三:手动反序列化(简单直接)
如果你的项目比较简单,不想自定义转换器,也可以在消费端手动处理反序列化,直接拿byte[]类型的payload,结合type header自己转:
import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.messaging.Sink; import org.springframework.messaging.handler.annotation.Header; import org.springframework.stereotype.Component; @Component public class SimpleMessageConsumer { private final ObjectMapper objectMapper; // 注入Spring容器的ObjectMapper public SimpleMessageConsumer(ObjectMapper objectMapper) { this.objectMapper = objectMapper; } @StreamListener(Sink.INPUT) public void handleMessage(@Header("type") String type, byte[] payload) { try { if ("customer".equals(type)) { Customer customer = objectMapper.readValue(payload, Customer.class); handleCustomer(customer); } else if ("address".equals(type)) { Address address = objectMapper.readValue(payload, Address.class); handleAddress(address); } else { System.err.println("未知的消息类型: " + type); } } catch (Exception e) { System.err.println("反序列化消息失败: " + e.getMessage()); } } private void handleCustomer(Customer customer) { // 处理Customer业务 System.out.println("处理Customer: " + customer); } private void handleAddress(Address address) { // 处理Address业务 System.out.println("处理Address: " + address); } }
五、一些注意事项
- 生产者发送消息时,一定要记得设置
typeheader,并且把消息的content-type设为application/json,不然转换器可能不生效。 - 如果你后续要加新的消息类型,只需要新增对应的实体类,然后在自定义转换器的switch里加个case就行,扩展性很好。
- 尽量用Spring容器里的
ObjectMapper,不要自己new,这样可以复用Spring的全局配置(比如日期格式化、空值处理等)。
内容来源于stack exchange
相关产品推荐
相关产品推荐

