【问题标题】:completionservice: how to kill all threads and return result through 5 seconds?完成服务:如何杀死所有线程并通过 5 秒返回结果?
【发布时间】:2009-07-08 09:11:38
【问题描述】:

我对 CompletionService 有一些问题。 我的任务:并行解析大约 300 个 html 页面,我只需要等待所有结果 5 秒, 然后 - 将结果返回给主代码。 我决定为此使用 CompletionService + Callable。 问题是如何停止由 CompletionService 引起的所有线程并从成功解析的页面返回结果? 在这段代码中删除了打印行,但我可以说 5 秒就足够了(有很好的结果,但程序等待所有线程完成)。我的代码执行了大约 2 分钟。

我的电话号码:

Collection<Callable<HCard>> solvers = new ArrayList<Callable<HCard>>();
for (final String currentUrl : allUrls) {
    solvers.add(new Callable<HCard>() {
        public HCard call() throws ParserException {
            HCard hCard = HCardParser.parseOne(currentUrl);                      
            if (hCard != null) {
                return hCard;
            } else {
                return null;
            }
        }
    });
}
ExecutorService execService = Executors.newCachedThreadPool();
Helper helper = new Helper();
List<HCard> result = helper.solve(execService, solvers);
//then i do smth with result list

我的调用代码:

public class Helper {
List<HCard> solve(Executor e, Collection<Callable<HCard>> solvers) throws InterruptedException {
    CompletionService<HCard> cs = new ExecutorCompletionService<HCard>(e);
    int n = solvers.size();

    Future<HCard> future = null;
    HCard hCard = null;
    ArrayList<HCard> result = new ArrayList<HCard>();

    for (Callable<HCard> s : solvers) {
        cs.submit(s);
    }
    for (int i = 0; i < n; ++i) {
        try {
            future = cs.take();
            hCard = future.get();
            if (hCard != null) {
                result.add(hCard);
            }
        } catch (ExecutionException e1) {
            future.cancel(true);
        }
    }
    return result;
}

我尝试使用:

  • awaitTermination(5000, TimeUnit.MILLISECONDS)
  • future.cancel(true)
  • execService.shutdownNow()
  • future.get(5000, TimeUnit.MILLISECONDS);
  • TimeOutException:我无法获得 TimeOutException。

请帮助我了解我的代码上下文。
提前致谢!

【问题讨论】:

  • 你使用缓存线程池有什么原因吗?在最坏的情况下它将启动 300 个线程。考虑使用 FixedThreadPool?
  • 已考虑。我通过探查器看到 - 我的代码没有区别。

标签: java multithreading concurrency


【解决方案1】:

您需要确保您提交的任务正确响应中断,即它们检查 Thread.isInterrupted() 或被视为“可中断”。

我不确定您是否需要完成服务。

ExecutorService service = ...

// Submit all your tasks
for (Task t : tasks) {
    service.submit(t);
}

service.shutdown();

// Wait for termination
boolean success = service.awaitTermination(5, TimeUnit.SECONDS);
if (!success) {
    // awaitTermination timed out, interrupt everyone
    service.shutdownNow();
}

此时,如果您的 Task 对象不响应中断,您无能为力

【讨论】:

  • +1 提到可中断性:这让我遇到了类似的问题(请参阅我的个人资料)。不确定 HCardParser.parseOne 在做什么,但在我的情况下,我调用的函数不受中断的影响。使用 HTTP_Connection 超时,或者你在做什么,你可能会更幸运地让你的 Callables 停止。
【解决方案2】:

问题是你总是得到每一个结果,所以代码总是会运行到完成。我会按照下面的代码使用 CountDownLatch 来完成。

另外,不要使用 Executors.newCachedThreadPool - 这很可能会产生大量线程(如果您的任务需要任何时间,最多可以产生 300 个线程,因为执行程序不会让空闲线程的数量降至零)。

类都是内联的,这样更容易 - 将整个代码块粘贴到一个名为 lotoftasks 的类中并运行它。

import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
import java.util.List;
import java.util.Random;
import java.util.concurrent.Callable;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

public class LotsOfTasks {
    private static final int SIZE = 300;

    public static void main(String[] args) throws InterruptedException {

        String[] allUrls = generateUrls(SIZE);

        Collection<Callable<HCard>> solvers = new ArrayList<Callable<HCard>>();
        for (final String currentUrl : allUrls) {
            solvers.add(new Callable<HCard>() {
                public HCard call() {
                    HCard hCard = HCardParser.parseOne(currentUrl);
                    if (hCard != null) {
                        return hCard;
                    } else {
                        return null;
                    }
                }
            });
        }
        ExecutorService execService = Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors()); // One thread per cpu, ideal for compute-bound
        Helper helper = new Helper();

        System.out.println("Starting..");
        long start = System.nanoTime();
        List<HCard> result = helper.solve(execService, solvers, 5);
        long stop = System.nanoTime();
        for (HCard hCard : result) {
            System.out.println("hCard = " + hCard);
        }

        System.out.println("Took: " + TimeUnit.SECONDS.convert((stop - start), TimeUnit.NANOSECONDS) + " seconds");
    }

    private static String[] generateUrls(final int size) {
        String[] urls = new String[size];
        for (int i = 0; i < size; i++) {
            urls[i] = "" + i;
        }
        return urls;
    }

    private static class HCardParser {
        private static final Random random = new Random();

        public static HCard parseOne(String currentUrl) {
            try {
                Thread.sleep(random.nextInt(1000)); // Wait for a random time up to 1 seconds per task (simulate some activity)
            } catch (InterruptedException e) {
                // ignore
            }
            return new HCard(currentUrl);
        }
    }

    private static class HCard {
        private final String currentUrl;

        public HCard(String currentUrl) {
            this.currentUrl = currentUrl;
        }

        @Override
        public String toString() {
            return "HCard[" + currentUrl + "]";
        }
    }

    private static class Helper {
        List<HCard> solve(ExecutorService e, Collection<Callable<HCard>> solvers, int timeoutSeconds) throws InterruptedException {

            final CountDownLatch latch = new CountDownLatch(solvers.size());

            final ConcurrentLinkedQueue<HCard> executionResults = new ConcurrentLinkedQueue<HCard>();

            for (final Callable<HCard> s : solvers) {
                e.submit(new Callable<HCard>() {
                    public HCard call() throws Exception {
                        try {
                            executionResults.add(s.call());
                        } finally {
                            latch.countDown();
                        }
                        return null;
                    }
                });
            }

            latch.await(timeoutSeconds, TimeUnit.SECONDS);

            final List<Runnable> unfinishedTasks = e.shutdownNow();
            System.out.println("There were " + unfinishedTasks.size() + " urls not processed");

            return Arrays.asList(executionResults.toArray(new HCard[executionResults.size()]));
        }
    }
}

我系统上的典型输出如下所示:

Starting..
There were 279 urls not processed
hCard = HCard[0]
hCard = HCard[1]
hCard = HCard[2]
hCard = HCard[3]
hCard = HCard[5]
hCard = HCard[4]
hCard = HCard[6]
hCard = HCard[8]
hCard = HCard[7]
hCard = HCard[10]
hCard = HCard[11]
hCard = HCard[9]
hCard = HCard[12]
hCard = HCard[14]
hCard = HCard[15]
hCard = HCard[13]
hCard = HCard[16]
hCard = HCard[18]
hCard = HCard[17]
hCard = HCard[20]
hCard = HCard[19]
Took: 5 seconds

【讨论】:

    【解决方案3】:

    我从未使用过 CompletionService,但我确信有一个 poll(timeunit,unit) 调用可以进行有限的等待。然后检查是否为空。测量等待的时间并在 5 秒后停止等待。大约:

    public class Helper {
    List<HCard> solve(Executor e, Collection<Callable<HCard>> solvers) 
    throws InterruptedException {
    CompletionService<HCard> cs = new ExecutorCompletionService<HCard>(e);
    int n = solvers.size();
    
    Future<HCard> future = null;
    HCard hCard = null;
    ArrayList<HCard> result = new ArrayList<HCard>();
    
    for (Callable<HCard> s : solvers) {
        cs.submit(s);
    }
    long timeleft = 5000;
    for (int i = 0; i < n; ++i) {
        if (timeleft <= 0) {
            break;
        }
        try {
            long t = System.currentTimeMillis();
            future = cs.poll(timeleft, TimeUnit.MILLISECONDS);
            timeleft -= System.currentTimeMillis() - t;
            if (future != null) {
                hCard = future.get();
                if (hCard != null) {
                    result.add(hCard);
                }
            } else {
               break;
            }
        } catch (ExecutionException e1) {
            future.cancel(true);
        }
    }
    return result;
    }
    

    虽然没有测试。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2017-04-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-08-12
      • 1970-01-01
      相关资源
      最近更新 更多