ExecutorService在Quarkus中阻塞主线程时无法并行执行任务的问题
问题描述
运行以下代码时,若注释掉最后一行的latch.await();,得到输出[2],任务能并行执行;取消注释后得到输出[3],任务变成串行执行。需要解决的问题是:如何在阻塞主线程等待所有任务完成的同时,让ExecutorService保持并行执行任务。
补充:代码运行在Quarkus框架中,属于@RequestScoped上下文资源。
代码[1]
private void setProfileData(String clientId, List<ProfileData> profileDataList) throws InterruptedException { final int chunkSize = 500; var chunks = IntStream.range(0, (profileDataList.size() + chunkSize - 1) / chunkSize) .mapToObj(i -> profileDataList.subList(i * chunkSize, Math.min(chunkSize * (i + 1), profileDataList.size()))).collect(Collectors.toList()); ExecutorService executor = Executors.newFixedThreadPool(8); CountDownLatch latch = new CountDownLatch(chunks.size()); for (var c : chunks) { executor.submit(() -> { try { log.error("Start write"); //simulate work on the chunk Thread.sleep(5000); log.error("End write"); } catch (InterruptedException e) { e.printStackTrace(); } finally { latch.countDown(); } }); } executor.shutdown(); //latch.await(); // wait for all tasks to complete }
输出[2](并行执行,注释latch.await()时)
2023-06-21 17:30:25,838 ERROR [com.fan.fun.ser.MigrationService] (pool-147-thread-1) Start write 2023-06-21 17:30:25,845 ERROR [com.fan.fun.ser.MigrationService] (pool-148-thread-1) Start write 2023-06-21 17:30:25,851 ERROR [com.fan.fun.ser.MigrationService] (pool-149-thread-1) Start write ... 2023-06-21 17:30:30,843 ERROR [com.fan.fun.ser.MigrationService] (pool-147-thread-1) End write 2023-06-21 17:30:30,849 ERROR [com.fan.fun.ser.MigrationService] (pool-148-thread-1) End write 2023-06-21 17:30:30,856 ERROR [com.fan.fun.ser.MigrationService] (pool-149-thread-1) End write
输出[3](串行执行,取消latch.await()注释时)
2023-06-21 17:32:22,571 ERROR [com.fan.fun.ser.MigrationService] (pool-177-thread-1) Start write 2023-06-21 17:32:27,576 ERROR [com.fan.fun.ser.MigrationService] (pool-177-thread-1) End write 2023-06-21 17:32:27,592 ERROR [com.fan.fun.ser.MigrationService] (pool-178-thread-1) Start write 2023-06-21 17:32:32,598 ERROR [com.fan.fun.ser.MigrationService] (pool-178-thread-1) End write 2023-06-21 17:32:32,615 ERROR [com.fan.fun.ser.MigrationService] (pool-179-thread-1) Start write 2023-06-21 17:32:37,620 ERROR [com.fan.fun.ser.MigrationService] (pool-179-thread-1) End write
解决方案
问题根源在于Quarkus的@RequestScoped上下文会绑定到请求线程,你自定义的ExecutorService创建的线程无法继承请求上下文,Quarkus的线程隔离机制会限制这些线程的执行,导致任务串行。解决方法如下:
使用Quarkus托管的ManagedExecutor
在你的类中通过依赖注入获取ManagedExecutor,它会自动处理上下文传播,确保任务能并行执行:import jakarta.enterprise.concurrent.ManagedExecutor; import jakarta.inject.Inject; @RequestScoped public class MigrationService { @Inject ManagedExecutor managedExecutor; private void setProfileData(String clientId, List<ProfileData> profileDataList) throws InterruptedException { final int chunkSize = 500; var chunks = IntStream.range(0, (profileDataList.size() + chunkSize - 1) / chunkSize) .mapToObj(i -> profileDataList.subList(i * chunkSize, Math.min(chunkSize * (i + 1), profileDataList.size()))).collect(Collectors.toList()); CountDownLatch latch = new CountDownLatch(chunks.size()); for (var c : chunks) { managedExecutor.submit(() -> { try { log.error("Start write"); //simulate work on the chunk Thread.sleep(5000); log.error("End write"); } catch (InterruptedException e) { e.printStackTrace(); } finally { latch.countDown(); } }); } latch.await(); // 此时任务会并行执行 } }异步风格替代方案:CompletableFuture + ManagedExecutor
如果更倾向于异步编程,可以结合CompletableFuture,不需要手动管理CountDownLatch:private void setProfileData(String clientId, List<ProfileData> profileDataList) throws InterruptedException, ExecutionException { final int chunkSize = 500; var chunks = IntStream.range(0, (profileDataList.size() + chunkSize - 1) / chunkSize) .mapToObj(i -> profileDataList.subList(i * chunkSize, Math.min(chunkSize * (i + 1), profileDataList.size()))).collect(Collectors.toList()); List<CompletableFuture<Void>> futures = chunks.stream() .map(chunk -> CompletableFuture.runAsync(() -> { log.error("Start write"); try { Thread.sleep(5000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } log.error("End write"); }, managedExecutor) ) .collect(Collectors.toList()); // 等待所有任务完成 CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).get(); }
关键在于避免自行创建线程池,改用Quarkus托管的线程池,这样才能和框架的上下文管理、线程调度机制兼容,保证任务并行执行的同时,主线程可以等待任务完成。
内容的提问来源于stack exchange,提问作者Julian
相关产品推荐
相关产品推荐

