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

