如何强制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
相关产品推荐
相关产品推荐

