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

如何从返回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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 14:18:10