【问题标题】:Testing PriorityBlockingQueue in ThreadPoolExecutor在 ThreadPoolExecutor 中测试 PriorityBlockingQueue
【发布时间】:2013-05-25 21:32:38
【问题描述】:

我用 PriorityBlockingQueue 实现了我的 ThreadPoolExecutor,如下例所示: https://stackoverflow.com/a/12722648/2206775

并写了一个测试:

PriorityExecutor executorService = (PriorityExecutor)  PriorityExecutor.newFixedThreadPool(16);
    executorService.submit(new Runnable() {
        @Override
        public void run() {
            try {
                Thread.sleep(1000);
                Thread.sleep(1000);
                System.out.println("1");
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }, 1);

    executorService.submit(new Runnable() {
        @Override
        public void run() {
            try {
                Thread.sleep(1000);
                Thread.sleep(1000);
                System.out.println("3");
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }, 3);

    executorService.submit(new Runnable() {
        @Override
        public void run() {
            try {
                Thread.sleep(1000);
                Thread.sleep(1000);
                System.out.println("2");
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }, 2);

    executorService.submit(new Runnable() {
        @Override
        public void run() {
            try {
                Thread.sleep(1000);
                Thread.sleep(1000);
                System.out.println("5");
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }, 5);

    executorService.submit(new Runnable() {
        @Override
        public void run() {
            try {
                Thread.sleep(1000);
                Thread.sleep(1000);
                System.out.println("4");
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }, 4);

    executorService.shutdown();
    try {
        executorService.awaitTermination(30, TimeUnit.MINUTES);
    } catch (InterruptedException e) {
        e.printStackTrace();
    }

但最后,我没有得到 1 2 3 4 5,我得到了这些数字的随机顺序。是考试有问题还是有其他问题?如果首先,如何正确测试?

【问题讨论】:

    标签: java multithreading threadpool priority-queue threadpoolexecutor


    【解决方案1】:

    你有 16 个线程,只有 5 个任务,这意味着它们都是并发执行的,优先级实际上无关紧要。

    只有当有任务等待执行时,优先级才重要。

    为了说明这一点,如果您将示例设置为仅使用 1 个线程,您将获得预期的输出。

    【讨论】:

      【解决方案2】:

      仅当池完全繁忙并且您提交了多个新任务时才会考虑优先级。如果你只用一个线程定义你的池,你应该得到预期的输出。在您的示例中,所有任务都是同时执行的,哪个先完成有点随机。

      顺便说一句,如果您的队列已满并且您提交了新任务,则链接的实现有问题并引发异常。

      请参阅下面的工作示例,说明您正在尝试实现的目标(我以一种简单的方式覆盖了 newTaskFor,只是为了使其工作 - 您可能想要改进该部分)。

      打印:1 2 3 4 5

      public class Test {
      
          public static void main(String[] args) {
              PriorityExecutor executorService = (PriorityExecutor) PriorityExecutor.newFixedThreadPool(1);
              executorService.submit(getRunnable("1"), 1);
              executorService.submit(getRunnable("3"), 3);
              executorService.submit(getRunnable("2"), 2);
              executorService.submit(getRunnable("5"), 5);
              executorService.submit(getRunnable("4"), 4);
      
              executorService.shutdown();
              try {
                  executorService.awaitTermination(30, TimeUnit.MINUTES);
              } catch (InterruptedException e) {
                  e.printStackTrace();
              }
          }
      
          public static Runnable getRunnable(final String id) {
              return new Runnable() {
                  @Override
                  public void run() {
                      try {
                          Thread.sleep(1000);
                          System.out.println(id);
                      } catch (InterruptedException e) {
                          e.printStackTrace();
                      }
                  }
              };
          }
      
          static class PriorityExecutor extends ThreadPoolExecutor {
      
              public PriorityExecutor(int corePoolSize, int maximumPoolSize,
                                      long keepAliveTime, TimeUnit unit, BlockingQueue<Runnable> workQueue) {
                  super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue);
              }
              //Utitlity method to create thread pool easily
      
              public static ExecutorService newFixedThreadPool(int nThreads) {
                  return new PriorityExecutor(nThreads, nThreads, 0L,
                                              TimeUnit.MILLISECONDS, new PriorityBlockingQueue<Runnable>());
              }
              //Submit with New comparable task
      
              public Future<?> submit(Runnable task, int priority) {
                  return super.submit(new ComparableFutureTask(task, null, priority));
              }
              //execute with New comparable task
      
              public void execute(Runnable command, int priority) {
                  super.execute(new ComparableFutureTask(command, null, priority));
              }
      
              @Override
              protected <T> RunnableFuture<T> newTaskFor(Callable<T> callable) {
                  return (RunnableFuture<T>) callable;
              }
      
              @Override
              protected <T> RunnableFuture<T> newTaskFor(Runnable runnable, T value) {
                  return (RunnableFuture<T>) runnable;
              }
          }
      
          static class ComparableFutureTask<T> extends FutureTask<T> implements Comparable<ComparableFutureTask<T>> {
      
              volatile int priority = 0;
      
              public ComparableFutureTask(Runnable runnable, T result, int priority) {
                  super(runnable, result);
                  this.priority = priority;
              }
      
              public ComparableFutureTask(Callable<T> callable, int priority) {
                  super(callable);
                  this.priority = priority;
              }
      
              @Override
              public int compareTo(ComparableFutureTask<T> o) {
                  return Integer.valueOf(priority).compareTo(o.priority);
              }
          }
      }
      

      【讨论】:

      • 你能解释一下,如果我们不覆盖两个 'newTaskFor' 方法,为什么链接的实现会抛出这个异常?
      • 因为默认实现将 ComparableFutureTask 包装到 RunnableFuture 中,并且当您的 PriorityQueue 尝试将 RunnableFuture 转换回 ComparableFutureTask 时,它会出现异常。 JDK 附带源代码,因此您可以在 IDE 中逐步调试此操作,以实时查看它。
      • 我今天花了两个小时来分析这个......如果你 100% 确定你不会打电话给ThreadPoolExecutor#submit(只有ThreadPoolExecutor#execute),你实际上可以跳过整个PriorityExecutorThreadPoolExecutor#execute 不会将Runnable 参数包装在任何东西中,这将使其在PriorityBlockingQueue 中具有可比性。根据您的ThreadPoolExecutor 的使用范围,这可能是一个可行的解决方案。不过,我建议添加对此的评论...最后但并非最不重要的一点是,AbstractExecutorService 是为此设计的!
      • 我试过这段代码,它总是运行第一个任务而不考虑优先级。例如,如果您首先提交优先级为 4 的任务,它将始终在开始时运行该任务。此外,如果将线程池大小增加到 2,即使有足够多的线程在等待,它也会根据插入顺序运行它们。
      • 我进行了改进以保留相同优先级的任务顺序:stackoverflow.com/a/42831172/1386911
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2021-06-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-05-02
      相关资源
      最近更新 更多