Spring Integration异步Flow实现:并行执行方法并返回指定响应
我在Spring Integration中采用Scatter-Gather模式,最终会调用saveResponseAndGenerateApplicationResponse()流程,该流程包含三个操作:
saveResponse()storeSomeData()createApplicationResponse()
当前的Spring Integration配置代码如下:
@Bean public IntegrationFlow flow() { return flow -> flow.handle(validatorService, "validateRequest") .split() .channel(c -> c.executor(Executors.newCachedThreadPool())) .scatterGather( scatterer -> scatterer .applySequence(true) .recipientFlow(saveRequestToDB()) .recipientFlow(getSomething()), gatherer -> gatherer.outputProcessor(aggregateResponse())) .to(saveResponseAndGenerateApplicationResponse()); } private IntegrationFlow saveResponseAndGenerateApplicationResponse() { return flow -> flow.enrichHeaders(h -> h.errorChannel("dbErrorChannel", true)) .handle(someService, "saveResponse") .enrichHeaders(h -> h.errorChannel("ResponseErrorChannel", true)) .handle(someService, "storeSomeData") .handle(someService, "createApplicationResponse") .logAndReply("response"); }
我的需求是让storeSomeData()与createApplicationResponse()并行执行,无需等待storeSomeData()的响应,直接将createApplicationResponse()的结果作为应用返回。Spring Integration是否有类似Java中CompletableFuture.runAsync()的原生异步特性,无需手动处理就能实现这个需求?
目前我通过以下方式实现且运行正常,但希望改用Spring Integration原生方案:
.handle( (payload, headers) -> { CompletableFuture.runAsync( () -> someService.storeSomeData( (ResponseModel) payload, headers.get("cusTx", String.class), headers.get("productType", String.class))); return someService.createApplicationResponse( (ResponseModel) payload, headers.get("httpStatusCode", Object.class)); })
Spring Integration确实提供了原生异步处理能力,完全可以替代手动使用CompletableFuture的实现,下面推荐两种适配你需求的方案:
方案一:发布订阅通道 + 异步分支
在saveResponse()执行完成后,用发布订阅通道将消息分发给两个独立分支:一个分支通过异步线程池执行storeSomeData()(不阻塞主流程),另一个分支直接执行createApplicationResponse()并返回结果。
修改后的saveResponseAndGenerateApplicationResponse()流程代码:
private IntegrationFlow saveResponseAndGenerateApplicationResponse() { return flow -> flow.enrichHeaders(h -> h.errorChannel("dbErrorChannel", true)) .handle(someService, "saveResponse") .enrichHeaders(h -> h.errorChannel("ResponseErrorChannel", true)) // 发布订阅通道拆分流程 .publishSubscribeChannel(subscribers -> subscribers // 异步执行storeSomeData,使用自定义线程池 .subscribe(subFlow -> subFlow .channel(c -> c.executor(Executors.newCachedThreadPool())) .handle(someService, "storeSomeData")) // 同步执行createApplicationResponse,作为响应返回 .subscribe(subFlow -> subFlow .handle(someService, "createApplicationResponse") .logAndReply("response"))); }
核心说明
publishSubscribeChannel会将消息发送给所有订阅的子流程,各分支独立执行- 为
storeSomeData()分支配置executor通道,使其在异步线程中运行,不会阻塞createApplicationResponse()的执行 createApplicationResponse()分支会直接处理消息并返回结果,无需等待异步分支完成
方案二:使用async()操作符简化异步逻辑
如果只需要单个异步操作,async()操作符是更简洁的选择,它内部基于TaskExecutor实现异步执行,无需手动拆分分支:
private IntegrationFlow saveResponseAndGenerateApplicationResponse() { return flow -> flow.enrichHeaders(h -> h.errorChannel("dbErrorChannel", true)) .handle(someService, "saveResponse") .enrichHeaders(h -> h.errorChannel("ResponseErrorChannel", true)) // 异步执行storeSomeData,不等待结果 .async(asyncFlow -> asyncFlow.handle(someService, "storeSomeData")) // 直接执行createApplicationResponse并返回 .handle(someService, "createApplicationResponse") .logAndReply("response"); }
核心说明
async()操作符会将后续的处理逻辑提交到默认TaskExecutor(或自定义线程池)中异步执行- 主流程会立即继续执行
createApplicationResponse(),完全不受异步分支的影响 - 若需要自定义线程池,可通过
async(asyncFlow -> asyncFlow.channel(c -> c.executor(yourCustomThreadPool)).handle(...))指定
方案优势
- 完全基于Spring Integration原生组件,代码风格更贴合框架规范
- 天然兼容框架的错误处理机制(配置的
errorChannel依然有效) - 便于后续扩展(比如新增异步分支、调整线程池参数)
内容的提问来源于stack exchange,提问作者Somnath Mukherjee

