【问题标题】:CompletableFuture with PriorityBlockingQueue带有 PriorityBlockingQueue 的 CompletableFuture
【发布时间】:2020-02-23 10:36:29
【问题描述】:

我正在尝试编写一个单线程执行器,当任务被安排并根据PriorityBlockingQueue 执行任务时返回CompletableFuture

我的任务如下所示:

  public interface DriverTask<V> {

    V call(WebDriver driver) throws Throwable;

    default TaskPriority getPriority() {
      return TaskPriority.LOW;
    }
  }

  public enum TaskPriority {
    HIGH,
    MEDIUM,
    LOW
  }

现在我的问题是,当我使用 CompletableFuture.supplyAsync 方法时,Executor 只得到一个 Runnable,我不知道如何让我的 Executor 知道原始任务的优先级。

是否有其他方法可以创建 CompletableFuture 以便我根据优先级执行它们?

【问题讨论】:

  • 我可以在你自己的 ExecutorService 中简单地完成一个 CompletableFuture。只需在您的 ExecutorService 中执行工作并调用 .complete() 这样您就可以编写基于优先级工作的代码......
  • @Alex 我会通过构造函数创建未来,然后在执行时调用 complete() 吗?或者还有什么我需要记住的。
  • 还有其他障碍,比如投掷Throwable 或需要WebDriver,这使得supplyAsync 不可行。谁应该提供参数,调用者还是执行者?
  • 所以WebDriver 是由调用者提供的,调用者也提供了DriverTask,或者它是范围内的某个(全局)变量?
  • 所以执行者知道WebDriver。我可以假设,例如为了示例解决方案的方法currDriver()

标签: java java-8 concurrency


【解决方案1】:

supplyAsync 等方法的原理是创建一个新的CompletableFuture 实例,这是一个设置和异步作业,最终将在未来complete

您可以为您的重要设置做同样的事情:

private static WebDriver currentDriver() {
    …
}
private static final ExecutorService BACKEND
    = new ThreadPoolExecutor(1, 1, 1, TimeUnit.MINUTES, new PriorityBlockingQueue<>());

public static <V> CompletableFuture<V> runAsync(DriverTask<V> dt) {
    CompletableFuture<V> result = new CompletableFuture<>();
    class Job implements Runnable, Comparable<Job>,
                         CompletableFuture.AsynchronousCompletionTask {
        public void run() {
            try {
                if(!result.isDone()) result.complete(dt.call(currentDriver()));
            }
            catch(Throwable t) { result.completeExceptionally(t); }
        }
        private TaskPriority priority() { return dt.getPriority(); }
        public int compareTo(Job o) { return priority().compareTo(o.priority()); }
    }
    BACKEND.execute(new Job());
    return result;
}

请注意,不需要实现CompletableFuture.AsynchronousCompletionTask;这只是标记那些旨在完成CompletableFutureRunnable 实现的约定。

自己实现逻辑的另一个优点是第一阶段不需要将异常包装在CompletionException 中。因此,当调用者链接exceptionally 时,它将看到原始的未包装异常。此外,join 的调用者会得到一个 CompletionException,它反映了以原始异常为原因的 join 调用的代码位置,包含更多有用的信息。

if(!result.isDone()) 在实际完成尝试之前的目的是如果在队列中等待时CompletableFuture 已被取消(或以其他方式完成),则跳过它。一旦开始完成尝试,取消不会中断它。这是CompletableFuture 的一般行为。

【讨论】:

  • 感谢您抽出这么多时间来帮助我。我按照您在此处显示的方式实现了它,但我仍在努力完全理解它。由于ThreadPoolExecutor 的执行方法只需要一个Runnable,执行程序如何确保执行顺序实际上是基于优先级的?底层的 PriorityBlockingQueue 是否只是假设 Runnable 也是 Comparable?没有类型安全问题吗?
  • 是的,PriorityBlockingQueue 假设有可比较的元素。这是类型系统的一个弱点,当使用默认构造函数时,我们对所有排序集合(TreeSet 等)都有,换句话说,既没有指定比较器也没有指定初始值。请注意,这与ScheduledThreadPoolExecutor​ 的工作方式相同,只是我们的优先级不是基于时间的。
猜你喜欢
  • 1970-01-01
  • 2015-11-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-06-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多