如何并行执行两个带不同超时时间的ExecutorService invokeAll()?
同时执行两个带不同超时时间的Callable任务列表的方案
你的示例代码中,两次invokeAll是串行执行的——第一个任务列表执行完成(或超时)后,才会启动第二个列表的任务;而逐个submit后调用get也是顺序等待每个任务完成,都无法实现两个任务列表同时并行执行且各自遵守超时时间的需求。
方案一:用额外线程池并行触发两个invokeAll
把两个任务列表的invokeAll操作本身作为任务提交到另一个线程池,让它们同时启动执行,这样两个任务列表的任务就能在主线程池里并行运行,各自的超时时间独立生效。
示例代码:
import java.util.ArrayList; import java.util.Arrays; import java.util.List; import java.util.concurrent.*; public class ParallelTaskGroups { public static void main(String[] args) throws InterruptedException, ExecutionException { // 用于执行实际业务任务的线程池,根据任务数量调整大小 ExecutorService taskExecutor = Executors.newFixedThreadPool(4); // 用于并行触发两个任务组的invokeAll操作 ExecutorService controlExecutor = Executors.newSingleThreadExecutor(); // 初始化两个任务列表(示例) List<Callable<String>> taskGroup1 = Arrays.asList( () -> { Thread.sleep(3000); return "Task1-1"; }, () -> { Thread.sleep(4000); return "Task1-2"; } ); List<Callable<String>> taskGroup2 = Arrays.asList( () -> { Thread.sleep(1000); return "Task2-1"; }, () -> { Thread.sleep(3000); return "Task2-2"; } ); // 封装第一个任务组的处理逻辑 Callable<List<Future<String>>> groupHandler1 = () -> taskExecutor.invokeAll(taskGroup1, 5, TimeUnit.SECONDS); // 封装第二个任务组的处理逻辑 Callable<List<Future<String>>> groupHandler2 = () -> taskExecutor.invokeAll(taskGroup2, 2, TimeUnit.SECONDS); // 并行启动两个任务组的invokeAll List<Future<List<Future<String>>>> groupFutures = controlExecutor.invokeAll( Arrays.asList(groupHandler1, groupHandler2) ); // 收集所有任务结果 List<Future<String>> allResults = new ArrayList<>(); for (Future<List<Future<String>>> groupFuture : groupFutures) { allResults.addAll(groupFuture.get()); } // 处理结果 for (Future<String> future : allResults) { if (future.isCancelled()) { System.out.println("任务已取消"); } else { try { System.out.println("获取结果: " + future.get()); } catch (InterruptedException | ExecutionException e) { System.err.println("任务执行异常: " + e.getMessage()); } } } // 关闭线程池 controlExecutor.shutdown(); taskExecutor.shutdown(); } }
方案二:直接用线程触发invokeAll
如果不想额外创建线程池,可以直接启动两个线程,分别执行两个任务组的invokeAll:
示例代码:
import java.util.ArrayList; import java.util.Arrays; import java.util.List; import java.util.concurrent.*; public class ParallelTaskGroupsWithThreads { public static void main(String[] args) throws InterruptedException { ExecutorService taskExecutor = Executors.newFixedThreadPool(4); List<Callable<String>> taskGroup1 = Arrays.asList( () -> { Thread.sleep(3000); return "Task1-1"; }, () -> { Thread.sleep(4000); return "Task1-2"; } ); List<Callable<String>> taskGroup2 = Arrays.asList( () -> { Thread.sleep(1000); return "Task2-1"; }, () -> { Thread.sleep(3000); return "Task2-2"; } ); List<Future<String>> group1Results = new ArrayList<>(); List<Future<String>> group2Results = new ArrayList<>(); // 线程1处理第一个任务组 Thread group1Thread = new Thread(() -> { try { group1Results.addAll(taskExecutor.invokeAll(taskGroup1, 5, TimeUnit.SECONDS)); } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.err.println("任务组1执行被中断"); } }); // 线程2处理第二个任务组 Thread group2Thread = new Thread(() -> { try { group2Results.addAll(taskExecutor.invokeAll(taskGroup2, 2, TimeUnit.SECONDS)); } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.err.println("任务组2执行被中断"); } }); // 启动两个线程 group1Thread.start(); group2Thread.start(); // 等待两个线程完成 group1Thread.join(); group2Thread.join(); // 合并结果并处理 List<Future<String>> allResults = new ArrayList<>(); allResults.addAll(group1Results); allResults.addAll(group2Results); for (Future<String> future : allResults) { if (future.isCancelled()) { System.out.println("任务已取消"); } else { try { System.out.println("获取结果: " + future.get()); } catch (InterruptedException | ExecutionException e) { System.err.println("任务执行异常: " + e.getMessage()); } } } taskExecutor.shutdown(); } }
注意事项
- 执行实际任务的线程池大小要合理设置,确保两个任务组的任务都能被及时调度,避免因线程不足导致任务等待。
- 注意捕获
InterruptedException和ExecutionException,避免异常导致程序崩溃。 - 线程池使用完毕后必须调用
shutdown()(或shutdownNow())释放资源。
内容的提问来源于stack exchange,提问作者Gunther
相关产品推荐
相关产品推荐

