Project Reactor limitRate操作符异常行为咨询
问题分析与解决
为什么会出现初始两次request(1)?
你遇到的初始连续两次请求,是Reactor中limitRate操作符的水位补充机制导致的:
- 当你使用
limitRate(1)时,默认的低水位(lowTide)是高水位(highTide=1)的一半向下取整,也就是0。 - 订阅触发后,
limitRate首先向上游请求1个元素,Flux.generate立刻返回onNext(1)。 - 此时
limitRate的剩余请求配额变为0,刚好等于低水位阈值,触发补充请求逻辑,再次向上游请求1个元素,于是产生了第二次request(1),紧接着Flux.generate返回onNext(2)。 - 由于
concatMap处理的是异步Mono(你用了delayElement做延迟),第一个元素的处理还未完成,第二个元素会被暂存在concatMap的内部队列中,直到第一个元素处理完毕才会开始处理。
这种行为会导致你的业务场景中提前获取下一条数据库记录,不符合“处理完一条再获取下一条”的预期。
解决方法:严格控制逐个请求与处理
要实现严格的串行处理逻辑,只需调整concatMap的预取参数,或者配合limitRate的精确配置即可:
方案1:设置concatMap的预取数量为1
修改concatMap的调用,使用重载方法指定预取数为1,限制它仅从上游请求1个元素,处理完成后再请求下一个:
.concatMap(integer -> Mono.just(integer * 2).delayElement(Duration.ofSeconds(2)), 1)
方案2:配合limitRate的精确参数
如果你需要更严格的流量控制,也可以同时显式设置limitRate的低水位为1,确保只有当元素被处理后才会补充请求:
.limitRate(1, 1) .concatMap(integer -> Mono.just(integer * 2).delayElement(Duration.ofSeconds(2)))
修改后的预期日志
调整后,日志会变为严格的“请求-处理-再请求”顺序:
10:xx:xx [main] INFO reactor.Flux.Generate.1 - | request(1) 10:xx:xx [main] INFO reactor.Flux.Generate.1 - | onNext(1) 10:xx:xx+2s [parallel-1] DEBUG com.DummyTest - x2: 2 10:xx:xx+2s [parallel-1] INFO reactor.Flux.Generate.1 - | request(1) 10:xx:xx+2s [parallel-1] INFO reactor.Flux.Generate.1 - | onNext(2) 10:xx:xx+4s [parallel-2] DEBUG com.DummyTest - x2: 4 ...
内容的提问来源于stack exchange,提问作者Daniel Yordanov
相关产品推荐
相关产品推荐

