Spring Integration 5.x/6.x中如何用Mono<?>或CompletableFuture实现ListenableFuture async=false的等效效果
我明白你的痛点——手里有一堆公共库提供的反应式服务激活器,不想把它们改成阻塞式的GenericHandler,但又要实现和ListenableFuture设置async=false一样的同步获取结果的效果,而且还要适配Spring Integration 6.x的新规范(毕竟ListenableFuture已经被弃用了)。
先说说你示例代码的问题:你设置了async=false但还是拿到Mono<"Good Morning">,是因为Spring Integration对反应式返回值的默认逻辑——这个async参数主要控制handler是否异步执行,而非自动解包Mono/Flux这类反应式类型,所以即使你关闭了异步,它也不会帮你提取Mono内部的实际payload。
下面给你两种场景下的解决方案,完全不用修改原有反应式handler:
针对Mono的解决方案(兼容5.x/6.x)
核心思路是同步解析Mono的结果,这和ListenableFutureasync=false时同步等待的行为完全一致:
方案1:在handle方法内直接同步解析
修改你的IntegrationFlow,在调用原handler后通过block()同步获取Mono的实际值:
@Bean public IntegrationFlow monoAsyncFalseFlow() { // @formatter:off return IntegrationFlows.fromSupplier(() -> "Good Morning") .handle((payload, headers) -> monoResponseHandler().handle(payload, headers).block(), ec -> ec.async(false)) .log(LoggingHandler.Level.INFO, m -> m.getPayload()) // 现在payload是字符串"Good Morning" .get(); // @formatter:on }
这里的block()会同步等待Mono完成并返回内部值,完美对应ListenableFutureasync=false的同步行为,而且完全不用改动原有的反应式handler。
方案2:封装通用包装handler(多Flow复用)
如果多个IntegrationFlow都需要这个逻辑,可以封装一个通用的包装类,避免重复代码:
@Bean public GenericHandler<String> syncMonoHandlerWrapper() { return (payload, headers) -> monoResponseHandler().handle(payload, headers).block(); } @Bean public IntegrationFlow monoAsyncFalseFlow() { // @formatter:off return IntegrationFlows.fromSupplier(() -> "Good Morning") .handle(syncMonoHandlerWrapper(), ec -> ec.async(false)) .log(LoggingHandler.Level.INFO, m -> m.getPayload()) .get(); // @formatter:on }
针对CompletableFuture的解决方案(适配6.x推荐用法)
在Spring Integration 6.x中,官方更推荐用CompletableFuture替代ListenableFuture。要实现async=false的等效效果,只需要调用join()方法同步获取结果:
假设你有返回CompletableFuture的handler:
@Bean public GenericHandler<String> completableFutureHandler() { return (p, h) -> CompletableFuture.completedFuture(p); }
对应的IntegrationFlow配置:
@Bean public IntegrationFlow completableFutureSyncFlow() { // @formatter:off return IntegrationFlows.fromSupplier(() -> "Good Morning") .handle((payload, headers) -> completableFutureHandler().handle(payload, headers).join(), ec -> ec.async(false)) .log(LoggingHandler.Level.INFO, m -> m.getPayload()) .get(); // @formatter:on }
join()和block()逻辑类似,都是同步等待结果,而且不会抛出检查型异常,更适配CompletableFuture的使用场景。
注意事项
block()/join()会阻塞当前线程,这和ListenableFutureasync=false的设计初衷完全一致——就是同步等待任务完成。如果你的系统对线程阻塞敏感,可能需要重新评估是否真的需要这种同步模式;- 如果你是全反应式架构,Spring Integration 6.x的反应式流支持已经很完善,建议尽量保持异步非阻塞的处理模式,只在必须同步获取结果的场景下使用上述方法。
备注:内容来源于stack exchange,提问作者Rayyan

