【问题标题】:Java Parallel Streams not working properly on SetJava Parallel Streams 在 Set 上无法正常工作
【发布时间】:2017-11-29 22:19:00
【问题描述】:

一点背景知识,我尝试使用 Java 8 并行流以异步方式调用多个 API。我希望调用每个 API,然后阻塞,直到返回所有 API。我遇到了一个有趣的情况,如果我尝试流式传输地图而不是列表,API 将不再在新线程中被调用。

如果我运行以下代码,每个服务都会在一个新线程中调用:

    List<GitUser> result1 = Arrays.asList(service1, service2, service3).parallelStream()
        .map(s->s.getGitUser())
        .collect(Collectors.toList());

但是,如果我使用地图来完成相同的任务,则每个服务都会被同步调用:

    Map<String, ParallelStreamPOCService> map = new HashMap<>();
    map.put("1", service1);
    map.put("2", service2);
    map.put("3", service3);

    List<GitUser> result2 = map.entrySet().parallelStream()
            .map(s->s.getValue().getGitUser())
            .collect(Collectors.toList());

这里是服务实现:

    public GitUser getGitUser() {
        LOGGER.info("Loading user " + userName);
        String url = String.format("https://api.github.com/users/%s", userName);
        GitUser results = restTemplate.getForObject(url, GitUser.class);
        try {
            TimeUnit.SECONDS.sleep(secondsToSleep);
        } catch (InterruptedException e) {
            throw new IllegalStateException(e);
        }
        LOGGER.error("Finished " + userName);
        return results;
    }

【问题讨论】:

  • 您的样本太小,无法得出任何明确的结论。尝试一千个元素,看看你是否观察到相同的东西。
  • 我试图通过告诉线程休眠 x 秒来在较小的集合中模拟这一点。如果我 service1 休眠 10 秒,service2 休眠 1 秒,我希望 service2 先返回。
  • 注意:你的意思是“方式”,而不是“庄园”。前者的意思是“方式”;后者的意思是“有土地的大乡间别墅”。
  • I would expect service2 to return first,你可能就在这里,但这意味着公平——并行流不能保证

标签: java multithreading concurrency java-stream


【解决方案1】:

this answer 中所述,这是有关如何拆分工作负载的实施细节。 HashMap 有一个内部后备数组,其容量高于条目(通常)。它根据数组元素进行拆分,知道这可能会产生不平衡的拆分,因为确定条目在数组中的分布方式可能会很昂贵。

当你知道只有几个元素时,最简单的解决方案是减少HashSet的容量(默认容量为16):

HashMap<Integer,String> map = new HashMap<>();
map.put(0, "foo");
map.put(1, "bar");
map.put(2, "baz");

map.values().parallelStream().forEach(v -> {
    LockSupport.parkNanos(TimeUnit.MILLISECONDS.toNanos(200));
    System.out.println(v+"\t"+Thread.currentThread());
});
foo Thread[main,5,main]
bar Thread[main,5,main]
baz Thread[main,5,main]
HashMap<Integer,String> map = new HashMap<>(4);
map.put(0, "foo");
map.put(1, "bar");
map.put(2, "baz");

map.values().parallelStream().forEach(v -> {
    LockSupport.parkNanos(TimeUnit.MILLISECONDS.toNanos(200));
    System.out.println(v+"\t"+Thread.currentThread());
});
foo Thread[ForkJoinPool.commonPool-worker-1,5,main]
baz Thread[main,5,main]
bar Thread[ForkJoinPool.commonPool-worker-1,5,main]

请注意,由于舍入问题,它仍然不会为每个元素使用一个线程。如前所述,HashMapSpliterator 不知道元素是如何分布在数组中的。但它知道一共有三个元素,所以它估计在拆分后每个工作负载中都有一半。三的一半四舍五入为一,因此Stream 实现假定即使尝试进一步细分这些工作负载也没有任何好处。

除了使用具有更多元素的并行流之外,没有简单的解决方法。不过,仅出于教育目的:

HashMap<Integer,String> map = new HashMap<>(4, 1f);
map.put(0, "foo");
map.put(1, "bar");
map.put(2, "baz");
map.put(3, null);

map.values().parallelStream()
   .filter(Objects::nonNull)
   .forEach(v -> {
    LockSupport.parkNanos(TimeUnit.MILLISECONDS.toNanos(200));
    System.out.println(v+"\t"+Thread.currentThread());
});
bar Thread[ForkJoinPool.commonPool-worker-1,5,main]
baz Thread[main,5,main]
foo Thread[ForkJoinPool.commonPool-worker-2,5,main]

通过插入第四个元素,我们消除了舍入问题。它还需要提供1f 的负载因子,以防止HashMap 增加容量,这会使我们回到第一方位(除非我们至少有八个内核)。

这是一个杂项,正如我们事先知道的那样,我们将浪费一个工作线程来检测我们的虚拟 null 条目。但它展示了工作负载拆分的工作原理。

在地图中拥有更多元素会自动消除这些问题。


Stream 不适用于阻塞或休眠的任务。对于这类任务,您应该使用ExecutorService,这也将允许使用比 CPU 内核更多的线程,这对于在整个执行时间内不使用 CPU 内核的任务是合理的。

ExecutorService es = Executors.newCachedThreadPool();
List<GitUser> result =
    es.invokeAll(
        Stream.of(service1, service2, service3)
              .<Callable<GitUser>>map(s -> s::getGitUser)
              .collect(Collectors.toList())
    ) .stream()
      .map(future -> {
            try { return future.get(); }
            catch (InterruptedException|ExecutionException ex) {
                throw new IllegalStateException(ex);
            }
        })
      .collect(Collectors.toList());

【讨论】:

  • 您能否详细说明“这对于在整个执行时间内不使用 CPU 内核的任务来说是合理的”是什么意思。您是指发出 HTTP 请求吗?
  • 发出 HTTP 请求是不使用 CPU 的任务之一,因为它意味着等待传入数据。一般来说,所有类型的 I/O 或等待事件或另一个线程都属于这一类。
【解决方案2】:

请注意,Java 文档声明 parallelStream() “返回一个可能以这个集合为源的并行流。允许此方法返回一个顺序流。”。似乎流中的元素数量太小无法分解为多个线程执行的组件,因此您的代码是按顺序执行的,而不是并行执行的。仅仅因为您使用parallelStream() 并不一定意味着它将并行执行,这完全取决于库确定的内容是否足以并行化以获得最佳性能。

除了前面提到的 Joe C 在 cmets 中所说的之外,您还需要大幅增加流源中的元素数量,才能真正接近看到性能方面的任何影响。

【讨论】:

  • 这是文档中非常有效的一点。但是,两种实现都具有相同数量的元素。
  • @PhilNinan 能否分享一下您是如何计算出第一种方法异步和第二种方法同步执行的?
  • 我实例化了每个服务,并以秒为单位:service1 = new Service(5)service2 = new Service(1)service3 = new Service(3)
  • 列表输出:这是列表的输出:`2017-11-30 12:12:41,700 ERROR main com.lbisoft.core.app.components.video.service.ParallelStreamPOCService - Finished user2 2017-11-30 12:12:43,618 错误 ForkJoinPool.commonPool-worker-2 com.lbisoft.core.app.components.video.service.ParallelStreamPOCService - 完成 user3 2017-11-30 12:12:45,628 错误 ForkJoinPool。 commonPool-worker-1 com.lbisoft.core.app.components.video.service.ParallelStreamPOCService - 完成 user11
  • 地图输出:2017-11-30 12:12:50,713 错误 ForkJoinPool.commonPool-worker-1 com.lbisoft.core.app.components.video.service.ParallelStreamPOCService - 完成 user1 2017- 11-30 12:12:51,790 错误 ForkJoinPool.commonPool-worker-1 com.lbisoft.core.app.components.video.service.ParallelStreamPOCService - 完成 user2 2017-11-30 12:12:54,947 错误 ForkJoinPool.commonPool-worker -1 com.lbisoft.core.app.components.video.service.ParallelStreamPOCService - 完成 user3
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-01-14
  • 1970-01-01
  • 2021-06-17
  • 2019-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-07-14
相关资源
最近更新 更多