使用@Scheduled的fixedDelay结合parallelStream与ForkJoinPool遇调度问题求助
解决@Scheduled fixedDelay未等待并行任务完成的问题
问题根源
你的代码中,runJob方法调用pool.submit()提交并行任务后,仅执行pool.shutdown()就直接返回,并未等待所有并行任务处理完毕。而@Scheduled(fixedDelay)是以调度方法run()执行结束为起点开始计时的,所以run()会快速完成,导致调度器每3秒就触发新任务,完全没等并行处理结束。
另外需要注意:你创建的ForkJoinPool(2)并没有被parallelStream实际使用——parallelStream默认依赖ForkJoinPool.commonPool(),除非代码运行在自定义ForkJoinPool的任务上下文内,否则无法通过这种方式限制并行数为2。
修正方案
方案一:等待自定义池任务全部完成
修改runJob,在提交任务后等待所有执行结束,同时调整并行任务的提交方式,确保使用自定义池控制并行数:
public void runJob() throws InterruptedException { ForkJoinPool pool = new ForkJoinPool(2); List<String> list = Arrays.asList("1","2","3"); // 将每个元素的处理包装为独立任务,提交到自定义池 List<ForkJoinTask<Void>> tasks = list.stream() .map(item -> ForkJoinTask.adapt(() -> execute(item))) .collect(Collectors.toList()); // 等待所有任务执行完成 pool.invokeAll(tasks); pool.shutdown(); // 确保池关闭前所有任务收尾 if (!pool.awaitTermination(1, TimeUnit.HOURS)) { pool.shutdownNow(); } }
方案二:用invoke替代submit等待任务完成
ForkJoinPool.invoke()会阻塞直到提交的任务全部执行完毕,这样runJob会等到并行处理结束后才返回:
public void runJob() { ForkJoinPool pool = new ForkJoinPool(2); List<String> list = Arrays.asList("1","2","3"); // invoke会等待任务完成后才返回 pool.invoke(() -> list.forEach(this::execute) // 若需并行,建议用方案一的任务提交方式,确保使用自定义池 ); pool.shutdown(); }
方案三:用@Async配合自定义线程池
如果项目已配置自定义线程池,可通过CompletableFuture实现并行处理并等待所有任务完成:
// 先配置自定义线程池 @Configuration public class AsyncConfig { @Bean("taskExecutor") public Executor taskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(2); executor.setMaxPoolSize(2); executor.setThreadNamePrefix("Parallel-"); executor.initialize(); return executor; } } // 业务代码修改 @Scheduled(fixedDelay = 3000) public void run() throws ExecutionException, InterruptedException { System.out.println("Run"); runJob(); } @Autowired @Qualifier("taskExecutor") private Executor taskExecutor; public void runJob() throws ExecutionException, InterruptedException { List<String> list = Arrays.asList("1","2","3"); List<CompletableFuture<Void>> futures = list.stream() .map(item -> CompletableFuture.runAsync(() -> execute(item), taskExecutor)) .collect(Collectors.toList()); // 等待所有异步任务执行完成 CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).get(); }
内容的提问来源于stack exchange,提问作者mayank bisht
相关产品推荐
相关产品推荐

