Spring Integration Scatter Gather并行流Aggregator超时配置不生效问题求解
解决方案
1. 基础超时配置(从分发请求开始计时)
你之前配置的groupTimeout是聚合器的分组超时,计时起点为第一个结果消息到达聚合器的时间,不符合你从分发请求就开始计时的需求,直接使用scatterGather的第三个参数gatherTimeout即可实现分发即计时,同时需开启sendPartialResultOnExpiry保证超时后返回已收集的结果,完整配置如下:
@Bean public IntegrationFlow queueFlow(LogicService service) { return f -> f.scatterGather(scatterer -> scatterer .applySequence(true) .recipientFlow(aFlow(service)) .recipientFlow(bFlow(service)) , aggregatorSpec -> aggregatorSpec // 开启超时后返回部分已收集结果 .sendPartialResultOnExpiry(true) // 核心参数:分发完成后开始计时,2秒超时 , 2000L) }
该配置会在分发完所有子流请求后立刻启动计时,到2秒后直接返回当前已收集到的所有结果,不需要等待所有子流执行完成。
2. 自定义超时逻辑优化
如果你需要基于请求头的自定义超时时间控制,可以搭配短间隔轮询触发释放策略检查,解决只有新消息到达才执行释放策略的问题:
@Bean public IntegrationFlow queueFlow(LogicService service) { return f -> f.scatterGather(scatterer -> scatterer .applySequence(true) .recipientFlow(aFlow(service)) .recipientFlow(bFlow(service)) , aggregatorSpec -> aggregatorSpec .groupConditionProvider(group -> { // 读取你在网关层添加的自定义请求头 Long startTime = group.getOne().getHeaders().get("requestStartTime", Long.class); Long customTimeout = group.getOne().getHeaders().get("customTimeout", Long.class); return System.currentTimeMillis() - startTime >= customTimeout; }) .releaseStrategy(group -> Boolean.TRUE.equals(group.getCondition())) // 每100ms触发一次释放策略检查,无需等待新消息到达 .groupTimeout(100L) .sendPartialResultOnExpiry(true) ) }
注意事项
- 你当前使用
Executors.newCachedThreadPool作为子流入场通道的线程池是正确的,需确保子流都是异步执行,不会阻塞scatterer的分发逻辑,否则超时计时起点会延迟 - 超时后仍在执行的子流结果会被自动丢弃,不需要额外配置,如果你需要主动中断子流执行,可以在子流的业务逻辑中添加超时中断判断
内容的提问来源于stack exchange,提问作者jandresboyaca
相关产品推荐
相关产品推荐

