Reactor提取器的性能阻塞机制探究:Mono.toFuture()与Flux.collectList()的阻塞疑问及代码分析
collectList() and toFuture() Great question—let's unpack these two common points of confusion with Reactor, using your code example as context.
First: How does Flux.collectList() work, and does it "block" a thread?
Let's start by clarifying a key distinction: collectList() itself is not a blocking operation (unless you pair it with a blocking call like Future.get() later). Instead, it's a terminal operator that triggers the entire reactive stream, and waits asynchronously for all elements to be emitted before emitting a single List downstream.
Here's the breakdown of its thread behavior:
- It does NOT tie up a thread to sit and wait for all elements. Reactor uses asynchronous callbacks under the hood. When the last element of the upstream
Fluxis processed, the thread that handled that final element (or a Reactor-internal thread, depending on the scheduler context) will assemble the full list and trigger the downstream operators (like yourdoOnNext(o -> System.out.println(...before map))andmap). - In your code, after
Flux.merge(m)finishes emitting all elements,collectList()will kick in on the same thread that emitted the last element frommerge. Then yourmapruns, andpublishOn(Schedulers.single())switches the thread to the single scheduler for the finaltoFuture()step.
The reason Reactor discourages collectList() isn't about thread blocking—it's because it forces all elements to be held in memory at once (bad for large datasets) and breaks the reactive stream's backpressure and asynchronous flow, turning a streaming process into a batch process.
Second: Does Mono.toFuture() block until the Mono completes?
No—Mono.toFuture() returns a CompletableFuture immediately, even if the Mono hasn't emitted a value yet. This Future starts in an incomplete state:
- When the
MonoemitsonNext()(with a value) oronComplete()(for an emptyMono), theFutureis marked as completed successfully. - If the
MonoemitsonError(), theFuturecompletes exceptionally.
In your code example:
- The inner
Mono.just(i).map(...).toFuture()returns aCompletableFutureright away.Mono.fromFuture()then wraps thisFutureback into a reactive type, so it integrates with the rest of the stream asynchronously. - The final
toFuture()at the end of your pipeline also returns immediately. ThisFuturewill only complete once the entire upstream pipeline is done: all elements processed byparallel()/runOn(), merged, collected into a list, mapped, and passed throughpublishOn(Schedulers.single()).
If you were to call mm.get() right after assigning it, that would block the current thread until the Future completes—but the toFuture() call itself is non-blocking.
Why Reactor Discourages These Operators
To circle back to your initial note: Reactor warns against collectList() and toFuture() because they push you away from the reactive, non-blocking paradigm. Instead of letting the stream process elements asynchronously with backpressure, these operators force you to wait for a single, aggregated result. For large datasets, collectList() can cause memory issues, and toFuture() often leads developers to use blocking get() calls that negate the benefits of reactive programming.
内容的提问来源于stack exchange,提问作者bbnoraining

