Spring Integration Async RabbitMQ通信返回ListenableFuture而非结果问题
问题根因
你遇到的问题是Service2的Amqp入站网关默认未开启异步返回值支持,拿到ListenableFuture类型的返回值后直接作为消息payload序列化返回,不会等待Future执行完成后提取结果。你之前配置的@ServiceActivator(async = "true")是已废弃的旧版本参数,仅控制方法是否异步执行,不会触发返回值解包逻辑。
解决方案
1. 修正Service2的入站网关配置
给Amqp.inboundGateway添加async(true)配置,开启异步返回值自动解包能力,框架会自动等待ListenableFuture执行完成后,再把实际结果返回给Service1。同时注意你原有代码存在括号未闭合、未配置消息转换器的问题,修正后配置如下:
@Bean public IntegrationFlow requestNewSessionFlow(ConnectionFactory connectionFactory, SessionProperties sessionProperties, MessageConverter messageConverter, RequestNewSessionHandler requestNewSessionHandler) { return IntegrationFlows.from(Amqp.inboundGateway(connectionFactory, sessionProperties.requestSessionProperties().queueName()) // 开启网关异步支持,自动解包ListenableFuture的执行结果 .async(true) // 绑定消息转换器,保证DTO序列化/反序列化正常 .messageConverter(messageConverter)) .handle(requestNewSessionHandler) .get(); }
2. 调整Service2的业务处理方法
去掉@ServiceActivator上已废弃的async参数,只要返回ListenableFuture即可被框架自动识别处理:
@ServiceActivator public ListenableFuture<SessionCreationResponseDTO> handleRequestNewSession() { SettableListenableFuture<SessionCreationResponseDTO> settableListenableFuture = new SettableListenableFuture<>(); // 异步会话创建逻辑执行完成后,调用set方法写入成功响应 // 若执行异常,调用setException方法写入异常信息 return settableListenableFuture; }
补充说明
- 该配置要求Spring Integration版本不低于4.3,低于该版本需要手动通过
Future.get()阻塞获取结果后再返回 - 框架同时支持返回
CompletableFuture类型,同样会自动解包执行结果 - 如果需要自定义异常处理逻辑,可以给入站网关配置
errorChannel,统一处理异步执行抛出的异常
内容的提问来源于stack exchange,提问作者Boris
相关产品推荐
相关产品推荐

