代码之家  ›  专栏  ›  技术社区  ›  Pavel Feldman

等待,直到任何未来<T>完成

  •  53
  • Pavel Feldman  · 技术社区  · 18 年前

    我运行的异步任务很少,需要等待至少一个任务完成(将来可能需要等待N个任务中的M个任务完成)。 目前,它们被表示为未来,所以我需要

    /**
     * Blocks current thread until one of specified futures is done and returns it. 
     */
    public static <T> Future<T> waitForAny(Collection<Future<T>> futures) 
            throws AllFuturesFailedException
    

    有这样的吗?或任何类似的,对未来不必要的。目前,我循环收集期货,检查一个期货是否完成,然后睡眠一段时间,然后再次检查。这看起来不是最好的解决方案,因为如果我长时间睡眠,那么会增加不必要的延迟,如果我短时间睡眠,那么会影响性能。

    我可以尝试使用

    new CountDownLatch(1)
    

    countdown.await()
    

    <T> RunnableFuture<T> AbstractExecutorService.newTaskFor(Callable<T> callable)
    

    并创建RunnableFuture的自定义实现,能够在任务完成时附加要通知的侦听器,然后将此类侦听器附加到所需的任务并使用CountDownLatch,但这意味着我必须为我使用的每一个ExecutorService重写newTaskFor,而且可能会有一些实现不扩展AbstractExecutorService。我也可以尝试为同样的目的包装服务,但我必须装饰所有产生未来的方法。

    WaitHandle.WaitAny(WaitHandle[] waitHandles)
    

    在c#中。对于这类问题有什么众所周知的解决方案吗?

    更新:

    7 回复  |  直到 18 年前
        1
  •  55
  •   pingw33n    13 年前
        3
  •  7
  •   Alex Miller    18 年前

    为什么不创建一个结果队列并在队列上等待呢?或者更简单地说,使用CompletionService,因为它是:ExecutorService+结果队列。

        4
  •  6
  •   Scott Stanchfield    18 年前

    使用wait()和notifyAll()实际上非常简单。

    package com.javadude.sample;
    
    public class Lock {}
    

    接下来,定义工作线程。他必须在完成处理后通知锁对象。请注意,通知必须位于锁定对象上的同步块锁定中。

    package com.javadude.sample;
    
    public class Worker extends Thread {
        private Lock lock_;
        private long timeToSleep_;
        private String name_;
        public Worker(Lock lock, String name, long timeToSleep) {
            lock_ = lock;
            timeToSleep_ = timeToSleep;
            name_ = name;
        }
        @Override
        public void run() {
            // do real work -- using a sleep here to simulate work
            try {
                sleep(timeToSleep_);
            } catch (InterruptedException e) {
                interrupt();
            }
            System.out.println(name_ + " is done... notifying");
            // notify whoever is waiting, in this case, the client
            synchronized (lock_) {
                lock_.notify();
            }
        }
    }
    

    最后,您可以编写您的客户机:

    package com.javadude.sample;
    
    public class Client {
        public static void main(String[] args) {
            Lock lock = new Lock();
            Worker worker1 = new Worker(lock, "worker1", 15000);
            Worker worker2 = new Worker(lock, "worker2", 10000);
            Worker worker3 = new Worker(lock, "worker3", 5000);
            Worker worker4 = new Worker(lock, "worker4", 20000);
    
            boolean started = false;
            int numNotifies = 0;
            while (true) {
                synchronized (lock) {
                    try {
                        if (!started) {
                            // need to do the start here so we grab the lock, just
                            //   in case one of the threads is fast -- if we had done the
                            //   starts outside the synchronized block, a fast thread could
                            //   get to its notification *before* the client is waiting for it
                            worker1.start();
                            worker2.start();
                            worker3.start();
                            worker4.start();
                            started = true;
                        }
                        lock.wait();
                    } catch (InterruptedException e) {
                        break;
                    }
                    numNotifies++;
                    if (numNotifies == 4) {
                        break;
                    }
                    System.out.println("Notified!");
                }
            }
            System.out.println("Everyone has notified me... I'm done");
        }
    }
    
        5
  •  4
  •   jdmichal    18 年前

    WaitHandle.WaitAny 方法

    public WaitableFuture<T>
        extends Future<T>
    {
        private CountDownLatch countDownLatch;
    
        WaitableFuture(CountDownLatch countDownLatch)
        {
            super();
    
            this.countDownLatch = countDownLatch;
        }
    
        void doTask()
        {
            super.doTask();
    
            this.countDownLatch.countDown();
        }
    }
    

    doTask() 方法但如果在执行之前无法控制未来的对象,我真的认为没有轮询就无法做到这一点。

    或者如果未来总是在它自己的线程中运行,并且你可以以某种方式获得该线程。然后,您可以生成一个新线程来连接其他线程,然后在连接返回后处理等待机制。。。这将是非常丑陋的,但会导致大量的开销。如果将来的一些对象没有完成,那么可能会有很多阻塞线程,这取决于死线程。如果不小心,这可能会泄漏内存和系统资源。

    /**
     * Extremely ugly way of implementing WaitHandle.WaitAny for Thread.Join().
     */
    public static joinAny(Collection<Thread> threads, int numberToWaitFor)
    {
        CountDownLatch countDownLatch = new CountDownLatch(numberToWaitFor);
    
        foreach(Thread thread in threads)
        {
            (new Thread(new JoinThreadHelper(thread, countDownLatch))).start();
        }
    
        countDownLatch.await();
    }
    
    class JoinThreadHelper
        implements Runnable
    {
        Thread thread;
        CountDownLatch countDownLatch;
    
        JoinThreadHelper(Thread thread, CountDownLatch countDownLatch)
        {
            this.thread = thread;
            this.countDownLatch = countDownLatch;
        }
    
        void run()
        {
            this.thread.join();
            this.countDownLatch.countDown();
        }
    }
    
        6
  •  1
  •   Alex - GlassEditor.com    5 年前

    如果你能用 CompletableFuture 那就有了 CompletableFuture.anyOf 这就是您想要的,只需在结果中调用join:

    CompletableFuture.anyOf(futures).join()
    

    你可以用 通过致电 CompletableFuture.supplyAsync 或 runAsync

        7
  •  0
  •   1800 INFORMATION    18 年前

    既然您不关心哪个线程完成,为什么不为所有线程使用一个WaitHandle并等待它呢?谁先完成谁就可以设置手柄。

        8
  •  -1
  •   Crowie Ralph    13 年前

    public class WaitForAnyRedux {
    
    private static final int POOL_SIZE = 10;
    
    public static <T> T waitForAny(Collection<T> collection) throws InterruptedException, ExecutionException {
    
        List<Callable<T>> callables = new ArrayList<Callable<T>>();
        for (final T t : collection) {
            Callable<T> callable = Executors.callable(new Thread() {
    
                @Override
                public void run() {
                    synchronized (t) {
                        try {
                            t.wait();
                        } catch (InterruptedException e) {
                        }
                    }
                }
            }, t);
            callables.add(callable);
        }
    
        BlockingQueue<Runnable> queue = new ArrayBlockingQueue<Runnable>(POOL_SIZE);
        ExecutorService executorService = new ThreadPoolExecutor(POOL_SIZE, POOL_SIZE, 0, TimeUnit.SECONDS, queue);
        return executorService.invokeAny(callables);
    }
    
    static public void main(String[] args) throws InterruptedException, ExecutionException {
    
        final List<Integer> integers = new ArrayList<Integer>();
        for (int i = 0; i < POOL_SIZE; i++) {
            integers.add(i);
        }
    
        (new Thread() {
            public void run() {
                Integer notified = null;
                try {
                    notified = waitForAny(integers);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                } catch (ExecutionException e) {
                    e.printStackTrace();
                }
                System.out.println("notified=" + notified);
            }
    
        }).start();
    
    
        synchronized (integers) {
            integers.wait(3000);
        }
    
    
        Integer randomInt = integers.get((new Random()).nextInt(POOL_SIZE));
        System.out.println("Waking up " + randomInt);
        synchronized (randomInt) {
            randomInt.notify();
        }
      }
    }