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

如何强制Reactor在非守护线程上执行subscribe()代码?

如何让Reactor在非守护线程执行subscribe代码?

这是系列问题的后续,核心需求是:运行WebClient异步代码时,JVM能执行subscribe里的打印语句,不立即退出,且不使用block()方法。

直接运行以下代码不会有任何输出:

@Test
void test() {
    WebClient.builder().build()
        .get()
        .uri("https://httpbin.org/status/404")
        .retrieve()
        .bodyToMono(String.class)
        .subscribe(
                s -> System.out.println("Success: " + s),
                t -> System.out.println("Failure: " + t.getMessage())
        );
}

原因是Reactor默认使用守护线程执行异步逻辑,JVM主线程结束后会直接退出,不会等待守护线程完成任务。

尝试添加subscribeOn(Schedulers.newParallel("my-scheduler", 5, false))(false指定非守护线程)或更换其他Scheduler(如bounded elastic)均无效,问题出在WebClient的HTTP请求逻辑默认运行在Netty的IO线程池,该线程池的线程是守护线程,subscribeOn仅能影响订阅阶段的线程,无法覆盖后续HTTP请求执行的线程。

以下是两种可靠的解决方案:

方案1:自定义Netty非守护线程池

直接修改WebClient依赖的Netty线程池,将线程设置为非守护线程:

@Test
void test() {
    // 自定义非守护线程的EventLoopGroup
    EventLoopGroup eventLoopGroup = new NioEventLoopGroup(
            5,
            new ThreadFactoryBuilder()
                    .setNameFormat("custom-io-%d")
                    .setDaemon(false) // 关键:设置为非守护线程
                    .build()
    );

    try {
        WebClient webClient = WebClient.builder()
                .clientConnector(new ReactorClientHttpConnector(
                        HttpClient.create().runOn(eventLoopGroup)
                ))
                .build();

        webClient.get()
                .uri("https://httpbin.org/status/404")
                .retrieve()
                .bodyToMono(String.class)
                .subscribe(
                        s -> System.out.println("Success: " + s),
                        t -> System.out.println("Failure: " + t.getMessage())
                );
    } finally {
        // 测试结束后优雅关闭线程池
        eventLoopGroup.shutdownGracefully();
    }
}

此时HTTP请求的IO操作会在自定义非守护线程上执行,JVM会等待这些线程完成后再退出,就能看到打印输出。

方案2:用CountDownLatch让主线程等待

如果不想修改WebClient线程池,可通过CountDownLatch让主线程等待异步任务完成:

@Test
void test() throws InterruptedException {
    CountDownLatch latch = new CountDownLatch(1);

    WebClient.builder().build()
            .get()
            .uri("https://httpbin.org/status/404")
            .retrieve()
            .bodyToMono(String.class)
            .doFinally(signalType -> latch.countDown()) // 任务完成后计数减1
            .subscribe(
                    s -> System.out.println("Success: " + s),
                    t -> System.out.println("Failure: " + t.getMessage())
            );

    latch.await(); // 主线程等待直到latch计数为0
}

这种方式无需修改线程类型,通过阻塞主线程确保异步任务执行完毕,同样能触发打印逻辑。

内容的提问来源于stack exchange,提问作者Sergey Zolotarev

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 09:43:39