基于Reactor的轮询无法终止问题求助
轮询无法停止的问题排查与修复
我需要实现轮询指定次数,直到收到特定响应的逻辑,写了如下代码:
Mono.defer(() -> webClient.getResponse()) // repeat GET call after every 500 ms .repeatWhen(repeat -> repeat.delayElements(Duration.ofMillis(500))) .repeat(4) // max repeat attempts .takeUntil(response -> response.equals("1")) //.retryWhen(Retry.max(3)) .log() .subscribe(response -> { System.out.println("response received: " + response); });
但现在轮询完全停不下来,原本预期最多运行4次后退出,日志输出一直在重复:
12:54:54.527 [parallel-7] INFO reactor.Flux.TakeUntil.1 - onNext(2) response received: 2 12:54:55.034 [parallel-8] INFO reactor.Flux.TakeUntil.1 - onNext(2) response received: 2 12:54:55.540 [parallel-9] INFO reactor.Flux.TakeUntil.1 - onNext(2) response received: 2 12:54:56.045 [parallel-10] INFO reactor.Flux.TakeUntil.1 - onNext(2) response received: 2 12:54:56.551 [parallel-1] INFO reactor.Flux.TakeUntil.1 - onNext(2) response received: 2 12:54:57.057 [parallel-2] INFO reactor.Flux.TakeUntil.1 - onNext(2) response received: 2 12:54:57.563 [parallel-3] INFO reactor.Flux.TakeUntil.1 - onNext(2) response received: 2 .......
问题原因
你同时使用了repeatWhen和repeat(4),这两个操作符的逻辑叠加导致了无限循环:
repeatWhen(repeat -> repeat.delayElements(...))本身是无限重复的,只要上游流正常完成,就会延迟后再次触发调用repeat(4)会让整个repeatWhen的流再重复4次,相当于把无限循环的逻辑又复制了4份,最终导致轮询永远停不下来
修复方案
把次数限制整合到repeatWhen的逻辑里,同时保留takeUntil的终止条件,就能实现「最多轮询N次,或者提前收到目标响应就停止」的需求:
Mono.defer(() -> webClient.getResponse()) // 控制重复逻辑:延迟500ms,最多重复4次(加上初始调用共5次) .repeatWhen(repeatSignal -> repeatSignal .delayElements(Duration.ofMillis(500)) .take(4)) // 收到目标响应"1"立即停止 .takeUntil(response -> response.equals("1")) .log() .subscribe(response -> { System.out.println("response received: " + response); });
或者用Flux.interval实现更直观的轮询逻辑:
// 立即执行第一次调用,之后每500ms轮询一次 Flux.interval(Duration.ZERO, Duration.ofMillis(500)) .flatMap(ignore -> webClient.getResponse()) // 最多执行5次调用(初始+4次轮询) .take(5) // 收到目标响应提前终止 .takeUntil(response -> response.equals("1")) .log() .subscribe(response -> { System.out.println("response received: " + response); });
这两种写法都能满足需求:要么在收到"1"时立即停止,要么在完成指定次数的轮询后自动退出。
内容的提问来源于stack exchange,提问作者xploreraj
相关产品推荐
相关产品推荐

