Spring Integration scatter-gather并行流前如何调用校验方法
Spring Integration scatter-gather模式前置校验实现方案
校验节点放置位置
校验逻辑必须放在主流程最前端、所有异步线程切换、split、scatterGather分发逻辑之前,也就是网关传入请求后第一个执行的节点,绝对不能放在split()或异步executor通道之后:
- 异步通道切换后,逻辑会在子线程池执行,抛出的异常无法直接阻断主线程的流程调度,容易出现校验异常抛出时,部分并行子流程已经被触发执行的问题
- 该位置的请求payload还是网关传入的原始
LionRequest类型,不需要额外类型转换就能直接传入校验方法,完全匹配校验方法的入参要求 - 该位置的逻辑是同步阻塞执行的,校验抛出异常时流程会直接终止,不会执行任何后续逻辑,异常会直接透传给网关调用方,不需要额外做异常传播配置。
具体实现方式
直接在主flow定义的最开头新增handle节点,调用你的校验方法即可,修改后的主流程代码如下:
@Bean public IntegrationFlow flow() { return flow -> // 新增校验节点作为流程第一个执行逻辑 .<LionRequest>handle((request, headers) -> { // 调用定义好的无返回值校验方法,校验失败会直接抛出异常 lionsService.validateLionRequest(request); // 校验通过返回原请求,继续执行后续流程 return request; }) .log() .split() .channel(c -> c.executor(Executors.newCachedThreadPool())) .convert(LoanProvisionRequest.class) .scatterGather( scatterer -> scatterer .applySequence(true) .recipientFlow(flow1()) .recipientFlow(flow2()) .recipientFlow(flow3()), gatherer -> gatherer.releaseLockBeforeSend(true)) .log() .aggregate(a -> a.outputProcessor(MessageGroup::getMessages)) .channel("output-flow"); }
其他可选实现与注意事项
- 如果要拆分复用校验逻辑,可以把校验逻辑单独定义为一个IntegrationFlow,在主流程开头用
.gateway(validateFlow())的方式调用,效果完全一致:校验子流程抛出异常时主流程会直接中断,异常透传给调用方。 - 不要尝试在scatterGather的scatterer前置拦截器里做校验,该拦截器执行时已经完成了分发前的上下文初始化,异常阻断效果不如直接在流程最前端加handle节点稳定。
- 额外提醒:当前配置类里的
long dbId = new SequenceGenerator().nextId();是配置类的成员变量,会在Spring容器启动时只初始化一次,所有请求都会共用同一个dbId,如果需要每个请求生成唯一id,要把id生成逻辑放到对应子流程的handle方法里,按请求维度生成。
内容的提问来源于stack exchange,提问作者Somnath Mukherjee
相关产品推荐
相关产品推荐

