升级Spring Cloud Stream至2.0.0.RC3消费旧版本消息时触发转换异常
嘿,我刚好踩过Spring Cloud Stream跨版本消息转换的坑,来帮你搞定这个问题!
问题根源
你遇到的MessageConversionException,本质是Spring Cloud Stream 2.0.x(也就是你用的2.0.0.RC3,属于Elmhurst版本)和旧版Ditmars.RELEASE的默认消息序列化/反序列化机制不兼容:
- Ditmars.RELEASE默认用的是Spring Integration的
JsonMessageConverter,它序列化消息时可能不会严格给消息头加上application/json的标记,或者字节流格式和新版不匹配。 - 从2.0.x开始,Spring Cloud Stream默认换成了
Jackson2JsonMessageConverter,这家伙对消息头的contentType要求很严格——要是收到没有正确标记的字节数组消息,直接就会抛出“转不了类型”的异常。
解决方案(按推荐程度排序)
方法1:给消费端加个兼容旧版的转换器(最省心,不用改旧服务)
在你的2.0.0.RC3服务里,手动配置一个旧版的JsonMessageConverter,让消费端同时兼容新旧两种消息格式:
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.integration.converter.JsonMessageConverter; import org.springframework.messaging.converter.MessageConverter; @Configuration public class StreamConverterConfig { @Bean public MessageConverter jsonMessageConverter() { return new JsonMessageConverter(); } }
这样不管是Ditmars生产的旧格式消息,还是2.0.x服务生产的新消息,消费端都能正常解析。
方法2:修改旧生产端的配置,适配新版格式(如果能改旧服务的话)
要是你有权限修改Ditmars.RELEASE的服务配置,直接给生产的消息加上正确的contentType头就行。在旧服务的application.yml里加这段配置:
spring: cloud: stream: bindings: # 替换成你的输出绑定名称 your-output-binding-name: producer: header-mode: headers contentType: application/json
这样旧服务生产的消息就带有标准的JSON类型标记,新版消费端就能直接识别了。
方法3:自定义转换器,手动处理字节数组转换(极端情况用)
要是上面两种方法都用不了,那就自己写个转换器,硬把字节数组转成你的UsersDeletedMessage:
import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.messaging.converter.AbstractMessageConverter; import org.springframework.messaging.MessageHeaders; import org.springframework.util.MimeTypeUtils; @Configuration public class CustomMessageConverterConfig { @Bean public AbstractMessageConverter customMessageConverter(ObjectMapper objectMapper) { return new AbstractMessageConverter(MimeTypeUtils.APPLICATION_OCTET_STREAM) { @Override protected boolean supports(Class<?> clazz) { // 指定只处理UsersDeletedMessage类型的转换 return UsersDeletedMessage.class.isAssignableFrom(clazz); } @Override protected Object convertFromInternal(Object payload, MessageHeaders headers, Class<?> targetClass) { if (payload instanceof byte[]) { try { // 手动把字节数组反序列化成目标类 return objectMapper.readValue((byte[]) payload, targetClass); } catch (Exception e) { throw new RuntimeException("反序列化消息失败", e); } } return payload; } }; } }
小建议
优先试方法1,不用动旧服务,对现有系统影响最小。测试的时候可以打印一下消息头,看看旧版消息有没有contentType字段,这能帮你更快定位问题~
内容的提问来源于stack exchange,提问作者Yossi Shasha
相关产品推荐
相关产品推荐

