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

Spring响应式编程:主线程未等待Mono订阅者完成任务的解决方法

解决Spring Reactor Mono并行执行后主线程未等待的问题

这个问题我之前也碰到过,核心就是你没让主线程等待Reactor的异步任务完成,而且手动实现Mono的方式也不太规范,咱们一步步来改:

1. 先修正Mono的创建方式

你现在手动重写subscribe方法的方式不符合Reactor的最佳实践,应该用Mono.fromCallable来包装阻塞操作(比如Thread.sleep和DB查询),这样能正确将任务提交到调度器上执行,同时规范异常处理:

private Mono<BaseResponse> getProfileDetails(long profileId) {
    return Mono.fromCallable(() -> {
        try {
            Thread.sleep(5000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new RuntimeException("任务被中断", e);
        }
        // DB Operation
        System.out.println("Inside getProfileDetails");
        return new BaseResponse();
    });
}

private Mono<Address> getAddressDetails(long profileId) {
    return Mono.fromCallable(() -> {
        try {
            Thread.sleep(5000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new RuntimeException("任务被中断", e);
        }
        // DB Operation
        System.out.println("Inside getAddressDetails");
        return new Address();
    });
}

2. 让主线程等待所有异步任务完成

原来的subscribe是非阻塞的,主线程会直接往下执行,导致列表还没被填充就打印了。我们需要用block()让主线程等待Reactor序列执行完成,同时用合适的操作符合并多个Mono并收集结果:

public BaseResponse getDetails(long profileId) {
    // 优先用Reactor自带的调度器,自动管理线程池,避免手动创建的资源泄露问题
    Scheduler scheduler = Schedulers.boundedElastic();
    // 如果你坚持用自定义线程池,也可以这么写:
    // ExecutorService executors = Executors.newFixedThreadPool(2);
    // Scheduler scheduler = Schedulers.fromExecutor(executors);

    Mono<BaseResponse> profileDetail = getProfileDetails(profileId).subscribeOn(scheduler);
    // 假设Address继承自BaseResponse,否则需要做类型转换
    Mono<BaseResponse> addressDetail = getAddressDetails(profileId).subscribeOn(scheduler);

    // 合并两个Mono,等待两者都完成后收集成列表
    List<BaseResponse> list = Flux.merge(profileDetail, addressDetail)
            .collectList()
            .block(); // 这里block会阻塞主线程,直到所有异步任务完成

    System.out.println("list: " + new Gson().toJson(list));

    // 如果用了自定义线程池,记得关闭
    // executors.shutdown();

    // 这里根据业务需求合并结果到返回的BaseResponse中
    BaseResponse response = new BaseResponse();
    // 示例:将addressDetail的内容合并到response
    if (list != null && list.size() > 1) {
        response.setAddress((Address) list.get(1));
    }
    return response;
}

为什么原来的代码会出问题?

  • subscribe方法是异步触发任务,但不会阻塞当前线程,所以主线程走到System.out.println时,两个Mono的任务还在后台执行,列表自然是空的。
  • 手动重写Mono.subscribe的逻辑,没有正确利用Reactor的调度机制,即使加了subscribeOn也可能不生效。
  • 主线程提前关闭了线程池,可能导致任务还没执行完就被中断。

额外注意事项

  • 如果是在WebFlux应用中,尽量不要用block(),会破坏响应式非阻塞特性;但如果你的getDetails是传统Spring MVC的同步方法,用block()是完全可以接受的。
  • 对于10个并行方法的场景,只需要把所有Mono放到Flux.merge里即可,它会自动并行执行所有任务,然后等待全部完成。

内容的提问来源于stack exchange,提问作者Manshavi Kumar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 08:47:45