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
相关产品推荐
相关产品推荐

