【问题标题】:Executor tasks for parallelizing a method用于并行化方法的执行器任务
【发布时间】:2018-08-23 01:43:35
【问题描述】:

所以我试图从磁盘中删除 listDir 中列出的 n 个文件,因此我将 listDir 分为 4 个部分,并让它从磁盘中并行删除。这里的目的是并行执行,以使其快速而不是顺序执行。 deleteObject(x,credential,token) 可以假设为一个 API,它最终从磁盘中删除一个对象,并且是一个原子操作。删除成功返回true,否则返回false

我这里有几个问题

  1. 所以我有 4 个并行方法,我正在通过 invokeAll 执行 Executors.newFixedThreadPool(4) 声明了 4 个线程 是否总是将 1 个线程分配给 1 个方法?
  2. 是否需要在 parallelDeleteOperation() 方法的 for 循环中同步和使用 volatile 迭代器“i”。我问这个的原因是假设第一个线程没有完成它的任务(删除 listDir1 并且 for 循环没有完成)并且假设在中途它得到了上下文切换并且第二个线程开始执行相同的任务(删除 listDir1)。只是想知道在这种情况下,第二个线程是否可以获得 IndexOutOfBound 异常。
  3. 将列表分成 4 部分并执行此操作,而不是让多个线程在一个非常大的列表上执行删除操作有什么好处。
  4. 如果 ExecutorService 的其中一项操作返回 false,那么整个 deleteMain() API 将返回 false

     private boolean deleteMain(String parent, List<Structure>
        listDir, String place, String node, Sequence<String>
                                   groups, String Uid) throws IOException {
    
    int noCores = Runtime.getRuntime().availableProcessors();
    List<List<Integer>> splittedList = splitList(listDir, noCores);
    
    System.out.println(splittedList.size());
    
    System.out.println("NoOfCores" + noCores);
    Set<Callable<Boolean>> callables = new HashSet<Callable<Boolean>>();
    for (int i = 0; i < splittedList.size(); i++) {
        List<Integer> l = splittedList.get(i);
        callables.add(new Callable<Boolean>() {
            @Override
            public Boolean call() throws Exception {
    
                return parallelDeleteOperation(parent, listDir, place, node, groups Uid);
            }
        });
    }
    
    ExecutorService service = Executors.newFixedThreadPool(noCores);
    
    try {
        List<Future<Boolean>> futures = service.invokeAll(callables);
    
    
        for (Future<Boolean> future : futures) {
            if (future.get() != true)
                return future.get();
        }
    } catch (InterruptedException e) {
        e.printStackTrace();
    } catch (ExecutionException e) {
        e.printStackTrace();
    }
    service.shutdown();
    return true;
    }
    
    private Boolean parallelDeleteOperation(String parent, List<Structure>
        listDir, String place, String node, Sequence<String>
                                                groups, String Uid) throws IOException {
    for (int i = 0; i < listDir.size(); i++) {
        final String name = listDir.get(i).filename;
    
        final String filePath = "/" + (parent.isEmpty() ? "" : (parent +
                "/")) + name;
        final DeleteMessage message = new DeleteMessage(name, place, node
                filePath);
        final boolean Status = delete(message, groups, Uid, place);
        if (Status != true)
            return Status;
    }
    return true;
    }
    

【问题讨论】:

  • listDir.parallelStream().forEach(this::deleteObject);替换整个东西。
  • 稍微修改了deleteObject代码。

标签: java multithreading future executorservice callable


【解决方案1】:
  1. 既然我有 4 个并行方法,我正在通过 invokeAll 执行,并且 Executors.newFixedThreadPool(4) 声明了 4 个线程,那么是否总是将 1 个线程分配给 1 个方法?

通常,如果任务数量等于池中的线程数是正确的,但不能保证。这在很大程度上取决于池中的任务。

        Executors.newFixedThreadPool(4).invokeAll(IntStream.range(0, 8).mapToObj(i -> (Callable<Integer>) () -> {
            System.out.println(Thread.currentThread().getName() + ": " + i);
            return 0;
        }).collect(Collectors.toList()));

如果池中有更多任务,如上,输出可以如下。处理任务没有明显的可预测规则:

pool-1-thread-1: 0
pool-1-thread-2: 1
pool-1-thread-3: 2
pool-1-thread-1: 4
pool-1-thread-1: 5
pool-1-thread-1: 6
pool-1-thread-1: 7
pool-1-thread-4: 3
  1. 是否需要在 parallelDeleteOperation() 方法的 for 循环中同步和使用 volatile 迭代器“i”。

不,你不需要。您已经将原始列表拆分为单独的四个列表

在您的代码中:

最终列表 listDir1 = listDir.subList(0, listDir.size() / 4);

至于你的第三个问题:

  1. 将列表分成 4 部分并执行此操作,而不是让多个线程在一个非常大的列表上执行删除操作有什么好处。

您最好在现实条件下进行一些测试。这太复杂了,不能说它更好与否。

当您在一个大列表中删除 时,竞争条件可能比预先拆分列表更严重,这可能会产生额外的开销。

此外,即使没有任何并行性,它的性能也不会差。而对于多用户系统,由于线程上下文切换开销,并行版本可能比顺序版本更糟糕。

你必须测试它们,测试助手代码可以是:

    Long start = 0L;
    List<Long> list = new ArrayList<>();
    for (int i = 0; i < 1_000; ++i) {
        start = System.nanoTime();
        // your method to be tested;
        list.add(System.nanoTime() - start);
    }
    System.out.println("Time cost summary: " + list.stream().collect(Collectors.summarizingLong(Long::valueOf)));
  1. 如果 ExecutorService 的操作之一返回 false,那么整个 deleteMain() API 将返回 false

我想重构你的代码,因为这使它更干净,也满足你的最后一个要求(No.4):

// using the core as the count;
// since your task is CPU-bound, we can directly use parallelStream;
private static void testThreadPool(List<Integer> listDir) {
    // if one of the tasks failed, you got isFailed == true;
    boolean isFailed = splitList(listDir, Runtime.getRuntime().availableProcessors()).stream()
            .parallel().map(YourClass::parallelDeleteOperation).anyMatch(ret -> ret == false); // if any is false, it gives you false
}


// split up the list into "count" lists;
private static <T> List<List<T>> splitList(List<T> list, int count) {
    List<List<T>> listList = new ArrayList<>();
    for (int i = 0, blockSize = list.size() / count; i < count; ++i) {
        listList.add(list.subList(i * blockSize, Math.min((i+1) * blockSize, list.size()));
    }
    return listList;
}

【讨论】:

  • 感谢 Hearen,我明白了,这对我来说很有意义。我正在考虑一种方法,如果我可以让 deleteObject(x,credential,token) 以并行方式执行......所以基本上 deleteObject(x,credential,token) 是一个从磁盘中删除文件的 API,它是一个原子操作......我的主列表 listDir 包含要删除的文件的名称。有没有办法在不将列表分成几部分的情况下做到这一点。所以基本上我正在寻找一次读取listDir(不拆分)并将这个deleteObject(x,credential,token)并行调用。我在这里没有选择。
  • listDir.stream().parallel().map(ClassName::deleteMethod).anyMatch(ret -&gt; ret==false) 应该可以做到这一点。
  • 使用直接 stream().parallel() 操作而不是使用拆分选项是否有任何其他优势,正如您基于没有核心所建议的那样。我在网上某处读到,即使并行操作也会根据 CPU 内核拆分流并执行操作。所以基本上现在我使用availableProcessors()根据核心数量拆分我的主列表,然后使用执行器调用所有以FixedThreadPool作为Cpu核心数量来执行它。
  • @vippu 它使用默认的 forkJoinPool 将由程序共享。但是对于整体性能来说,拆分大列表应该比自己的调度方式要好。
  • 当然,您的回答解决了我的很多疑问。我仍然是新的开发人员,也是 stackoverflow 的新手。似乎我的支持不会显示为公开,但是是的,我已对您的回答打了一个肯定的复选标记
【解决方案2】:

我认为将列表拆分为 4 个子列表没有任何好处:一旦线程终止其列表,它将处于空闲状态。如果您为输入列表中的每个元素提交一个任务,则所有四个线程都将处于活动状态,直到队列为空。

更新: 正如有人指出的那样,您可以使用parallelStream,它更简单,可能更快;但如果你想保留ExecutorService,你可以这样做:

    int noCores = Runtime.getRuntime().availableProcessors();
    List<Future<Boolean>> futures = new ArrayList<>();
    ExecutorService service = Executors.newFixedThreadPool(noCores);
    try {
        for (Structure s: listDir) {
            String name = s.filename;
            String filePath = "/" + (parent.isEmpty() ? "" : (parent
                + "/")) + name;
            Future<Boolean> result = service.submit(()-> {
                final DeleteMessage message = new DeleteMessage(
                        name, place, node, filePath);
                return delete(message, groups, Uid, place);
            });
            futures.add(result);
        }
    } finally {
        service.shutdown();
    }

【讨论】:

  • 我同意,但这需要同步列表读取操作,否则从列表中读取的 4 个线程将不按顺序进行?
  • @vippu 不,实际上它不会:您不需要使用任务中的列表
  • 谢谢,我对您的回复投了赞成票,但似乎不会显示。
猜你喜欢
  • 2021-03-24
  • 2016-10-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多