如何在等待所有虚拟线程结束时正确使用SimpleAsyncTaskExecutor
优化虚拟线程并行任务等待方案
你的需求是并行执行数据库查询填充ConcurrentHashMap,待所有任务完成后串行发送邮件,当前用ExecutorCompletionService的思路可行,但可以优化实现逻辑,同时解决异常处理的潜在问题,以下是几种更优方案:
一、改进现有ExecutorCompletionService实现
当前代码核心逻辑没问题,但需补充任务执行异常的处理,避免单个任务失败导致数据不完整或流程异常:
@Component public class NotificacoesListener { // 将Executor配置为Spring Bean,复用实例而非每次创建 @Bean public Executor virtualThreadExecutor() { SimpleAsyncTaskExecutor executor = new SimpleAsyncTaskExecutor(); executor.setVirtualThreads(true); executor.setConcurrencyLimit(4); return executor; } @Autowired private Executor virtualThreadExecutor; @Scheduled(cron = MySettings.NOTIFICATION_SUB_SYSTEM_EVERY_NIGHT, zone = "GMT-3:00") @Transactional public void start() { System.out.println("Starting at " + Utils.formatDate(LocalDate.now()) + " ..."); ConcurrentHashMap<User, String> map = new ConcurrentHashMap<>(); ExecutorCompletionService<Void> completionService = new ExecutorCompletionService<>(virtualThreadExecutor); int taskCount = instancias.size(); for (Verifica instancia : instancias) { completionService.submit(() -> { instancia.verifica(tenants, map); return null; }); } // 等待所有任务完成,同时捕获任务执行异常 for (int i = 0; i < taskCount; i++) { try { completionService.take().get(); // 触发任务异常抛出 } catch (InterruptedException e) { Thread.currentThread().interrupt(); // 恢复线程中断状态 throw new RuntimeException("任务被中断", e); } catch (ExecutionException e) { throw new RuntimeException("任务执行失败", e.getCause()); } } sendEmails(map); System.out.println("END"); } }
优化点:
- 把
Executor配置为单例Bean,符合Spring资源管理规范 - 使用泛型
ExecutorCompletionService<Void>,代码可读性更强 - 显式调用
Future.get()捕获任务执行异常,避免失败任务被静默忽略 - 处理
InterruptedException时恢复线程中断状态,符合并发编程规范
二、使用CompletableFuture(Java 19+ 推荐)
如果项目基于Java 19及以上版本,直接用CompletableFuture结合虚拟线程API会更简洁,无需手动管理CompletionService:
@Component public class NotificacoesListener { @Scheduled(cron = MySettings.NOTIFICATION_SUB_SYSTEM_EVERY_NIGHT, zone = "GMT-3:00") @Transactional public void start() { System.out.println("Starting at " + Utils.formatDate(LocalDate.now()) + " ..."); ConcurrentHashMap<User, String> map = new ConcurrentHashMap<>(); List<CompletableFuture<Void>> futures = new ArrayList<>(); for (Verifica instancia : instancias) { CompletableFuture<Void> future = CompletableFuture.runAsync( () -> instancia.verifica(tenants, map), Thread.ofVirtual().name("db-query-", 0).factory() // 原生虚拟线程工厂 ); futures.add(future); } // 等待所有任务完成,任一任务失败立即终止流程 try { CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join(); } catch (CompletionException e) { throw new RuntimeException("部分任务执行失败", e.getCause()); } sendEmails(map); System.out.println("END"); } }
优势:
- 代码更简洁,无需手动维护任务计数器
- 依赖Java原生虚拟线程API,无需Spring特定Executor
CompletableFuture.allOf()原生支持批量等待,还可通过exceptionally()等方法灵活处理异常
三、直接使用Future列表等待
如果不想用CompletionService或CompletableFuture,也可以直接收集所有Future对象逐个等待:
@Component public class NotificacoesListener { @Autowired private Executor virtualThreadExecutor; @Scheduled(cron = MySettings.NOTIFICATION_SUB_SYSTEM_EVERY_NIGHT, zone = "GMT-3:00") @Transactional public void start() { System.out.println("Starting at " + Utils.formatDate(LocalDate.now()) + " ..."); ConcurrentHashMap<User, String> map = new ConcurrentHashMap<>(); List<Future<Void>> futures = new ArrayList<>(); for (Verifica instancia : instancias) { Future<Void> future = virtualThreadExecutor.submit(() -> { instancia.verifica(tenants, map); return null; }); futures.add(future); } // 逐个等待任务完成 for (Future<Void> future : futures) { try { future.get(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("任务中断", e); } catch (ExecutionException e) { throw new RuntimeException("任务失败", e.getCause()); } } sendEmails(map); System.out.println("END"); } }
注意:
- 此方式按任务提交顺序等待,而
ExecutorCompletionService按任务完成顺序获取结果,若仅需等待全部完成,两种方式均适用;若需任务完成后立即处理结果,CompletionService更合适
关键注意事项
无论采用哪种方案,都需关注:
- 异常处理:必须捕获任务执行异常,避免单个任务失败导致整个流程静默异常
- 中断处理:遇到
InterruptedException时恢复线程中断状态,避免后续逻辑忽略中断信号 - 并发限制:合理设置虚拟线程并发数,避免数据库连接池过载
内容的提问来源于stack exchange,提问作者KenobiBastila
相关产品推荐
相关产品推荐

