【问题标题】:Scheduling Tasks to Execute Sequentially After a Given Delay安排任务在给定延迟后按顺序执行
【发布时间】:2017-11-23 13:42:10
【问题描述】:

我有一堆工作线程,我想在给定延迟后按顺序执行。我想实现以下行为:

延迟 -> 工人 1 -> 延迟 -> 工人 2 - 延迟 -> 工人 3 -> ...

我想出了这个解决方案:

long delay = 5;
for(String value : values) {
    WorkerThread workerThread = new WorkerThread(value);
    executorService.schedule(workerThread, delay, TimeUnit.SECONDS);
    delay = delay + 5;
}

executorService 的创建位置如下:

private final ScheduledExecutorService executorService = Executors.newSingleThreadScheduledExecutor();

还有没有其他方法可以在 Java 中使用 ExecutorService 来实现这一点?

【问题讨论】:

  • 您可以有一个工作线程,它每delay 秒执行一次,并从队列中获取value 进行处理。
  • 您当前的解决方案相对于调度的开始延迟了每个工人的开始,而您的问题布局建议在一名工人完成与您下一次开始之间存在固定延迟。你能澄清一下你真正想要的吗?
  • @bowmore 我的目标是模拟状态机。所以,我会在工作人员和下一次开始之间有一个固定的延迟。

标签: java multithreading executorservice


【解决方案1】:

看到您的问题,我想出了另一个解决方案。假设值是一个可以更改的队列。这是一个有效的解决方案。我稍微修改了您的 WorkerThread 并在其中添加了一个回调对象。希望这会有所帮助。

private final Queue<String> values = new LinkedList<>();

private final ScheduledExecutorService executorService = Executors.newSingleThreadScheduledExecutor();

private void start() {
    AtomicLong delay = new AtomicLong(5);
    String value = values.poll();

    if (value != null) {
        WorkerThread workerThread = new WorkerThread(value, new OnCompleteCallback() {
            @Override
            public void complete() {
                String valueToProcessNext = values.poll();

                if (valueToProcessNext != null) {
                    executorService.schedule(new WorkerThread(valueToProcessNext, this), delay.addAndGet(5), TimeUnit.SECONDS);
                }
            }
        });
        executorService.schedule(workerThread, delay.get(), TimeUnit.SECONDS);
    }
}

class WorkerThread implements Runnable {

    private final String value;

    private final OnCompleteCallback callback;

    WorkerThread(String value, OnCompleteCallback callback) {
        this.value = value;
        this.callback = callback;
    }

    @Override
    public void run() {
        try {
            System.out.println(value);
        } finally {
            callback.complete();
        }
    }
}

interface OnCompleteCallback {
    void complete();
}

【讨论】:

    【解决方案2】:

    如果应该使用ExecutorService,除了您的解决方案之外,什么都不会出现。但是您可能会发现 CompletableFuture 更有用,因为它提供了类似的行为,但相对于任务完成有延迟,而不是开始调度。

    CompletableFuture<Void> completableFuture = CompletableFuture.completedFuture(null);
    String[] values = new String[]{"a", "b", "c"};
    for (String value : values) {
        completableFuture
                .thenRun(() -> {
                    try {
                        Thread.sleep(5000);
                    } catch (InterruptedException e) {
    
                    }
                })
                .thenRun(() -> System.out.println(value));
    }
    completableFuture.get();
    

    【讨论】:

      【解决方案3】:

      您可以在每个工作人员之间使用DelayQueue。并用这个类装饰你的工人:

      public class DelayedTask implements Runnable {
      
          private final Runnable task;
          private final DelayQueue<Delayed> waitQueue;
          private final DelayQueue<Delayed> followerQueue;
      
          public DelayedTask(Runnable task, DelayQueue<Delayed> waitQueue, DelayQueue<Delayed> followerQueue) {
              this.task = Objects.requireNonNull(task);
              this.waitQueue = Objects.requireNonNull(waitQueue);
              this.followerQueue = followerQueue;
          }
      
          @Override
          public void run() {
              try {
                  waitQueue.take();
                  try {
                      task.run();
                  } finally {
                      if (followerQueue != null) {
                          followerQueue.add(new Delay(3, TimeUnit.SECONDS));
                      }
                  }
              } catch (InterruptedException e) {
                  Thread.currentThread().interrupt();
              }
          }
      }
      

      还有一个简单的延迟实现

      class Delay implements Delayed {
          private final long nanos;
      
          Delay(long amount, TimeUnit unit) {
              this.nanos = TimeUnit.NANOSECONDS.convert(amount, unit) + System.nanoTime();
          }
      
          @Override
          public long getDelay(TimeUnit unit) {
              return unit.convert(nanos - System.nanoTime(), TimeUnit.NANOSECONDS);
          }
      
          @Override
          public int compareTo(Delayed other) {
              return Long.compare(nanos, other.getDelay(TimeUnit.NANOSECONDS));
          }
      }
      

      允许这种用法:

          ExecutorService executorService = Executors.newFixedThreadPool(1);
      
      // ....
      
          DelayQueue<Delayed> currentQueue = new DelayQueue<>();
          currentQueue.add(new Delay(3, TimeUnit.SECONDS));
          for (String value : values) {
              DelayedTask delayedTask = new DelayedTask(new WorkerThread(value), currentQueue, currentQueue = new DelayQueue<>());
              executorService.submit(delayedTask);
          }
      

      【讨论】:

        猜你喜欢
        • 2013-11-11
        • 2014-05-10
        • 1970-01-01
        • 2015-09-14
        • 2013-04-29
        • 2015-07-24
        • 2014-09-15
        • 1970-01-01
        • 2016-12-01
        相关资源
        最近更新 更多