【问题标题】:Which ThreadPool in Java should I use?我应该使用 Java 中的哪个 ThreadPool?
【发布时间】:2010-07-15 07:50:20
【问题描述】:

有大量的任务。 每个任务都属于一个组。要求是每组任务应该像在单个线程中执行一样串行执行,并且应该在多核(或多cpu)环境中最大化吞吐量。注意:还有大量的组,与任务的数量成正比。

天真的解决方案是使用 ThreadPoolExecutor 并同步(或锁定)。但是线程会互相阻塞,吞吐量没有最大化。

有更好的主意吗?或者是否存在满足要求的第三方库?

【问题讨论】:

  • “但是,线程会互相阻塞,吞吐量没有最大化。”。您的意思是各个任务正在访问共享数据结构或资源,这就是争用的原因吗?
  • 你提前知道小组的所有任务吗?这在选择解决方案时很重要(队列与无队列)

标签: java scala threadpool actor akka


【解决方案1】:

一种简单的方法是将所有组任务“连接”成一个超级任务,从而使子任务连续运行。但这可能会导致其他组延迟,除非其他组完全完成并在线程池中腾出一些空间,否则这些组将无法启动。

作为替代方案,可以考虑将小组的任务链接起来。下面的代码说明了这一点:

public class MultiSerialExecutor {
    private final ExecutorService executor;

    public MultiSerialExecutor(int maxNumThreads) {
        executor = Executors.newFixedThreadPool(maxNumThreads);
    }

    public void addTaskSequence(List<Runnable> tasks) {
        executor.execute(new TaskChain(tasks));
    }

    private void shutdown() {
        executor.shutdown();
    }

    private class TaskChain implements Runnable {
        private List<Runnable> seq;
        private int ind;

        public TaskChain(List<Runnable> seq) {
            this.seq = seq;
        }

        @Override
        public void run() {
            seq.get(ind++).run(); //NOTE: No special error handling
            if (ind < seq.size())
                executor.execute(this);
        }       
    }

优点是没有使用额外的资源(线程/队列),并且任务的粒度比幼稚方法中的要好。缺点是应该提前知道所有组的任务

--编辑--

为了使这个解决方案通用和完整,您可能需要决定错误处理(即即使发生错误,链是否继续),实现 ExecutorService 并将所有调用委托给底层执行者。

【讨论】:

  • 也许我们还应该添加一个地图,以便我们可以找到指定Task的TaskChain并将其添加到它的TaskChain中。
  • @James:你是对的。通过一些简单的同步,链将能够即时接收新任务(它们实际上将充当调用者的队列)。我想我会让解决方案更加通用和有用,并在我的博客中写下它:)
【解决方案2】:

我建议使用任务队列:

  • 对于每组任务,您都创建了一个队列并将该组中的所有任务插入其中。
  • 现在您的所有队列都可以并行执行,而一个队列中的任务可以串行执行。

快速谷歌搜索表明 java api 本身没有任务/线程队列。但是,有很多关于编码的教程。如果您知道一些,每个人都可以随意列出好的教程/实现:

【讨论】:

  • 谢谢戴夫。如果组的数量很大,那么线程的数量就会达到限制。
  • @James 不一定。仅仅因为您有 n 个组并不意味着您需要创建 n 个线程来执行它们。只需创建您认为合适的尽可能多的线程,它们就会以循环方式或串行方式处理队列。
【解决方案3】:

我基本同意 Dave 的回答,但如果您需要跨所有“组”划分 CPU 时间,即所有任务组应该并行进行,您可能会发现这种构造很有用(使用删除作为“锁定”。这在我的情况下工作得很好,虽然我想它往往会使用更多的内存):

class TaskAllocator {
    private final ConcurrentLinkedQueue<Queue<Runnable>> entireWork
         = childQueuePerTaskGroup();

    public Queue<Runnable> lockTaskGroup(){
        return entireWork.poll();
    }

    public void release(Queue<Runnable> taskGroup){
        entireWork.offer(taskGroup);
    }
 }

 class DoWork implmements Runnable {
     private final TaskAllocator allocator;

     public DoWork(TaskAllocator allocator){
         this.allocator = allocator;
     }

     pubic void run(){
        for(;;){
            Queue<Runnable> taskGroup = allocator.lockTaskGroup();
            if(task==null){
                //No more work
                return;
            }
            Runnable work = taskGroup.poll();
            if(work == null){
                //This group is done
                continue;
            }

            //Do work, but never forget to release the group to 
            // the allocator.
            try {
                work.run();
            } finally {
                allocator.release(taskGroup);
            }
        }//for
     }
 }

然后您可以使用最佳线程数来运行DoWork 任务。这是一种循环负载平衡..

您甚至可以做一些更复杂的事情,通过使用它而不是 TaskAllocator 中的简单队列(剩余任务较多的任务组往往会被执行)

ConcurrentSkipListSet<MyQueue<Runnable>> sophisticatedQueue = 
    new ConcurrentSkipListSet(new SophisticatedComparator());

SophisticatedComparator 在哪里

class SophisticatedComparator implements Comparator<MyQueue<Runnable>> {
    public int compare(MyQueue<Runnable> o1, MyQueue<Runnable> o2){
        int diff = o2.size() - o1.size();
        if(diff==0){
             //This is crucial. You must assign unique ids to your 
             //Subqueue and break the equality if they happen to have same size.
             //Otherwise your queues will disappear...
             return o1.id - o2.id;
        }
        return diff;
    }
 }

【讨论】:

  • +1 任务队列允许您使用任何适合您需要的调度算法。
  • 看起来你正在重新实现一个线程池。为什么不使用标准的 ThreadPoolExecutor 以及我的解决方案中的一些额外功能?我的解决方案不需要队列,也不需要同步。
  • @Eyal:如果可以按顺序使用任务组,我同意你的看法。但是,如果它们必须并行消耗,这是必要的。
  • 在我的解决方案中,组是并行执行的,每个组是串行执行的,就像在您的解决方案中一样。我们的解决方案之间的最大区别(如果我理解正确的话)是您的解决方案允许动态向现有组添加新任务,而我的解决方案要简单得多,因为它假设每当一个组开始执行时,它的所有任务都是已知的前进。
  • 哦,好吧,这就是你重新提交的原因.. 聪明 :) 我想 TPE 只允许 BlockingQueue 的事实可能会受到限制,但现在我明白你的意思了..
【解决方案4】:

Actor 也是这种特定类型问题的另一种解决方案。 Scala 有演员,也有 Java,由 AKKA 提供。

【讨论】:

    【解决方案5】:

    我遇到了与您类似的问题,我使用了与 Executor 配合使用的 ExecutorCompletionService 来完成任务集合。 这是从 Java7 开始的 java.util.concurrent API 的摘录:

    假设您有一组求解某个问题的求解器,每个求解器都返回某种类型的值 Result,并希望同时运行它们,以某种方法处理每个返回非空值的结果使用(结果 r)。你可以这样写:

    void solve(Executor e, Collection<Callable<Result>> solvers)
            throws InterruptedException, ExecutionException {
        CompletionService<Result> ecs = new ExecutorCompletionService<Result>(e);
        for (Callable<Result> s : solvers)
            ecs.submit(s);
        int n = solvers.size();
        for (int i = 0; i < n; ++i) {
            Result r = ecs.take().get();
            if (r != null)
                use(r);
        }
    }
    

    因此,在您的场景中,每个任务都将是一个 Callable&lt;Result&gt;,并且任务将被分组在一个 Collection&lt;Callable&lt;Result&gt;&gt; 中。

    参考: http://docs.oracle.com/javase/7/docs/api/java/util/concurrent/ExecutorCompletionService.html

    【讨论】:

      猜你喜欢
      • 2014-03-25
      • 2011-06-25
      • 1970-01-01
      • 1970-01-01
      • 2016-05-19
      • 2010-10-18
      • 2016-05-17
      • 1970-01-01
      • 2016-09-28
      相关资源
      最近更新 更多