【问题标题】:Why Hashmap.values().parallelStream() does not run in parallel while wrap them in ArrayList could work?为什么 Hashmap.values().parallelStream() 不能并行运行,而将它们包装在 ArrayList 中可以工作?
【发布时间】:2019-11-28 09:50:36
【问题描述】:

hashmap 有两个键值对,它们不会被不同的线程并行处理。


import java.util.stream.Stream;
import java.util.Map;
import java.util.HashMap;

class Ideone
{
    public static void main (String[] args) throws java.lang.Exception
    {
        Map<String, Integer> map = new HashMap<>();
        map.put("a", 1);
        map.put("b", 2);
        map.values().parallelStream()
              .peek(x -> System.out.println("processing "+x+" in "+Thread.currentThread()))
              .forEach(System.out::println);
    }
}

输出:

processing 1 in Thread[main,5,main]
1
processing 2 in Thread[main,5,main]
2

网址:https://ideone.com/Hkxkoz

ValueSpliterator 应该尝试将 HashMap 的数组拆分为大小为 1 的插槽,这意味着两个元素应该在不同的线程中处理。

来源:https://www.codota.com/code/java/methods/java8.util.HMSpliterators$ValueSpliterator/%3Cinit%3E

将它们包裹在ArrayList 中后,它可以按预期工作。

        new ArrayList(map.values()).parallelStream()
              .peek(x -> System.out.println("processing "+x+" in "+Thread.currentThread()))
              .forEach(System.out::println);

输出:

processing 1 in Thread[ForkJoinPool.commonPool-worker-3,5,main]
1
processing 2 in Thread[main,5,main]
2

【问题讨论】:

  • 请注意,parallelStream() 的文档说“返回一个可能并行的流,此集合作为其源。允许此方法返回一个顺序流。”。无法保证会使用多个线程。
  • ...考虑到流是一个大小的流,而且非常小,不返回并行流是一个非常明智的选择。
  • 有代码分析吗?高级文档文档在这里没有帮助。
  • @Kamel 为什么要关心细节?您已经提供了根据规范工作的玩具代码(“可能是并行流”),并且您已经得到了一个可能的解释,为什么您的流不是并行的(因为它很小)。如果您愿意,当然可以深入研究源代码,但我认为您不会从中获得任何更重要的信息。
  • @Kayaman 有人做过……

标签: java java-stream spliterator


【解决方案1】:

正如this answer 中所解释的,该问题与HashMap 的容量可能大于其大小以及实际值基于其哈希码分布在后备数组中这一事实有关。

所有基于数组的拆分器的拆分逻辑基本相同,无论您是通过数组、ArrayList 还是 HashMap 流式传输。为了在尽力而为的基础上获得平衡拆分,每个拆分将是(索引)范围的一半,但在 HashMap 的情况下,范围内的实际元素数量与范围大小不同。

原则上,每个基于范围的拆分器都可以拆分为单个元素,但是,客户端代码,即 Stream API 实现,到目前为止可能不会拆分。是否尝试拆分的决定是由预期的元素数量和 CPU 内核数量决定的。

参加以下课程

public static void main(String[] args) {
    Map<String, Integer> map = new HashMap<>();
    map.put("a", 1);
    map.put("b", 2);

    for(int depth: new int[] { 1, 2, Integer.MAX_VALUE }) {
        System.out.println("With max depth: "+depth);
        Tree<Spliterator<Map.Entry<String, Integer>>> spTree
            = split(map.entrySet().spliterator(), depth);
        Tree<String> valueTree = spTree.map(sp -> "estimated: "+sp.estimateSize()+" "
            +StreamSupport.stream(sp, false).collect(Collectors.toList()));
        System.out.println(valueTree);
    }
}

private static <T> Tree<Spliterator<T>> split(Spliterator<T> sp, int depth) {
    Spliterator<T> prefix = depth-- > 0? sp.trySplit(): null;
    return prefix == null?
        new Tree<>(sp): new Tree<>(null, split(prefix, depth), split(sp, depth));
}

public static class Tree<T> {
    final T value;
    List<Tree<T>> children;

    public Tree(T value) {
        this.value = value;
        children = Collections.emptyList();
    }
    public Tree(T value, Tree<T>... ch) {
        this.value = value;
        children = Arrays.asList(ch);
    }
    public <U> Tree<U> map(Function<? super T, ? extends U> f) {
        Tree<U> t = new Tree<>(value == null? null: f.apply(value));
        if(!children.isEmpty()) {
            t.children = new ArrayList<>(children.size());
            for(Tree<T> ch: children) t.children.add(ch.map(f));
        }
        return t;
    }
    public @Override String toString() {
        if(children.isEmpty()) return value == null? "": value.toString();
        final StringBuilder sb = new StringBuilder(100);
        toString(sb, 0, 0);
        return sb.toString();
    }
    public void toString(StringBuilder sb, int preS, int preEnd) {
        final int myHandle = sb.length() - 2;
        sb.append(value == null? "": value).append('\n');
        final int num = children.size() - 1;
        if (num >= 0) {
            if (num != 0) {
                for (int ix = 0; ix < num; ix++) {
                    int nPreS = sb.length();
                    sb.append(sb, preS, preEnd);
                    sb.append("\u2502 ");
                    int nPreE = sb.length();
                    children.get(ix).toString(sb, nPreS, nPreE);
                }
            }
            int nPreS = sb.length();
            sb.append(sb, preS, preEnd);
            final int lastItemHandle = sb.length();
            sb.append("  ");
            int nPreE = sb.length();
            children.get(num).toString(sb, nPreS, nPreE);
            sb.setCharAt(lastItemHandle, '\u2514');
        }
        if (myHandle > 0) {
            sb.setCharAt(myHandle, '\u251c');
            sb.setCharAt(myHandle + 1, '\u2500');
        }
    }
}

你会得到:

With max depth: 1

├─estimated: 1 [a=1, b=2]
└─estimated: 1 []

With max depth: 2

├─
│ ├─estimated: 0 [a=1, b=2]
│ └─estimated: 0 []
└─
  ├─estimated: 0 []
  └─estimated: 0 []

With max depth: 2147483647

├─
│ ├─
│ │ ├─
│ │ │ ├─estimated: 0 []
│ │ │ └─estimated: 0 [a=1]
│ │ └─
│ │   ├─estimated: 0 [b=2]
│ │   └─estimated: 0 []
│ └─
│   ├─
│   │ ├─estimated: 0 []
│   │ └─estimated: 0 []
│   └─
│     ├─estimated: 0 []
│     └─estimated: 0 []
└─
  ├─
  │ ├─
  │ │ ├─estimated: 0 []
  │ │ └─estimated: 0 []
  │ └─
  │   ├─estimated: 0 []
  │   └─estimated: 0 []
  └─
    ├─
    │ ├─estimated: 0 []
    │ └─estimated: 0 []
    └─
      ├─estimated: 0 []
      └─estimated: 0 []

开启ideone

因此,如前所述,如果我们拆分得足够深,拆分器可以拆分为单个元素,但是,两个元素的估计大小并不表明值得这样做。在每次拆分时,它会将估计值减半,虽然您可能会说它对于您感兴趣的元素是错误的,但对于这里的大多数拆分器来说实际上是正确的,因为当下降到最大级别时,大多数拆分器代表一个空范围把它们分开是浪费资源。

正如在另一个答案中所说,该决定是关于平衡拆分工作(或一般准备)和并行化的预期工作,Stream 实现无法提前知道。如果您事先知道每个元素的工作量将非常高,为了证明更多的准备工作是合理的,您可以使用,例如new ArrayList&lt;&gt;(map.[keySet|entrySet|values]()) .parallelStream() 强制执行平衡拆分。通常,对于较大的地图,问题会小得多。

【讨论】:

    【解决方案2】:

    谢谢Holger's answer,我会在这里补充更多细节。

    根本原因来自HashMap.values() 的 sizeEstimate 不准确。 默认情况下,HashMap 的容量为 16,有 2 个元素,由数组支持。 Spliterator 的估计大小为 2。

    每次,每次拆分都会将数组减半。在这种情况下,数组的 16 长度被分成两部分,每半 8 个,每半的估计大小为 1。由于元素是根据哈希码放置的,不幸的是,两个元素位于同一半。

    那么forkjoin框架认为1是below the sizeThreshold,就会停止分裂,开始处理任务。

    同时arrayList不存在这个问题,因为estimatedSize总是准确的。

    【讨论】:

      猜你喜欢
      • 2010-10-16
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-06-22
      • 1970-01-01
      相关资源
      最近更新 更多