【问题标题】:Can I use the work-stealing behaviour of ForkJoinPool to avoid a thread starvation deadlock?我可以使用 ForkJoinPool 的工作窃取行为来避免线程饥饿死锁吗?
【发布时间】:2014-10-26 18:11:12
【问题描述】:

如果池中的所有线程都在等待同一个池中的队列任务完成,则在普通线程池中会发生线程饥饿死锁ForkJoinPool 通过从 join() 调用内部窃取其他线程的工作来避免这个问题,而不是简单地等待。例如:

private static class ForkableTask extends RecursiveTask<Integer> {
    private final CyclicBarrier barrier;

    ForkableTask(CyclicBarrier barrier) {
        this.barrier = barrier;
    }

    @Override
    protected Integer compute() {
        try {
            barrier.await();
            return 1;
        } catch (InterruptedException | BrokenBarrierException e) {
            throw new RuntimeException(e);
        }
    }
}

@Test
public void testForkJoinPool() throws Exception {
    final int parallelism = 4;
    final ForkJoinPool pool = new ForkJoinPool(parallelism);
    final CyclicBarrier barrier = new CyclicBarrier(parallelism);

    final List<ForkableTask> forkableTasks = new ArrayList<>(parallelism);
    for (int i = 0; i < parallelism; ++i) {
        forkableTasks.add(new ForkableTask(barrier));
    }

    int result = pool.invoke(new RecursiveTask<Integer>() {
        @Override
        protected Integer compute() {
            for (ForkableTask task : forkableTasks) {
                task.fork();
            }

            int result = 0;
            for (ForkableTask task : forkableTasks) {
                result += task.join();
            }
            return result;
        }
    });
    assertThat(result, equalTo(parallelism));
}

但是当使用ExecutorService 接口到ForkJoinPool 时,工作窃取似乎不会发生。例如:

private static class CallableTask implements Callable<Integer> {
    private final CyclicBarrier barrier;

    CallableTask(CyclicBarrier barrier) {
        this.barrier = barrier;
    }

    @Override
    public Integer call() throws Exception {
        barrier.await();
        return 1;
    }
}

@Test
public void testWorkStealing() throws Exception {
    final int parallelism = 4;
    final ExecutorService pool = new ForkJoinPool(parallelism);
    final CyclicBarrier barrier = new CyclicBarrier(parallelism);

    final List<CallableTask> callableTasks = Collections.nCopies(parallelism, new CallableTask(barrier));
    int result = pool.submit(new Callable<Integer>() {
        @Override
        public Integer call() throws Exception {
            int result = 0;
            // Deadlock in invokeAll(), rather than stealing work
            for (Future<Integer> future : pool.invokeAll(callableTasks)) {
                result += future.get();
            }
            return result;
        }
    }).get();
    assertThat(result, equalTo(parallelism));
}

粗略看一下ForkJoinPool的实现,所有常规的ExecutorService API都是使用ForkJoinTasks实现的,所以我不确定为什么会发生死锁。

【问题讨论】:

  • 我不认为偷工作能避免死锁。一旦陷入僵局,就无法取得进展。工作窃取只是通过允许线程在队列为空时从其他队列窃取来避免不平衡的队列。
  • @markspace 在ForkJoinTask 的实现中,join() 尝试从双端队列运行其他作业而不是停止,这样可以避免死锁。由于ForkJoinPool.invokeAll()Callables 转换为ForkJoinTasks,我希望它也能工作。

标签: java multithreading concurrency java.util.concurrent fork-join


【解决方案1】:

您几乎是在回答自己的问题。解决方案是声明“ForkJoinPool 通过从join() 调用内部窃取其他线程的工作来避免此问题”。每当线程因为ForkJoinPool.join()以外的其他原因而被阻塞时,这种工作窃取不会发生,线程只是等待而不做任何事情。

这样做的原因是,在 Java 中,ForkJoinPool 不可能阻止其线程阻塞,而是给它们一些其他的工作。线程本身需要避免阻塞,而是要求池完成它应该做的工作。而且这仅在ForkJoinTask.join() 方法中实现,在任何其他阻塞方法中均未实现。如果你在ForkJoinPool 中使用Future,你也会看到饥饿死锁。

为什么工作窃取只在ForkJoinTask.join() 中实现,而没有在 Java API 中的任何其他阻塞方法中实现?嗯,有很多这样的阻塞方法(Object.wait()Future.get()java.util.concurrent 中的任何并发原语、I/O 方法等),它们与 ForkJoinPool 无关,这只是一个API 中的任意类,因此向所有这些方法添加特殊情况将是糟糕的设计。它还会导致可能非常令人惊讶和不希望的效果。例如,假设用户将任务传递给等待FutureExecutorService,然后发现该任务在Future.get() 中挂起很长时间,只是因为正在运行的线程窃取了一些其他(长时间运行的)工作项而不是等待Future 并在结果可用后立即继续。一旦一个线程开始处理另一个任务,它就不能返回到原来的任务,直到第二个任务完成。因此,其他阻塞方法不进行工作窃取实际上是一件好事。对于ForkJoinTask,不存在这个问题,因为主任务是否尽快继续并不重要,重要的是所有任务一起尽可能高效地处理。

也不可能在ForkJoinPool 中实现您自己的工作窃取方法,因为所有相关部分都不公开。

然而,实际上还有第二种方法可以防止饥饿死锁。这称为托管屏蔽。它不使用工作窃取(避免上面提到的问题),但还需要将要被阻塞的线程积极配合线程池。使用托管阻塞,线程在调用潜在阻塞方法之前告诉线程池它可能被阻塞,并且在阻塞方法完成时通知线程池。然后线程池知道存在饥饿死锁的风险,并且如果它的所有线程当前都处于某个阻塞操作中并且还有其他任务要执行,则可能会产生额外的线程。请注意,由于额外线程的开销,这比工作窃取效率低。如果你用普通的期货和托管阻塞来实现递归并行算法,而不是用ForkJoinTask和工作窃取,额外线程的数量会变得非常大(因为在算法的“划分”阶段,很多任务将创建并提供给立即阻塞并等待子任务结果的线程)。但是,仍然可以防止饥饿死锁,并且避免了一个任务必须等待很长时间,因为它的线程同时开始在另一个任务上工作的问题。

Java 的ForkJoinPool 也支持托管阻塞。要使用它,需要实现接口ForkJoinPool.ManagedBlocker,以便从该接口的block 方法中调用任务想要执行的潜在阻塞方法。那么任务可能不会直接调用阻塞方法,而是需要调用静态方法ForkJoinPool.managedBlock(ManagedBlocker)。该方法处理阻塞前后与线程池的通信。如果当前任务没有在ForkJoinPool 中执行,它也可以工作,那么它只是调用阻塞方法。

我在 Java API(针对 Java 7)中发现的唯一一个真正使用托管阻塞的地方是 Phaser 类。 (这个类是一个像互斥锁和锁存器一样的同步屏障,但更灵活和强大。)因此在ForkJoinPool 任务中与Phaser 同步应该使用托管阻塞并且可以避免饥饿死锁(但ForkJoinTask.join() 仍然是可取的,因为它使用工作窃取而不是托管阻塞)。无论您是直接使用ForkJoinPool 还是通过其ExecutorService 接口,这都有效。但是,如果您使用任何其他 ExecutorService(例如 Executors 类创建的那些),它将不起作用,因为它们不支持托管阻塞。

在 Scala 中,托管阻塞的使用更为广泛(descriptionAPI)。

【讨论】:

  • 感谢您的回答,非常全面。不过,挑剔的是,ForkJoinTask 实现在 get() 中的窃取与在 join() 中的窃取相同。我的问题中的死锁主要来自尝试在没有ForkJoinPool.managedBlock() 的情况下进行同步(实际上,这两个示例都在 Java 7 上死锁)。改用Phasers,两者都可以工作。
  • 也许你可以澄清Java's thread's blocking state 是如何对应现代 Scala 意义上的阻塞线程的。是否每个阻塞操作都等待该链接上描述的 Java 监视器锁?还是线程池以其他方式跟踪阻塞状态?
  • @matt 我不确定您所说的“现代 Scala 意义上的阻塞线程”是什么意思。 State.BLOCKED 和托管阻塞之间的联系只是一个知道它可能很快处于 BLOCKED 状态的线程应该通过调用 managedBlock 提前告诉 ForkJoinPool 这件事。 ForkJoinPool 不会尝试检测线程当前是否被阻塞。如果一个线程使用managedBlock 但实际上并没有被阻塞,那么 ForkJoinPool 仍然会增加线程数。
  • 好的。一些描述暗示 ForkJoinPool 可能会或可能不会产生一个新线程,所以我想知道它是否有一些更深层次的决定智慧。我猜根据你,不是。在现代意义上,有些人将其称为“反应式编程”,对于面向用户的应用程序,也可以将长时间的 CPU 密集型任务视为阻塞。我正是这个意思。我认为您的回答也适用于这种情况。只关心除OufOfMemoryError 之外的线程数没有限制。这似乎带来了麻烦——为什么没有上限?
  • @matt 正如我在回答中指出的那样,ForkJoinPool 并不是为了对单个任务进行“反应”,而是为了尽快交付所有任务的最终结果可能,并且为了实现这一点,它接受为了更大的利益而延迟单个任务(在偷工作时)。
【解决方案2】:

我知道你在做什么,但我不知道为什么。屏障的想法是独立的线程可以等待彼此到达一个共同点。你没有独立的线程。线程池,F/J,用于Data Parallelism

你正在做一些更符合Task Parallelism的事情

F/J 继续的原因是框架创建“继续线程”以在所有工作线程都在等待时继续从双端队列获取工作。

【讨论】:

  • 障碍只是为了确保每个任务都安排在单独的线程上。而且我不认为“继续线程”是答案,如果您打印出Thread.currentThread().getId(),您会看到ForkableTasks 之一与调用其余线程的线程在同一线程中运行,并且只有4总共使用了线程。
  • 你不能保证哪个线程通过工作窃取来处理什么任务。所有任务都转储到同一个提交队列中。根据您使用的版本(Java7/8),阻塞的工作线程被替换为“继续”或“补偿”线程。你做的不是F/J(数据并行)的强项。
【解决方案3】:

您将实施细节与合同保证混淆了。您在文档中的哪个位置发现join 会窃取工作并防止在您的情况下出现死锁?相反,文档说:

方法 join() 及其变体仅适用于完成依赖是非循环的;也就是说,并行计算可以描述为有向无环图(DAG)。否则,由于任务循环等待,执行可能会遇到死锁。

您的示例是循环的。调用barrier.await() 的任务相互依赖。

文档进一步指出:

但是,此框架支持其他方法和技术(例如使用 Phaser、helpQuiesce() 和 complete(V)),这些方法和技术可用于为非静态结构化为 DAG 的问题构建自定义子类。

Phaser 状态的文档:

Phasers 也可以被在 ForkJoinPool 中执行的任务使用。如果池的 parallelismLevel 可以容纳同时被阻止的最大方数,则可以确保进度。

这仍然不清楚(因为它没有明确描述与 join 的交互),但这可能意味着 Phaser 旨在像您的第一个示例一样工作。

【讨论】:

    猜你喜欢
    • 2021-06-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-06-03
    • 2023-01-19
    相关资源
    最近更新 更多