如何判断Flux完成原因,且在上游提前终止时抛出错误?
如何判断Flux的终止原因并在特定场景抛出异常?
假设有如下示例流:
Flux.range(1, 10) .log() .takeUntil(i -> i == 5) .doOnNext(System.out::println) .blockLast();
这段代码会输出1、2、3、4、5后终止,因为takeUntil的条件i == 5被满足。
但如果修改源流使其提前终止:
Flux.range(1, 3) .log() .takeUntil(i -> i == 5) .doOnNext(System.out::println) .blockLast();
它会输出1、2、3后终止——此时takeUntil的条件并未触发,是源流自身结束导致整个流终止。
我希望在源流在takeUntil条件满足前就完成的情况下抛出错误,用Reactor实现此需求的最佳方式是什么?
之前尝试在takeUntil前添加doOnComplete的方式不可靠,还容易出错,而且可能会因为源的预取机制失效:
Flux.range(1, 3) .log() .doOnComplete(() -> { throw new RuntimeException(); }) .takeUntil(i -> i == 5) .doOnNext(System.out::println) .blockLast();
最佳实现方式
可以通过跟踪takeUntil的条件是否被触发来实现需求,推荐两种可靠的方案:
方案一:原子变量跟踪状态 + 终止时校验
用原子变量记录takeUntil的条件是否被满足,在流终止时校验状态,若未满足则抛出异常:
AtomicBoolean conditionMet = new AtomicBoolean(false); Flux.range(1, 3) .log() .takeUntil(i -> { boolean met = i == 5; if (met) { conditionMet.set(true); } return met; }) .doOnNext(System.out::println) .doFinally(signalType -> { if (!conditionMet.get() && signalType == SignalType.ON_COMPLETE) { throw new RuntimeException("源流在takeUntil条件满足前已完成"); } }) .blockLast();
这种方式直接在takeUntil的判断逻辑里标记状态,然后在doFinally中结合终止信号类型做校验,完全规避了预取机制带来的问题,逻辑清晰且可靠。
方案二:利用takeUntilOther触发错误
借助takeUntilOther的特性,当源流正常结束时主动抛出错误:
Flux.range(1, 3) .log() .takeUntilOther( // 源流完成时,发出错误信号 Mono.defer(() -> Mono.error(new RuntimeException("源流在takeUntil条件满足前已完成"))) .delaySubscription(Duration.ZERO) // 确保源流先处理完毕 ) .takeUntil(i -> i == 5) .doOnNext(System.out::println) .blockLast();
takeUntilOther会监听传入的Mono:如果takeUntil的条件先被满足,流会正常终止,错误Mono不会被触发;如果源流先完成,defer里的错误信号会被触发,直接终止流并抛出异常。
为什么之前的doOnComplete不可靠?
takeUntil会订阅源流,当源流完成时doOnComplete会被触发,但此时takeUntil已经终止了流的传递,抛出的异常可能无法被正确捕获;再加上源流的预取机制,doOnComplete的执行时机可能和预期不符,导致逻辑失效。
内容的提问来源于stack exchange,提问作者ChrisDekker
相关产品推荐
相关产品推荐

