You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 11:58:32