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

如何在等待所有虚拟线程结束时正确使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 21:08:14