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

升级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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 11:42:45