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

Spring WebFlux中如何用响应式算子替代while循环的阻塞逻辑?

重构响应式代码消除阻塞与额外线程

原代码通过额外线程+block()实现循环查询,既违反了响应式非阻塞原则,又容易引发线程池耗尽问题。下面用Reactor原生算子重构,完全实现非阻塞的循环查询逻辑,同时满足你的三个终止条件:

重构后的代码

首先定义一个辅助类封装单次查询结果(也可以用Reactor的Tuple2替代):

// 辅助类:封装单次数据库查询的判断结果
private static class CheckResult {
    boolean shouldTerminate;
    long fechaBd;

    CheckResult(boolean shouldTerminate, long fechaBd) {
        this.shouldTerminate = shouldTerminate;
        this.fechaBd = fechaBd;
    }
}

然后改造httpush方法:

@PostMapping("/httpush")
public Mono<ResponseEntity<HttpushResponseUtil>> httpush(@RequestParam("timestamp") String timestamp, @RequestParam("token") String token) {
    // 转换传入的timestamp为秒级时间戳
    long fechaActual = "null".equals(timestamp) ? 0 : Long.parseLong(timestamp);
    // 循环计数(响应式环境用原子类保证线程安全)
    AtomicInteger loopCount = new AtomicInteger(0);

    // 单次数据库查询+终止判断逻辑
    Mono<CheckResult> singleCheck = service.ObtenerUpdateTimePorSerial(token, "update_time")
            .map(updateTimeList -> {
                if (updateTimeList.isEmpty()) {
                    return new CheckResult(false, 0);
                }

                HuellaTemp huellaTemp = updateTimeList.get(0);
                // 条件1:状态为"Muestras Restantes: 0"时终止
                if ("Muestras Restantes: 0".equals(huellaTemp.getStatusPlantilla())) {
                    return new CheckResult(true, 0);
                }

                // 计算数据库中的秒级时间戳
                long fechaBd = huellaTemp.getUpdate_time() != null 
                        ? huellaTemp.getUpdate_time().getTime() / 1000 
                        : 0;
                // 条件2:数据库时间大于传入的timestamp时终止
                boolean timeConditionMet = fechaBd > fechaActual;
                return new CheckResult(timeConditionMet, fechaBd);
            })
            // 每次查询后延迟1秒,避免频繁压库(对应原代码elapsedTime递增逻辑)
            .delayElement(Duration.ofSeconds(1));

    // 循环执行查询,直到满足终止条件
    return singleCheck
            .repeatUntil(result -> {
                int currentCount = loopCount.incrementAndGet();
                // 三个终止条件满足任意一个就停止循环
                return result.shouldTerminate || currentCount >= 600;
            })
            // 循环结束后执行最终的指纹查询并返回响应
            .then(service.ObtenerHuellaHttpush())
            .map(httpush -> ResponseEntity
                    .ok()
                    .contentType(MediaType.APPLICATION_JSON)
                    .body(httpush));
}

关键改进说明

  • 移除了手动创建的线程和block()调用,所有操作都基于Reactor异步流处理,完全非阻塞
  • 用repeatUntil替代命令式while循环,天然适配响应式编程模型
  • 通过delayElement控制查询频率,避免短时间内重复访问数据库
  • 用AtomicInteger管理循环次数,保证响应式环境下的线程安全
  • 用then算子保证循环逻辑完成后才执行最终的查询,严格控制执行顺序

内容的提问来源于stack exchange,提问作者Jose Ricardo Castañeda

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 07:05:24