You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 08:53:29