【问题标题】:Track progress of concurrently executed tasks跟踪并发执行任务的进度
【发布时间】:2021-09-29 11:21:20
【问题描述】:

我有一堆任务要并行执行。任务分为阶段,每个阶段的任务取决于前一阶段任务的结果。因此,一个阶段上的所有任务必须在递归地进入下一个阶段之前完成执行。代码如下:

java.util.concurrent.Executor executor;

void process(Stage stage) throws InterruptedException {
    Task[] tasks = stage.tasks();
    CountDownLatch countDownLatch = new CountDownLatch(tasks.length);

    logger.info("Execute {} tasks on stage {}", tasks.length, stage.index());
    for (Task task : tasks) {
        executor.execute(() -> {
            this.execute(task);
            countDownLatch.countDown();
        });
    }

    countDownLatch.await(); // Wait for all tasks to finish
    if (! stage.isLast()) {
        process(stage.next()); // Execute tasks on the next stage
    }
}

到目前为止,一切正常。但是,我不想记录每个阶段的进度,而是定期在单独的TimerTask

class Progress extends TimerTask {
    @Override
    public void run() {
        logger.info("Executed {}/{} tasks on stage {}/{}",
            executedTasksOnStage, totalTasksOnStage, currentStageIndex, totalStageCount);
    }
}

如何以线程安全的方式将上面记录的变量从process() 方法传递给Progress 对象?我熟悉简单的原子计数器,但在这种情况下,有多个变量一起更新,然后再单独更新。这也是一个设计问题,如果您能提供代码示例,我将不胜感激。

【问题讨论】:

  • 有什么理由不在舞台上积累统计数据等?这似乎是维护您的状态信息的自然场所。一个更高的概念,管道,可以跟踪它处于哪个阶段。然后报告可以简单地查询管道和阶段的执行细节。
  • 感谢您的意见!我将发布我的解决方案作为答案。

标签: java concurrency timertask


【解决方案1】:

我能想到的最简单的解决方案是static volatile boolean threadOneStageComplete = false;,当线程一个完成一个阶段并以一个while循环结束每个阶段时设置为真,它测试!(threadOneStageComplete && threadTwoStageComplete && ...)并让线程休眠任何时间在你的情况下感觉。确保在进入下一阶段时再次将变量设置为 false,并可能将每个任务放入一个类中,这样您就不会拥有一堆静态 volatile 布尔值,而是每个任务对象中都有一个 volatile 属性。

【讨论】:

    【解决方案2】:

    好的,多亏了 Ben Manes,我提出了以下解决方案:我引入了一个 Pipeline 对象,它包含对当前阶段的原子引用,由 PipelineStage 包裹。包装器还包含一个CountDownLatch,用于跟踪从包装阶段执行的任务数。

    class PipelineStage {
        final Stage stage;
        final CountDownLatch latch;
    
        PipelineStage(Stage stage) {
            this.stage = stage;
            latch = new CountDownLatch(stage.tasks().length);
        }
    
        int index() {
            return stage.index();
        }
    
        int totalTaskCount() {
            return stage.tasks().length;
        }
    
        int executedTaskCount() {
            return totalTaskCount() - (int) latch.getCount();
        }
    
        Stream<Runnable> tasks() {
            return Stream.of(stage.tasks()).map(task -> () -> {
                task.run();
                latch.countDown();
            });
        }
    
        void await() throws InterruptedException {
            latch.await();
        }
    
        PipelineStage next() {
            return new PipelineStage(stage.next());
        }
    }
    
    class Pipeline {
        final AtomicReference<PipelineStage> stage;
    
        Pipeline(PipelineStage stage) {
            this.stage = new AtomicReference<>(stage);
        }
    
        PipelineStage current() {
            return stage.get();
        }
    
        PipelineStage advance() {
            return stage.updateAndGet(PipelineStage::next);
        }
    }
    

    现在原问题中的process() 方法变为:

    java.util.concurrent.Executor executor;
    
    void process(Pipeline pipeline) throws InterruptedException {
        PipelineStage stage = pipeline.current();
    
        while (stage != null) {
            stage.tasks().forEach(executor::execute);
            stage.await();
            stage = pipeline.advance();
        }
    }
    

    最后,管道可以传递给任意数量的并发阅读器,例如进度记录器:

    @Override
    public void run() {
        final PipelineStage stage = pipeline.current();
    
        logger.debug("Executed {}/{} tasks on stage {}",
            stage.executedTaskCount(), stage.totalTaskCount(), stage.index())
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2013-02-14
      • 2014-09-18
      • 2016-05-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-08-18
      • 1970-01-01
      相关资源
      最近更新 更多