【问题标题】:Why ParallelStream won't use all commonPool's thread in recursion?为什么 ParallelStream 不会在递归中使用所有 commonPool 的线程?
【发布时间】:2021-12-01 21:58:58
【问题描述】:

当我运行以下代码时,8 个可用线程中只有 2 个运行,谁能解释为什么会这样?如何更改代码以利用所有 8 个线程?

Tree.java:

package il.co.roy;

import java.util.HashSet;
import java.util.Objects;
import java.util.Set;

public class Tree<T>
{
    private final T data;
    private final Set<Tree<T>> subTrees;

    public Tree(T data, Set<Tree<T>> subTrees)
    {
        this.data = data;
        this.subTrees = subTrees;
    }

    public Tree(T data)
    {
        this(data, new HashSet<>());
    }

    public Tree()
    {
        this(null);
    }

    public T getData()
    {
        return data;
    }

    public Set<Tree<T>> getSubTrees()
    {
        return subTrees;
    }

    @Override
    public boolean equals(Object o)
    {
        if (this == o)
            return true;
        if (o == null || getClass() != o.getClass())
            return false;
        Tree<?> tree = (Tree<?>) o;
        return Objects.equals(data, tree.data) &&
                Objects.equals(subTrees, tree.subTrees);
    }

    @Override
    public int hashCode()
    {
        return Objects.hash(data, subTrees);
    }

    @Override
    public String toString()
    {
        return "Tree{" +
                "data=" + data +
                ", subTrees=" + subTrees +
                '}';
    }

    public void sendCommandAll()
    {
        if (data != null)
            System.out.println("[" + Thread.currentThread().getName() + "] sending command to " + data);
        try
        {
            Thread.sleep(5000);
        } catch (InterruptedException e)
        {
            e.printStackTrace();
        }
        if (data != null)
            System.out.println("[" + Thread.currentThread().getName() + "] tree with data " + data + " got " + true);
        subTrees.parallelStream()
//              .map(Tree::sendCommandAll)
                .forEach(Tree::sendCommandAll);
//              .reduce(true, (aBoolean, aBoolean2) -> aBoolean && aBoolean2);
    }
}

(不管我用forEach还是reduce)。

Main.java:

package il.co.roy;

import java.util.Set;
import java.util.concurrent.Executor;
import java.util.concurrent.Executors;
import java.util.stream.Collectors;
import java.util.stream.IntStream;

public class Main
{
    public static void main(String... args)
    {
        System.out.println("Processors: " + Runtime.getRuntime().availableProcessors());


        final Tree<Integer> root = new Tree<>(null,
                Set.of(new Tree<>(1,
                        IntStream.range(2, 7)
                                        .boxed()
                                        .map(Tree::new)
                                        .collect(Collectors.toSet()))));

        root.sendCommandAll();

//      IntStream.generate(() -> 1)
//              .parallel()
//              .forEach(i ->
//              {
//                  System.out.println(Thread.currentThread().getName());
//                  try
//                  {
//                      Thread.sleep(5000);
//                  } catch (InterruptedException e)
//                  {
//                      e.printStackTrace();
//                  }
//              });
    }
}

main 方法中,我创建了一个具有以下结构的树:\

root (data is `null`)
  |- 1
     |- 2
     |- 3
     |- 4
     |- 5
     |- 6

sendCommandAll 函数处理每个子树(并行)只有当它的父完成处理时。 但结果如下:

处理器:8
[main] 向 1 发送命令
[主] 数据 1 的树得到了 true
[main] 向 6 发送命令
[ForkJoinPool.commonPool-worker-2] 向 5 发送命令
[main] 数据 6 的树得到了 true
[ForkJoinPool.commonPool-worker-2] 数据为 5 的树为 true
[ForkJoinPool.commonPool-worker-2] 向 4 发送命令
[ForkJoinPool.commonPool-worker-2] 数据 4 的树得到了 true
[ForkJoinPool.commonPool-worker-2] 向 3 发送命令
[ForkJoinPool.commonPool-worker-2] 数据 3 的树得到了 true
[ForkJoinPool.commonPool-worker-2] 向 2 发送命令
[ForkJoinPool.commonPool-worker-2] 数据为 2 的树为真

(作为记录,当我执行Main.java 中的注释代码时,JVM 使用了所有可用的 7(+1)个线程commonPool

如何改进我的代码?

【问题讨论】:

  • 我不认为我遇到了同样的问题:首先我使用JDK 17,即使我使用自定义ForkJoinPool,并行度为20,只有2个线程处于活动状态
  • 正如答案所说 - 它不能保证工作。此外,如果我重写您的代码以使用 List 而不是 Set 我会看到池中使用了更多线程。
  • 正如this answer 的后半部分所解释的,HashMaps(进而HashSets)具有少量元素,与它们的(默认)容量相比可能会分配它们的工作糟糕,取决于哈希码分布。您可以使用new ArrayList&lt;&gt;(subTrees).parallelStream() 解决此问题,但您的方法还有其他缺陷,例如在开始遍历孩子之前完成工作/等待。您应该将迭代逻辑与实际操作分开。
  • 谢谢@Holger,它确实解决了我的问题,您能否将您的评论改写为获得积分和积分的官方答案:-)

标签: java multithreading java-8 java-stream java-17


【解决方案1】:

正如this answer 的后半部分所解释的,处理HashMaps 或HashSets 时的线程利用率取决于后备数组中元素的分布,这取决于哈希码。尤其是元素数量较少的情况下,与(默认)容量相比,这可能会导致糟糕的工作拆分。

一个简单的解决方法是使用new ArrayList&lt;&gt;(subTrees).parallelStream() 而不是subTrees.parallelStream()

但请注意,您的方法在处理子节点之前执行当前节点的实际工作(在使用sleep 模拟的示例中),这也会降低潜在的并行性。

你可以使用

public void sendCommandAll() {
    if(subTrees.isEmpty()) {
        actualSendCommand();
        return;
    }
    List<Tree<T>> tmp = new ArrayList<>(subTrees.size() + 1);
    tmp.addAll(subTrees);
    tmp.add(this);
    tmp.parallelStream().forEach(t -> {
        if(t != this) t.sendCommandAll(); else t.actualSendCommand();
    });
}

private void actualSendCommand() {
    if (data != null)
        System.out.println("[" + Thread.currentThread().getName()
                         + "] sending command to " + data);
    try {
        Thread.sleep(5000);
    } catch (InterruptedException e) {
        e.printStackTrace();
    }
    if (data != null)
        System.out.println("[" + Thread.currentThread().getName()
                         + "] tree with data " + data + " got " + true);
}

这允许在处理子节点的同时处理当前节点。

【讨论】:

  • 那么,简单的解决方案new ArrayList&lt;&gt;(subTrees).parallelStream()就可以了。我将为可能没有此限制的未来读者保留其他解决方案……
  • 您阅读链接的答案了吗?挑战在于在不花太多时间分析情况的情况下拆分工作。我调试了您的案例,并且元素确实都聚集在数组的一个区域中,因此将数组拆分为相等大小的范​​围会导致几个全空范围。 TreeMap 不应该有这个问题,因为节点已经(几乎)平衡了,这允许递归地将一半传递给另一个工作线程,但是,它不是完全平衡的,这可能仍然会导致小型集合的并行性降低。跨度>
  • @RoyAsh 这不是错误,也不是特定于基于哈希的集合。并行处理是一种权衡;拆分和加入是有成本的,我们希望通过投入更多的 CPU 来解决这个问题。一般来说,拆分为一个元素并不是最优的(而且,在小型数据集上使用并行性也不是最优的。)拆分启发式算法针对具有大量 CPU 密集型计算的大型数据集进行了调整,以实现有效的并行性。
  • @RoyAsh 试试 LinkedList,情况更糟。这是拆分启发式和集合拓扑以及拆分器实现质量的复杂功能。
  • @RoyAsh ArrayList 知道所有元素都存储在其支持数组中,从索引 0 到 4。将它们分发给工人很容易。但是HashSet 默认有一个长度为 16 的后备数组,并且这五个元素存储在该数组的 somewhere 中,具体取决于它们的哈希码。然后,它获得了“为另一个工人选择(大约)一半的任务”实际上没有迭代的任务。请注意,该算法也必须适用于 1 亿的数组长度。所以它需要一半的数组,希望能紧紧抓住一半的元素。然后,每个工作人员递归地重复此操作。
猜你喜欢
  • 2020-08-29
  • 1970-01-01
  • 2013-09-14
  • 2021-04-25
  • 2022-08-03
  • 2017-10-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多