Reactor Flux阻塞疑问:提交请求需等响应才继续及线程复用问题
问题描述
我创建了一个订阅时会触发API调用的Flux,该API需要数秒才能返回响应。同时我把流切换到了一个包含2个线程的调度器:
private Scheduler sc = Schedulers.newBoundedElastic(2, Integer.MAX_VALUE, "newwww-sch"); Flux.just(i) .flatMap(x -> { try { System.out.println(Thread.currentThread().getName() + " Do Request " + i); return Mono.just(spectrumAvailabilityApi.getSpectrumAvailability(spectrumAvailabilityRequest)); } catch (ApiException e) { throw new RuntimeException(e); } }) .subscribeOn(sc) .subscribe(s -> System.out.println(Thread.currentThread().getName() + " Got Response " + i));
当我用并行方式触发这个流程时,发现它提交2个请求后就停滞了,必须等至少一个请求收到响应才会继续提交新请求,日志如下:
newwww-sch-1 Do Request 11 newwww-sch-2 Do Request 12 newwww-sch-2 Got Response 12 newwww-sch-1 Got Response 11 newwww-sch-2 Do Request 13 newwww-sch-1 Do Request 14 newwww-sch-1 Got Response 14 newwww-sch-2 Got Response 13
另外,提交请求的线程和处理响应的线程是同一个,这是否符合预期?我的预期是只要调度器线程可用(或者线程在等待服务器响应时),就会持续提交请求。
问题分析与解决
1. 提交2个请求后停滞的原因
flatMap默认并发数为2,加上你用subscribeOn(sc)将整个流的执行绑定到仅含2个线程的调度器上,而你的API调用是同步阻塞方法——调用后会占用当前线程直到响应返回,两个线程被完全占住后,没有多余线程处理后续请求,自然要等线程释放才会继续提交。
2. 提交与响应线程相同是否符合预期
符合预期。subscribeOn(sc)指定了整个流的所有逻辑(包括flatMap里的请求发起、subscribe里的响应处理)都在该调度器的线程上执行。由于API调用是同步阻塞的,线程在等待响应期间不会释放,响应返回后自然由同一个线程继续处理后续逻辑。
3. 实现“持续提交请求”的调整方案
要达成预期,需做两个关键修改:
- 隔离同步阻塞调用:将同步API调用放到专门的阻塞调度器(如
Schedulers.boundedElastic())中执行,避免占用主调度器线程。可以用Mono.fromCallable()包装同步调用,再通过subscribeOn指定阻塞调度器。 - 调整
flatMap并发数:根据需求设置更大的并发数(如flatMap(..., 10)),让flatMap可以同时处理更多元素,配合异步化的阻塞调用,线程不会被长时间占用,就能持续提交请求。
修改后的示例代码:
private Scheduler sc = Schedulers.newBoundedElastic(2, Integer.MAX_VALUE, "newwww-sch"); // 专门处理阻塞调用的调度器 private Scheduler blockingScheduler = Schedulers.boundedElastic(); Flux.range(11, 4) .flatMap(x -> { return Mono.fromCallable(() -> { System.out.println(Thread.currentThread().getName() + " Do Request " + x); return spectrumAvailabilityApi.getSpectrumAvailability(spectrumAvailabilityRequest); }) .subscribeOn(blockingScheduler); // 把阻塞调用隔离到专门调度器 }, 10) // 设置flatMap并发数为10 .subscribeOn(sc) .subscribe(s -> System.out.println(Thread.currentThread().getName() + " Got Response " + s));
内容的提问来源于stack exchange,提问作者quintin
相关产品推荐
相关产品推荐

