为何Flux.generate()上游的doOnRequest从未被触发?
问题描述
以下是一段Reactor测试代码:
@Test void testTooMuchBuffering2() { var counter = new AtomicInteger(0); var lastCollected = new AtomicInteger(-1); Flux.<List<Integer>>generate( sink -> { int val = counter.getAndAdd(1000); if (val >= 10_000) { sink.complete(); return; } // 从“数据库”读取1000行数据 System.out.println("Generating " + (val + 1000)); sink.next(IntStream.range(val, val + 1000).boxed().toList()); }) .doOnSubscribe(s -> System.out.println("Subscribed")) .doOnRequest(r -> System.out.println("flatMapIterable: " + r)) // 从未执行 .flatMapIterable(Function.identity()) .doOnRequest(r -> System.out.println("blockLast: " + r)) .blockLast(); }
测试输出如下:
Subscribed blockLast: 9223372036854775807 Generating 1000 Generating 2000 Generating 3000 Generating 4000 Generating 5000 Generating 6000 Generating 7000 Generating 8000 Generating 9000 Generating 10000
其中flatMapIterable上方的doOnRequest中的Lambda表达式从未被执行。为何会出现该情况?在其他场景中存在请求过多的问题,但此处却无任何请求触发,原因是什么?
问题原因解析
1. Flux.generate的无背压特性是核心根源
Flux.generate是同步且不遵循背压机制的生成器:一旦被订阅,它就会持续调用生成逻辑,直到sink.complete()触发,全程无视下游是否发送请求信号。
在你的代码里,generate会直接生成所有10批数据(直到val >=10000),完全不需要等待下游的请求——这就导致flatMapIterable根本没必要向上游(generate的Flux)发送请求,自然不会触发上方的doOnRequest。
2. doOnRequest的监听逻辑
doOnRequest仅监听当前操作符的下游订阅者向它发起的请求:
- 最下游的
blockLast()会向上游发送一个最大值请求(Long.MAX_VALUE,也就是输出里的9223372036854775807),这个请求会被flatMapIterable下方的doOnRequest捕获,所以你能看到对应的输出。 - 但
flatMapIterable的上游(generate)是主动推送数据的,不需要它发请求就已经把所有数据推过来了,所以flatMapIterable上方的doOnRequest永远不会执行。
3. 和其他“请求过多”场景的区别
其他场景里的请求过多,通常是因为上游是支持背压的数据源(比如数据库查询、消息队列),它们会根据下游的请求量来生成数据。但generate是主动推送的特例,它不遵守背压规则,不需要下游请求就会输出所有数据,自然不存在“下游向上游发请求”的过程,也就不会触发对应的doOnRequest。
如果想让generate遵守背压,可以改用Flux.create并手动处理背压逻辑,或者在generate后添加onBackpressureBuffer之类的背压处理操作符。
内容的提问来源于stack exchange,提问作者RonVe
相关产品推荐
相关产品推荐

