Spring Integration Supplier<?>类型转换异常:消息转为byte[]而非Order对象
升级Spring Boot/Spring Integration后消息类型转换异常排查与解决
将Java 8+Spring Boot 1.8+spring-integration-java-dsl 1.2.3应用升级为Java 21+Spring Boot 3.3.5+spring-integration-core 6.3.4后,发送到RabbitMQ队列的消息消费时变为byte[]而非OrderMessage对象,触发类型转换异常。旧版本可直接获取OrderMessage对象,怀疑是配置缺失或Supplier<?>实现错误。
旧版配置与代码
application.yml (old)
spring: rabbitmq: host: localhost #127.0.0.1 port: 5672 username: guest password: guest cloud: fail-fast: false stream: bindings: order-delay-output-channel: destination: orderDelayExchange producer: requiredGroups: orderDelayQueue rabbit: bindings: order-delay-output-channel: producer: auto-bind-dlq: true ttl: 30000
Output定义
public interface OrderProcessor { @Output( "order-delay-output-channel") MessageChannel orderDelayOutputChannel(); }
IntegrationFlow
@Configuration @EnableBinding( OrderProcessor.class ) public class IntegrationConfig { @Bean public IntegrationFlow orderDelayFlow() { return IntegrationFlows.from(ORDER_DELAY_CHANNEL) .log(Level.TRACE, this.getClass().getName() + ".orderDelayFlow") .enrichHeaders(h -> h.header("source", "orderDelayFlow")) .<OrderMessage>handle((orderMessage, headers) -> { final Channel channel = (Channel) headers.get(AmqpHeaders.CHANNEL); final Long deliveryTag = (Long) headers.get(AmqpHeaders.DELIVERY_TAG); try { channel.basicAck(deliveryTag, false); } catch (final IOException exception) { throw new MessagingException("ack failed", exception); } return orderMessage; }) //旧版本发送到order-delay-output-channel后,可在其他IntegrationFlow中获取OrderMessage对象 .channel("order-delay-output-channel") .get(); } }
新版配置与代码
application.yml (new)
spring: cloud: function: definition: orderDelayOutputChannel; stream: defaultBinder: rabbit bindings: orderDelayOutputChannel-out-0: #order-delay-output-channel: destination: orderDelayExchange producer: requiredGroups: orderDelayQueue rabbit: bindings: orderDelayOutputChannel-out-0: producer: auto-bind-dlq: true ttl: 1000
IntegrationFlow
@Configuration @EnableBinding( OrderProcessor.class ) public class IntegrationConfig { @Bean public IntegrationFlow orderDelayFlow() { return IntegrationFlows.from(ORDER_DELAY_CHANNEL) .log(Level.TRACE, this.getClass().getName() + ".orderDelayFlow") .enrichHeaders(h -> h.header("source", "orderDelayFlow")) .<OrderMessage>handle((orderMessage, headers) -> { final Channel channel = (Channel) headers.get(AmqpHeaders.CHANNEL); final Long deliveryTag = (Long) headers.get(AmqpHeaders.DELIVERY_TAG); try { channel.basicAck(deliveryTag, false); } catch (final IOException exception) { throw new MessagingException("ack failed", exception); } return orderMessage; }) //新版本发送到order-delay-output-channel后,消费时得到byte[],触发转换异常! .channel("order-delay-output-channel") .get(); } }
Output定义
@Configuration public class OrderProcess { @Bean("orderDelayOutputChannel") public Supplier<Message<OrderMessage>> orderDelayOutputChannel() { return () -> null; } }
其他尝试实现(基于Consumer建议)
public class OrderMessageProcessor { //使用此实现发送消息时出现"No Dispatcher"错误 @MessagingGateway public interface orderDelayOutputChannel extends Supplier<Message<OrderSend>> {} }
原因分析与解决办法
核心原因
- 新旧模型混用:Spring Cloud Stream 3.x+ 已弃用
@EnableBinding旧模型,混用后消息序列化逻辑不兼容,导致对象未被正确序列化直接以字节形式发送。 - Supplier实现错误:返回
null的Supplier无法正确绑定消息通道,框架无法识别需要序列化的对象类型。 - 配置不匹配:函数式通道绑定配置与实际通道引用不一致,未指定消息内容类型,导致默认以字节传输。
解决步骤
1. 移除旧版@EnableBinding,规范函数式实现
删除@EnableBinding(OrderProcessor.class)及旧OrderProcessor接口,重新定义生产者函数:
@Configuration public class OrderProcess { // 推荐:直接返回对象流,由框架自动处理消息封装与序列化 @Bean public Supplier<Flux<OrderMessage>> orderDelayOutputChannel() { return Flux::empty; } }
2. 调整IntegrationFlow,绑定正确的函数式通道
将IntegrationFlow中的输出通道改为函数式模型的标准通道名{函数名}-out-0:
@Configuration public class IntegrationConfig { @Bean public IntegrationFlow orderDelayFlow() { return IntegrationFlows.from(ORDER_DELAY_CHANNEL) .log(Level.TRACE, this.getClass().getName() + ".orderDelayFlow") .enrichHeaders(h -> h.header("source", "orderDelayFlow")) .<OrderMessage>handle((orderMessage, headers) -> { final Channel channel = (Channel) headers.get(AmqpHeaders.CHANNEL); final Long deliveryTag = (Long) headers.get(AmqpHeaders.DELIVERY_TAG); try { channel.basicAck(deliveryTag, false); } catch (final IOException exception) { throw new MessagingException("ack failed", exception); } return orderMessage; }) // 绑定到函数式输出通道 .channel("orderDelayOutputChannel-out-0") .get(); } }
3. 修正application.yml配置
移除冗余配置,指定全局消息内容类型为JSON,确保序列化/反序列化一致:
spring: rabbitmq: host: localhost port: 5672 username: guest password: guest cloud: function: definition: orderDelayOutputChannel stream: defaultBinder: rabbit # 全局指定消息内容类型,强制JSON序列化 default.content-type: application/json bindings: orderDelayOutputChannel-out-0: destination: orderDelayExchange producer: requiredGroups: orderDelayQueue rabbit: bindings: orderDelayOutputChannel-out-0: producer: auto-bind-dlq: true ttl: 1000
4. 确保OrderMessage可序列化
OrderMessage需实现Serializable接口,或确保Jackson能正确序列化该类:
public class OrderMessage implements Serializable { // 类字段、getter/setter、构造方法等 }
5. 正确使用@MessagingGateway(可选)
若需使用网关发送消息,需绑定到正确的函数式通道:
@MessagingGateway public interface OrderDelayGateway { @Gateway(requestChannel = "orderDelayOutputChannel-out-0") void sendOrderDelayMessage(OrderMessage message); }
内容的提问来源于stack exchange,提问作者raVen
相关产品推荐
相关产品推荐

