为何Buffer处理后的Flux使用publishOn无法并行执行耗时任务?
为什么分块后的Flux映射操作无法并行执行?
我有一个所有值同时到达的Flux,希望对这些值分块后并行执行耗时任务,但以下代码中每个映射操作都在同一个线程执行而非并行,这是为什么?
测试代码
@Test public void test() throws InterruptedException { Flux.just("a", "b", "c", "d", "e", "f", "g") .buffer(2) .publishOn(Schedulers.parallel()) .map(this::doTimeintensiveStuff) .doOnNext(val -> { log.info("[" + Thread.currentThread().getName() + "] " + val); }) .blockLast(); } private String doTimeintensiveStuff(List<String> input) { return String.join(", ", input); // 这里只是占位,实际是耗时任务 }
日志输出
Aug. 07, 2023 4:44:50 PM de.test.ParallelTest lambda$1 INFORMATION: [parallel-1] a, b Aug. 07, 2023 4:44:50 PM de.test.ParallelTest lambda$1 INFORMATION: [parallel-1] c, d Aug. 07, 2023 4:44:50 PM de.test.ParallelTest lambda$1 INFORMATION: [parallel-1] e, f Aug. 07, 2023 4:44:50 PM de.stest.ParallelTest lambda$1 INFORMATION: [parallel-1] g
原因
publishOn仅负责指定下游操作的执行线程,但不会改变Reactor序列串行处理元素的默认特性。即使切换到线程池,整个处理流程还是会按顺序把每个分块交给线程池中的某一个线程执行,不会自动将分块分配到多个线程并行处理。
解决方法
要实现分块后的并行处理,需要用flatMap配合线程池。flatMap可以为每个分块创建独立的异步任务,并支持并行执行这些任务:
修改后的示例代码:
@Test public void test() throws InterruptedException { Flux.just("a", "b", "c", "d", "e", "f", "g") .buffer(2) // flatMap实现并行处理,可指定并行度控制并发数 .flatMap(chunk -> Mono.fromCallable(() -> doTimeintensiveStuff(chunk)) .subscribeOn(Schedulers.parallel()) , 4) // 第二个参数为最大并行度,可根据需求调整 .doOnNext(val -> { log.info("[" + Thread.currentThread().getName() + "] " + val); }) .blockLast(); } private String doTimeintensiveStuff(List<String> input) { // 模拟耗时任务 try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } return String.join(", ", input); }
关键说明
flatMap会为每个分块生成独立的Mono,通过subscribeOn为每个Mono分配线程池中的线程,从而实现多任务并行。- 可以通过
flatMap的第二个参数限制最大并行数,避免因任务过多导致线程过载。
内容的提问来源于stack exchange,提问作者flxkrmr
相关产品推荐
相关产品推荐

