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

Java多线程依赖处理器调度:等待前置Callable任务完成

问题描述

我需要在两个线程中执行若干处理器,其中部分处理器相互独立可随时运行,部分存在依赖关系。当执行流程到达依赖处理器时,需检查所有前置Callable任务是否已执行完成,且待当前依赖处理器执行完成后再执行后续任务。

尝试用while(future.isDone())检测任务状态时出现无限循环,现需解决两个问题:

  1. 如何检测线程/Callable任务是否已启动;
  2. 执行串行依赖处理器时,等待所有现有任务完成后再启动当前任务,且待当前任务完成后再取下一个任务。

现有代码

主线程方法

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 23:05:20