Java多线程依赖处理器调度:等待前置Callable任务完成
问题描述
我需要在两个线程中执行若干处理器,其中部分处理器相互独立可随时运行,部分存在依赖关系。当执行流程到达依赖处理器时,需检查所有前置Callable任务是否已执行完成,且待当前依赖处理器执行完成后再执行后续任务。
尝试用while(future.isDone())检测任务状态时出现无限循环,现需解决两个问题:
- 如何检测线程/Callable任务是否已启动;
- 执行串行依赖处理器时,等待所有现有任务完成后再启动当前任务,且待当前任务完成后再取下一个任务。
现有代码
主线程方法
PackageExportGraph executeMultiThread(PackageExportGraph exportGraphInp, PackageExportContext exportContextnInp) throws WTException { Map<PackageExportDependencyProcessor, Boolean> processorToParallelExecutionMap = new LinkedHashMap<>(); this.processorQueue = new LinkedBlockingQueue<>(); ExecutorService execService = null; try { int threads = 2;// 2 countDownLatch = new CountDownLatch(threads); execService = ExecutorServiceFactory.getDefault().newExecutorService(threads, true); boolean isThread1Started = false; ThreadedDepProcessor thread1 = new ThreadedDepProcessor( exportGraphInp, countDownLatch, processorToParallelExecutionMap, processorQueue, exportContextnInp, isThread1Started); threadList.add(thread1); thread1.addListener(this); boolean isThread2Started = false; ThreadedDepProcessor thread2 = new ThreadedDepProcessor( exportGraphInp, countDownLatch, processorToParallelExecutionMap, processorQueue, exportContextnInp, isThread2Started); threadList.add(thread2); thread1.addListener(this); List<Future<LinkedBlockingQueue>> futureList = new ArrayList<>(); for (ThreadedDepProcessor thread : threadList) { Future f = execService.submit(thread); System.out.println("f " + f); futureList.add(f); } int currentidx = 0; for (PackageExportDependencyProcessor processor : origOrderedList) { if (!processorToParallelExecutionMap.containsKey(processor)) { System.out.println(" parallel threadStatusMap values 1 - " + threadStatusMap.values()); System.out.println("Adding parallel - " + processor); if (currentidx > 0) { while (threadStatusMap.containsValue(false)) { System.out.println("Waiting"); System.out.println("threadStatusMap values - " + threadStatusMap.values()); Thread.sleep(1000); } Thread.sleep(2000); // execService.awaitTermination(5, TimeUnit.SECONDS); System.out.println("Size - " + futureList.size()); for (Future f : futureList) { System.out.println("futureList is done " + f.isDone()); System.out.println("Getting future Object"); if (f.isDone()) { continue; } Object o = f.get(10, TimeUnit.SECONDS); System.out.println(o); /* * Object object = f.get(10, TimeUnit.SECONDS); System.out.println("Obj " + object); */ } processorQueue.put(processor); Thread.sleep(2000); } else { processorQueue.put(processor); Thread.sleep(2000); System.out.println("Size - " + futureList.size()); for (Future f : futureList) { System.out.println("futureList is done " + f.isDone()); System.out.println("Getting future Object"); if (f.isDone()) { continue; } Object o = f.get(10, TimeUnit.SECONDS); System.out.println(o); /* * Object object = f.get(10, TimeUnit.SECONDS); System.out.println("Obj " + object); */ } // execService.awaitTermination(5, TimeUnit.SECONDS); while (threadStatusMap.containsValue(false)) { System.out.println("Waiting"); System.out.println("threadStatusMap values - " + threadStatusMap.values()); Thread.sleep(1000); } } for (ThreadedDepProcessor thread : threadList) { System.out.println("Finished adding dependents" + thread.finishedAddingDependents.get()); } } else { System.out.println("Adding non-parallel - " + processor); processorQueue.put(processor); } currentidx++; } } catch (WTException | RuntimeException exc) { if (Objects.nonNull(execService)) { execService.shutdown(); } throw exc; } catch (Exception exc) { throw new WTException(exc); } finally { System.out.println("shutting down"); try { countDownLatch.await(); } catch (InterruptedException e) { e.printStackTrace(); } execService.shutdown(); } return exportGraphInp; }
Callable实现代码
@Override public LinkedBlockingQueue call() throws WTException, InterruptedException { try { System.out.println("Started - "); isThreadStarted = true; while (!processorQueue.isEmpty()) { nextEligible = processorQueue.take(); if (Objects.isNull(nextEligible)) { finishedAddingDependents.set(true); break; } System.out.println("calling addDependentObjects for processor - " + nextEligible); nextEligible.addDependentObjects(exportGraph, exportContext); nextEligible = null; // notifyListeners(); } } catch (Exception e) { System.out.println("Error occured " + e); e.printStackTrace(); return processorQueue; } finally { countDownLatch.countDown(); System.out.println("countDownLatch now - " + countDownLatch.getCount()); } return processorQueue; }
解决方案
1. 检测Callable任务是否已启动
原代码中传入的isThreadStarted是值类型,修改后无法同步到主线程,需改用原子布尔变量来标记启动状态:
public class ThreadedDepProcessor implements Callable<LinkedBlockingQueue> { // 用原子变量存储启动状态,线程安全 private final AtomicBoolean isStarted = new AtomicBoolean(false); // 其他成员变量... @Override public LinkedBlockingQueue call() throws WTException, InterruptedException { try { // 任务启动时标记为true isStarted.set(true); System.out.println("Started - "); while (true) { PackageExportDependencyProcessor nextEligible = processorQueue.take(); if (nextEligible == null) { // 收到结束信号后退出循环 finishedAddingDependents.set(true); break; } System.out.println("calling addDependentObjects for processor - " + nextEligible); nextEligible.addDependentObjects(exportGraph, exportContext); } } catch (Exception e) { System.out.println("Error occured " + e); e.printStackTrace(); return processorQueue; } finally { countDownLatch.countDown(); System.out.println("countDownLatch now - " + countDownLatch.getCount()); } return processorQueue; } // 提供外部查询方法 public boolean isTaskStarted() { return isStarted.get(); } }
主线程中检测启动状态:
// 提交任务后,等待线程启动 for (ThreadedDepProcessor thread : threadList) { while (!thread.isTaskStarted()) { Thread.sleep(100); // 短暂等待,避免忙等消耗资源 } System.out.println("线程已启动: " + thread); }
2. 串行依赖处理器的同步执行
原代码用Thread.sleep()和轮询future.isDone()的方式不可靠,且容易出现无限循环。正确的做法是用Future.get()阻塞等待任务完成,区分并行/串行任务的处理逻辑:
修改主线程逻辑
PackageExportGraph executeMultiThread(PackageExportGraph exportGraphInp, PackageExportContext exportContextnInp) throws WTException { Map<PackageExportDependencyProcessor, Boolean> processorToParallelExecutionMap = new LinkedHashMap<>(); ExecutorService execService = null; try { int threads = 2; execService = ExecutorServiceFactory.getDefault().newExecutorService(threads, true); List<Future<?>> parallelFutures = new ArrayList<>(); for (PackageExportDependencyProcessor processor : origOrderedList) { if (processorToParallelExecutionMap.containsKey(processor)) { // 并行任务:提交到线程池,保存Future Future<?> future = execService.submit(() -> { processor.addDependentObjects(exportGraphInp, exportContextnInp); }); parallelFutures.add(future); System.out.println("添加并行处理器: " + processor); } else { // 串行依赖任务:先等待所有已提交的并行任务完成 System.out.println("等待所有并行任务完成..."); for (Future<?> future : parallelFutures) { try { future.get(); // 阻塞直到任务完成,任务异常会抛出ExecutionException } catch (ExecutionException e) { // 转成业务异常抛出 throw new WTException(e.getCause()); } } parallelFutures.clear(); // 清空已完成的并行任务记录 // 执行当前串行处理器,确保完成后再处理下一个 System.out.println("执行串行处理器: " + processor); processor.addDependentObjects(exportGraphInp, exportContextnInp); } } // 处理最后一批并行任务 System.out.println("等待剩余并行任务完成..."); for (Future<?> future : parallelFutures) { try { future.get(); } catch (ExecutionException e) { throw new WTException(e.getCause()); } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new WTException(e); } catch (Exception exc) { throw new WTException(exc); } finally { if (execService != null) { execService.shutdown(); try { // 等待线程池关闭 if (!execService.awaitTermination(60, TimeUnit.SECONDS)) { execService.shutdownNow(); } } catch (InterruptedException e) { execService.shutdownNow(); } } } return exportGraphInp; }
关键说明
- 用
Future.get()替代轮询:get()会阻塞直到任务完成,无需手动轮询,避免无限循环和不必要的睡眠。 - 区分并行/串行任务:并行任务批量提交,串行任务前先等待所有并行任务完成,确保依赖顺序。
- 异常处理:捕获
ExecutionException并转换为业务异常,避免任务失败被忽略。
内容的提问来源于stack exchange,提问作者Dungeon Master
相关产品推荐
相关产品推荐

