如何从返回Mono的函数正确生成Flux?(响应式迭代场景)
响应式迭代生成Flux:规避block与链式flatMap问题
问题背景
我需要编写响应式迭代数据库的代码(下一次查询依赖上一次查询结果),现将问题简化为:基于以下模拟耗时操作的函数,生成包含其多次调用结果的Flux:
// 模拟耗时调用 private Mono<Integer> calculateNext(Integer value) { return Mono.defer(() -> Mono.just(value + 1)) .delayElement(Duration.ofSeconds(1L)); }
我尝试了两种实现,但都存在问题:
方案一:含block()的实现(可行但违背响应式设计)
这段代码能生成1到1000的递增整数,但使用block()阻塞了异步流,违背响应式设计初衷,还可能引发性能问题:
Flux.generate( () -> Mono.just(0), (previousResult, sink) -> { var nextResult = calculateNext(previousResult.block()); sink.next(nextResult); return nextResult; } ).flatMap(it -> (Mono<Integer>) it) .takeUntil(it -> it > 1000) .doOnNext(it -> LOG.info("it = " + it)) .subscribeOn(Schedulers.boundedElastic()) .blockLast();
方案二:链式flatMap的实现(可行但存在风险)
该实现也能运行,但相当于链式调用上千次flatMap,可能带来性能损耗或栈溢出风险:
(previousResult, sink) -> { var nextResult = previousResult.flatMap(it -> calculateNext(it)); sink.next(nextResult); return nextResult; }
推荐解决方案:使用Mono.expand递归生成Flux
Mono.expand是专门处理递归依赖场景的响应式操作符,能完美适配"下一次调用依赖上一次结果"的需求,既无阻塞操作,也不会产生深层嵌套的flatMap链:
Mono.just(0) .expand(value -> calculateNext(value)) .takeUntil(it -> it > 1000) .doOnNext(it -> LOG.info("it = " + it)) .subscribeOn(Schedulers.boundedElastic()) .blockLast();
实现说明
Mono.just(0)作为初始值启动迭代流程expand会对每个元素调用calculateNext,将返回的Mono继续展开,形成连续的异步数据流takeUntil(it -> it > 1000)控制迭代终止条件,当元素超过1000时停止- 全程保持响应式异步特性,完全符合Reactor设计规范
备选方案:改进Flux.generate的异步处理
若坚持使用Flux.generate,可通过异步状态维护避免block(),但代码复杂度更高:
Flux.<Integer>generate(() -> { AtomicReference<Integer> current = new AtomicReference<>(0); return sink -> { int currentValue = current.get(); if (currentValue > 1000) { sink.complete(); return; } calculateNext(currentValue) .subscribe(next -> { sink.next(next); current.set(next); }, sink::error); }; }) .doOnNext(it -> LOG.info("it = " + it)) .subscribeOn(Schedulers.boundedElastic()) .blockLast();
内容的提问来源于stack exchange,提问作者Alpharius
相关产品推荐
相关产品推荐

