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

Spring WebFlux基于DB资源通过.repeatWhen()实现重复订阅问题咨询

问题根因分析

  • 第一个问题:repeatWhen 接收的函数会在整个流订阅时就立即执行,生成用于控制重复逻辑的发布者。你当前的写法直接返回了一个单次查询的Mono,没有和repeatWhen传入的重复触发信号(即参数repeat,类型是Flux<Long>,每次上游process()执行完成后会推送一个元素代表可以触发下一次重复校验)绑定,因此查询逻辑会在流初始化阶段就提前执行。
  • 第二个问题:repeatWhen的控制逻辑规则是:控制发布者发出元素则触发重复,控制发布者发出完成信号则终止重复、整体流正常结束。你当前用filter过滤后空元素场景返回空Mono,等于直接给控制发布者推送了完成信号,所以会直接终止整个流程,不符合你要的「跑完一次process()再判断是否重复」的逻辑。

最终实现代码

private Consumer<Signal<String>> processOnNewThread() {
    return signal -> {
        final var resourceId = signal.get();
        if (resourceId == null) return;

        this.process(resourceId)
            .repeatWhen(repeat -> repeat
                // 每次process执行完成后收到触发信号,再执行状态查询
                .flatMap(trigger -> Mono.defer(() -> repo.findById(resourceId))
                    // 转换为是否需要继续重复的布尔标识
                    .map(resource -> !resource.getElems().isEmpty())
                    // 资源不存在时默认不需要重复
                    .defaultIfEmpty(false)
                )
                // 只要标识为false就终止控制流,停止重复
                .takeWhile(shouldRepeat -> shouldRepeat)
            )
            .collectList()
            .contextWrite(signal.getContextView())
            .subscribeOn(Schedulers.boundedElastic())
            .subscribe();
    };
}

private Flux<String> process(String resourceId) { ... } 

逻辑说明

  • 重复触发信号和数据库查询绑定后,只有当一次process()完全执行完成后,才会触发数据库查询,不会出现查询提前执行的问题,同时每次查询都会拉取最新的资源状态。
  • 当查询到资源不存在、或者资源的元素数组为空时,会返回false标识,takeWhile收到false后会终止控制发布者,repeatWhen检测到控制流完成就会停止重复,整个流程正常结束。

内容的提问来源于stack exchange,提问作者Dmitry Kokora

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 13:45:05