【问题标题】:Java fixed size Thread Pool and optimal usage of all CPU coresJava 固定大小的线程池和所有 CPU 内核的最佳使用率
【发布时间】:2014-04-12 11:59:30
【问题描述】:

如何始终将 8 个线程用于“昂贵”的部分?

我有一个数字运算问题,为此我创建了一个简单的框架。我的问题是找到一种优雅而简单的方法来优化使用所有 CPU 内核。

为了获得良好的性能,我使用了一个固定大小为 8 的线程池。这个想法是使用与硬件线程一样多的线程来获得最佳性能。

框架的简化伪代码用法如下:

interface Task {
  data[] compute(data[]);
}

Task task = new Loop(new Chain(new DoX(), new DoY(), new Split(2, new DoZ())));
result = task.compute(data);
  • 循环任务将循环直到满足某些终止条件
  • Chain Task 将链式任务(例如上面的 r = t1.compute(r); r = t2.compute(r); r = t3.compute(r); return r;)
  • Split Task 会拆分数据并对部分执行任务(例如,创建 2 个部分并返回新数据[] {t1.compute(part1), t1.compute(part2)})

目前线程是在Split Task中实现的。所以拆分任务会将 t1.compute(part1) 和 t1.compute(part2) 的计算交给线程池。

方法一,可能完全死锁

我的第一种方法是拆分任务有一个期货数组,并一个接一个地调用 get()。但这意味着如果拆分任务在另一个拆分任务中,future.get() 中的阻塞等待将阻塞外部拆分任务从线程池中获取的线程。所以我真正工作的线程不到 8 个。如果这个层次很深,我可能没有人工作,永远等待。

1) 我假设future.get() 不会将线程返回到线程池,对吧?因此,如果这样做,我将在 future.get() 中等待,但没有更多的线程开始工作? [我不能轻易测试,因为我已经改变了方法]

方法 2,目前的方法,至少有人在工作

我目前的方法(好不了多少)是用当前线程进行拆分的最后一部分(partN)。如果完成,我检查 partN-1 是否已经启动,如果是,我会等待 future.get() 中的所有任务,否则当前线程也会执行 partN-1,如果需要 partN-2 ... 所以现在我应该总是在池中至少有一个线程在工作。

但由于问题 1) 的答案可能是 future.get() 会阻塞我的线程,使用这种方法我将只有很少的工作线程在深层层次结构上。

方法3,我看到的唯一解决方案

我假设我必须使用 2 个线程池,一个用于辛勤工作,一个用于所有等待。因此,我将有一个固定大小的线程池用于辛勤工作,而(动态?)一个用于等待。

3.a.:但这意味着拆分任务只能从等待池中产生线程,而执行实际工作的任务将从工作池中产生一个新线程并等待它完成。丑陋,但应该工作。丑陋,因为目前整个线程支持都在拆分任务中,但是使用此解决方案,其他正在努力工作的任务必须了解线程。

3.b.:另一种方法是 Split 产生工作线程,但在 split 内部,每个等待必须由等待线程完成,而当前线程同时也执行工作线程任务。有了这个,所有线程支持都在拆分任务类中,但我不确定如何实现。

2a)如何在不阻塞当前线程的情况下等待任务?

2b) 是否可以将当前线程返回到工作线程池,让等待者线程等待,然后在等待后继续上一个当前线程或工作线程池中的线程?怎么样?

其他解决方案

不要使用固定大小的线程池。

3) 我的想法有 8 个线程是错误的吗?但是,如果层次结构可以很深,那么有多少呢? JVM 并行启动许多任务并在它们之间进行大量切换是不是存在风险?

4)我错过了什么或者你会怎么解决这个问题?

非常感谢和问候


[编辑]

接受的解决方案以及为什么我尝试不同的方法(基于方法 2)

我接受了 ForkJoinPool 作为正确的解决方案。

但是,一些细节以及可能的开销和失控让我想尝试另一种方法。但是我想得越多,我就越会回到使用 ForkJoinPool (原因见最后的注释)。很抱歉文字量太大。

http://docs.oracle.com/javase/7/docs/api/java/util/concurrent/ForkJoinPool.html

“但是,面对阻塞的 IO 或其他非托管同步,不能保证此类调整。”

"最大运行线程数为32767"

http://homes.cs.washington.edu/~djg/teachingMaterials/grossmanSPAC_forkJoinFramework.html

“ForkJoin 框架的文档建议创建并行子任务,直到基本计算步骤的数量超过 100 且小于 10,000。”

“艰苦的工作”任务从磁盘读取大量数据,与 10,000 次基本计算相去甚远。实际上,我可以将其 fork/join 降低到可接受的水平,但现在这工作量太大了,因为那部分代码相当复杂。

我认为方法 3a 基本上是 ForkJoin 的一种实现,除了我会有更多的控制权并且可能会更少的开销并且上面提到的问题不应该存在(但不会自动适应操作系统提供的 CPU 资源,但我会强制操作系统可以在必要时给我我想要的东西)。

我可能会尝试使用方法 2 并进行一些更改:这样我可以使用确切的线程号并且我没有任何等待线程,如果我理解正确,ForkJoinPool 似乎可以使用等待线程。

当前线程执行作业,直到此 Split 实例中的所有作业都由工作线程运行(因此像以前一样在 Split 节点中窃取工作),但随后它不会调用 future.get(),而只是检查是否所有期货已准备好使用 future.isDone()。如果没有全部完成,它将从线程池中窃取一个作业并执行它,然后再次检查期货。这样,只要有一个作业没有运行,我就永远不会等待。

丑陋的:如果没有工作可以偷,我将不得不睡一小会,然后再次检查期货或从池中偷一个新工作(有没有办法等待多个期货全部完成超时不会取消计算,如果它触发?)

所以我认为我必须在每个拆分任务中为 ThreadPool 使用 Completion Service,然后我可以使用超时轮询并且不需要休眠。

假设:完成服务中的线程池仍然可以像普通线程池一样使用(例如作业窃取)。一个 ThreadPool 可以在多个 Completion Services 中。

我认为这是问题中详述的问题的最佳解决方案。但是,这有一个小问题,请参阅以下内容。

注意:

再次查看“硬”任务后,我发现它们可以针对许多实例进行并行化。因此,在那里添加线程也是下一个合乎逻辑的步骤。这些始终是叶节点,并且它们所做的工作最好通过完成服务完成(在某些情况下,子作业可以有不同的运行时,但任何两个结果都可以构建一个新作业)。要使用 ForkJoinPool 完成它们,我必须使用 managedBlock() 并实现 ForkJoinPool.ManagedBlocker,这会使代码更加复杂。但是,与此同时,在这些离开节点中使用 CompletionService 意味着我基于方法 2 的解决方案可能也需要等待线程,所以我最好使用 ForkJoinPool。

【问题讨论】:

  • 如果你有一个任务依赖于n 其他任务,你总是可以共享一个CountDownLatchn 作为这个依赖任务之间的构造函数参数(这将.await())和其他人(这将是.countDown())。
  • 谢谢。但我认为这与我的问题无关。

标签: java concurrency threadpool future


【解决方案1】:

我不得不离开ForkJoinPool,它没有以最佳方式使用线程。 虽然它对 Loop 和 Split 节点工作得很好,但如果我想并行化实际工作发生的叶节点,它就不再工作了。当我将它们添加为RecursiveTask 时,大多数线程都处于空闲状态。出于某种原因,join() 调用不会窃取叶子中的工作(jdk1.7.0_45)。它正在等待。在我的情况下,所有工作都在叶子中,因此为叶子使用自定义 RecursiveTask 子类比仅将其用于 Loop 和 Split 节点更糟糕(因为它在部分工作之后等待,否则它在所有工作之后等待工作)。我不认为我用ForkJoinPool 错了,如果你用谷歌搜索你会发现有类似问题的人。

我现在做了一个简单的解决方案:2 个线程池,1 个用于实际工作的固定大小,以及一个用于所有 Loop 和 Split 节点的缓存。我创建了FakeRecursiveTask(扩展这个而不是原来的),所以我不必更改代码(循环和拆分)。我使用HardWork 作为叶子的基类,很明显它是不同的,只需调用doHardWork(work)

有了这个解决方案,我所有的工作线程一直都被完全使用。由于树的大小有限,我永远不应该用完辅助线程。实际上,在我的情况下,它主要使用与工作线程相同数量的辅助线程(在我的情况下是 8 个)。

public class ThreadPool3 {
        private static int maxNumWorkerThreads;
        private static ExecutorService workerPool = null;
        private static ExecutorService helperPool = null;

        public static void initThreadPool(int maxNumWorkerThreads_) {
                int availProcessors = Runtime.getRuntime().availableProcessors();
                if (maxNumWorkerThreads_ <= 0) {
                        maxNumWorkerThreads_ = availProcessors;
                }
                maxNumWorkerThreads = maxNumWorkerThreads_;

                if (availProcessors != maxNumWorkerThreads) {
                        System.out.println("WARN: maxNumWorkerThreads (" + maxNumWorkerThreads + ") != availProcessors (" + availProcessors + ")");
                }
                workerPool = Executors.newFixedThreadPool(maxNumWorkerThreads);
                BlockingQueue<Runnable> workQueue = new SynchronousQueue<Runnable>();
                helperPool = new ThreadPoolExecutor(0, 4 * maxNumWorkerThreads, 60, TimeUnit.MINUTES, workQueue, Executors.defaultThreadFactory(),
        new ThreadPoolExecutor.CallerRunsPolicy());
        }


        public static abstract class HardWork implements Callable<Void> {
                @Override
                public abstract Void call() throws Exception;
        }

        public static void doHardWork(List<HardWork> tasks) throws Exception {
                workerPool.invokeAll(tasks);
        }


        /**
        * fake ForkJoinPoolInterface:
        *
        */
        public static abstract class FakeRecursiveTask<T> implements Callable<T> {
                private Future<T> resultFuture = null;

                /**
                * fake interface:
                */
                public abstract T compute();

                /**
                * fake interface:
                */
                public T invoke() {
                        return compute();
                }

                /**
                * fake interface:
                */
                public void fork() {
                        resultFuture = helperPool.submit(this);
                }

                /**
                * fake interface:
                */
                public T join() {
                        try {
                                return resultFuture.get();
                        }
                        catch (Exception e) {
                                throw new RuntimeException(e);
                        }
                }

                @Override
                public T call() throws Exception {
                        return compute();
                }
        }


        public static void shutdownThreadPool() {
                if (workerPool != null) {
                        workerPool.shutdown();
                }
                if (helperPool != null) {
                        helperPool.shutdown();
                }
        }
}

【讨论】:

    【解决方案2】:

    为了完全避免死锁,不要使用同步的 Future.get()。改用异步方法 CompletableFuture.then 和 CompletableFuture.both,在 Java8 中可用。这些方法不会阻塞,而是在数据可用时提交新任务。如果您不想使用 Java8,请查看 Guava 库,它(我相信)具有等效的功能。存在其他异步库,例如我的https://github.com/rfqu/df4j。它的优点是可以重复使用任务对象,因此必须创建较少数量的对象。如果您提供更详细的问题描述(例如,以普通的顺序形式,或使用无限数量的线程),我可以帮助您使用 df4j 实现您的程序。

    【讨论】:

    • 谢谢。但我认为,如果只是在拆分任务中使用它,它就像问题中的方法 2。因为这也是“异步”的,所以它在等待之前会做尽可能多的工作。如果我理解正确,这与此解决方案相同,如果异步部分完成,则必须等待(我认为方法 2 的等待时间更短,因为它窃取了任务)。但 CompletableFuture 似乎很复杂,所以也许我没有明白。
    【解决方案3】:

    您似乎有一个并行的“分而治之”类型的问题,您正在递归地将问题拆分为要使用可用内核“解决”的子问题。

    您是正确的,创建线程的 niave 实现可能会使用大量资源,并且使用有界线程池很可能会死锁。

    第三种选择是在 Java 7 中实现的“fork/join”模型。这在 Oracle Java 教程 (here) 中有描述,但我认为 Dan Grossman 的讲义更好地解释了它:

    【讨论】:

    • 谢谢,我认为这在理论上是正确的解决方案。请参阅对我的问题的编辑,为什么我可能会尝试使用其他解决方案。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-08-02
    • 2022-01-23
    • 2021-10-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多