使用Mono.fromSupplier处理长耗时操作时map与flatMap的对比
问题解答
你的写法有没有优化效果?
没有(除非补充调度器配置)。原因是:
Mono.fromSupplier()默认会在当前订阅线程(比如WebFlux的IO线程)执行computeStuff,和直接用map一样会阻塞线程,根本没达到异步解耦的目的。- 你只把同步操作包装成了Mono,但没指定它在独立的线程池执行,本质还是同步阻塞。
正确的写法
要让耗时操作不占用WebFlux的IO线程,需要给Mono指定专门的调度器,推荐用Schedulers.boundedElastic()(专门处理阻塞/耗时同步任务的弹性线程池):
fluxOfSomething.flatMap(value -> Mono.fromSupplier(() -> computeStuff(value)) .subscribeOn(Schedulers.boundedElastic()) // 关键:指定耗时任务的执行线程池 )
如果computeStuff本身是异步方法(返回Mono),那直接返回它就行,不用额外包装:
// 假设computeStuffAsync返回Mono<Result> fluxOfSomething.flatMap(this::computeStuffAsync)
怎么验证优化效果?
可以通过两个简单测试来验证:
1. 打印线程名,看执行线程是否隔离
在computeStuff里加一行打印:
private Result computeStuff(Value value) { System.out.println("computeStuff running on thread: " + Thread.currentThread().getName()); // 模拟耗时操作 try { Thread.sleep(1000); } catch (InterruptedException e) {} return new Result(value); }
- 用原
map写法:打印的线程名会是reactor-http-nio-*(WebFlux的IO线程),会阻塞这些线程导致后续请求排队。 - 用加了
subscribeOn的flatMap写法:打印的线程名会是boundedElastic-*,IO线程不会被阻塞。
2. 统计批量执行的总耗时
用StepVerifier测试批量处理的时间:
@Test void testMapVsFlatMap() { // 创建包含10个元素的Flux Flux<Integer> flux = Flux.range(1, 10); // 测试map写法的耗时 long mapStart = System.currentTimeMillis(); StepVerifier.create(flux.map(this::computeStuff)) .expectNextCount(10) .verifyComplete(); long mapEnd = System.currentTimeMillis(); System.out.println("Map total time: " + (mapEnd - mapStart) + "ms"); // 大概10*1000=10000ms左右,同步串行 // 测试优化后的flatMap写法的耗时 long flatMapStart = System.currentTimeMillis(); StepVerifier.create(flux.flatMap(value -> Mono.fromSupplier(() -> computeStuff(value)) .subscribeOn(Schedulers.boundedElastic()) )) .expectNextCount(10) .verifyComplete(); long flatMapEnd = System.currentTimeMillis(); System.out.println("FlatMap total time: " + (flatMapEnd - flatMapStart) + "ms"); // 大概1000ms左右,并行执行(受线程池大小限制) }
优化后的flatMap写法总耗时会远低于map写法,因为耗时任务是并行在弹性线程池执行的。
内容的提问来源于stack exchange,提问作者C-Otto
相关产品推荐
相关产品推荐

