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

Java Reactor:如何非急切获取将Flux<Mono<T>>转为Flux<T>并遵守背压?

解决Reactor中flatMap/concatMap的急切获取与背压问题

我完全懂你遇到的这个坑——当你用flatMap、concatMap这类操作符时,Reactor默认的预取行为直接绕过了订阅者的初始需求,一下子拉取256条数据,这对慢订阅者来说完全违背了背压的初衷,快生产者的资源消耗问题也会立刻凸显出来。

问题根源

Reactor里的flatMap、concatMap等操作符,默认会启用预取(prefetch)策略,默认值是256。这个设计是为了提升吞吐量:提前拉取一批数据,让操作符能并行/顺序处理,减少等待时间。但在你的场景(快生产者+慢订阅者)里,这种预取反而成了灾难——生产者会一次性生成大量数据,订阅者根本处理不过来,内存占用直接飙升。

解决方案:修改预取参数

其实这些操作符都支持自定义预取数量,你只需要在调用时显式指定预取数值,就能让它严格遵守订阅者的背压需求。比如你希望和初始request(3)匹配,或者更保守地设置为1(完全跟着订阅者的节奏走)。

修改后的代码示例

把原来的flatMap(Mono::just)改成带预取参数的版本:

Flux.defer(() -> Flux.range(1, 1000))
    .doOnRequest(i -> System.out.println("Requested: " + i))
    .doOnNext(v -> System.out.println("Emitted: " + v))
    .flatMap(Mono::just, 1) // 显式指定预取为1,严格遵循背压
    .subscribe(new BaseSubscriber<Object>() {
        @Override
        protected void hookOnSubscribe(final Subscription subscription) {
            subscription.request(3);
        }
        @Override
        protected void hookOnNext(final Object value) {
            System.out.println("Received: " + value);
        }
    });

效果验证

修改后你会看到输出回到了和不用flatMap时一致的状态:

Requested: 3
Emitted: 1
Received: 1
Emitted: 2
Received: 2
Emitted: 3
Received: 3

此时flatMap会完全跟着订阅者的请求节奏走,不会提前拉取额外数据。

针对Spring WebClient的场景

如果你是用Spring WebClient作为快生产者,比如从接口拉取数据流,只需要在后续的flatMap中指定预取参数即可:

webClient.get()
         .uri("/high-speed-data-stream")
         .retrieve()
         .bodyToFlux(YourDataClass.class)
         .flatMap(data -> processSlowly(data), 1) // 控制预取,适配慢订阅者
         .subscribe(new YourSlowSubscriber());

额外提示

  • 如果你需要严格的顺序处理,concatMap同样支持预取参数,concatMap(Mono::just, 1)也能达到相同的背压效果,因为它本身就是逐个处理元素的。
  • 预取数值的选择:如果订阅者处理速度稳定,你可以设置成和初始request值匹配的数字;如果不确定,设置为1是最保守且安全的选择,完全不会破坏背压链路。

内容的提问来源于stack exchange,提问作者Mario Zsilak

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 16:12:33