【问题标题】:Nested parallel streams in JavaJava中的嵌套并行流
【发布时间】:2020-07-01 05:34:46
【问题描述】:

我想了解 Java 中嵌套流之间的排序约束。

示例 1:

public static void main(String[] args) {
    IntStream.range(0, 10).forEach(i -> {
        System.out.println(i);
        IntStream.range(0, 10).forEach(j -> {
            System.out.println("    " + i + " " + j);
        });
    });
}

此代码确定性地执行,因此内部循环在每个j 上运行forEach,然后外部循环在下一个i 上运行自己的forEach

0
    0 0
    0 1
    0 2
    0 3
    0 4
    0 5
    0 6
    0 7
    0 8
    0 9
1
    1 0
    1 1
    1 2
    1 3
    1 4
    1 5
    1 6
    1 7
    1 8
    1 9
2
    2 0
    2 1
    2 2
    2 3
...

示例 2:

public static void main(String[] args) {
    IntStream.range(0, 10).parallel().forEach(i -> {
        System.out.println(i);
        IntStream.range(0, 10).parallel().forEach(j -> {
            System.out.println("    " + i + " " + j);
        });
    });
}

如果像第二个示例一样将流设为parallel(),我可以想象内部工作人员在等待线程在外部工作队列中可用时阻塞,因为外部工作队列线程必须在完成时阻塞内部流,默认线程池只有有限数量的线程。但是,死锁似乎没有发生:

6
5
8
    8 6
0
1
    6 2
7
    1 6
    8 5
    7 6
    8 8
2
    0 6
    0 2
    0 8
    5 2
    5 4
    5 6
    0 5
    2 6
    7 2
    7 5
    7 8
    6 4
    8 9
    1 5
 ...

两个流共享相同的默认线程池,但它们生成不同的工作单元。每个外部工作单元只能在该外部工作单元的所有内部单元完成后才能完成,因为每个并行流的末尾都有一个完成障碍。

如何在工作线程共享池中管理这些内部和外部流之间的协调,而不出现任何形式的死锁?

【问题讨论】:

  • “两个流”:不,这里有多个管道(至少 1+10)。为什么你认为这里需要在外部和内部流任务之间进行协调?唯一可以确定的事情是,对于任何给定的i 迭代,相应的j 迭代将在i 迭代完成之前运行,没有任何形式的保证顺序(对于i 或@ 987654334@处决)
  • @ernest_k 因为默认线程池中的工作线程数量有限,所有这些线程最终都会阻塞在外部流中,等待每个对应的内部流完成。但是,如果默认线程池中没有可以处理其项目的可用(空闲/静止)工作线程,则内部流无法完成——这会导致死锁。不过,请参阅 akuzminykh 的回复——显然这种死锁的可能性已被检测到,并根据需要创建了额外的工作线程。

标签: java multithreading java-stream


【解决方案1】:

并行流后面的线程池是公共池,可以通过ForkJoinPool.commonPool() 获得。它通常使用 NumberOfProcessors - 1 个工人。为了解决您所描述的依赖关系,如果(某些)当前工作人员被阻塞并且可能出现死锁,它能够动态创建额外的工作人员。

但是,这不是您的情况的答案。

ForkJoinPool 中的任务有两个重要功能:

  • 他们可以创建子任务并将当前任务拆分成更小的部分(分叉)。
  • 他们可以等待子任务(加入)。

当一个线程执行这样一个任务A并加入一个子任务B时,它不仅仅等待阻塞子任务完成它的执行,而是执行另一个任务C 同时。当 C 完成时,线程返回到 A 并检查 B 是否完成。请注意,BC 可以(并且很可能是)相同的任务。如果 B 完成,则 A 已成功等待/加入它(非阻塞!)。如果前面的解释不清楚,请查看this指南。

现在,当您使用并行流时,流的范围会递归地拆分为多个任务,直到任务变得非常小,以至于它们可以更有效地按顺序执行。这些任务被放入公共池中的工作队列(每个工作人员有一个)。所以,IntStream.range(0, 100).parallel().forEach 所做的是递归地拆分范围,直到它不再值得。每个最终任务,或者更确切地说是一堆迭代,都可以使用forEach 中提供的代码顺序执行。此时,公共池中的工作人员可以执行这些任务,直到所有任务都完成并且流可以返回。请注意,调用线程通过加入子任务来帮助执行!

现在,在您的情况下,这些任务中的每一个都使用并行流本身。程序相同;将其拆分为较小的任务并将这些任务放入公共池中的工作队列中。从ForkJoinPool 的角度来看,这些只是已经存在的任务之上的附加任务。工人只是继续执行/加入任务,直到一切都完成并且外部流可以返回。

这就是您在输出中看到的:没有确定性的行为,没有固定的顺序。也不会发生死锁,因为在给定的用例中不会有阻塞线程。

您可以使用以下代码查看说明:

    public static void main(String[] args) {
        IntStream.range(0, 10).parallel().forEach(i -> {
            IntStream.range(0, 10).parallel().forEach(j -> {
                for (int x = 0; x < 1e6; x++) { Math.sqrt(Math.log(x)); }
                System.out.printf("%d %d %s\n", i, j, Thread.currentThread().getName());
                for (int x = 0; x < 1e6; x++) { Math.sqrt(Math.log(x)); }
            });
        });
    }

您应该注意到主线程参与了内部迭代的执行,因此它没有(!)阻塞。普通的泳池工人只是一个接一个地挑选任务,直到全部完成。

【讨论】:

  • 谢谢,您的解释很有道理。但是,无论我为两个嵌套流设置的结束范围有多大,我都无法使 getPoolSize() 在 16 线程 CPU 上返回超过 15 个。因此,我不确定这里的死锁问题是如何解决的。
  • @LukeHutchison 请尝试以下操作: 在 print 语句前后添加for (int x = 0; x &lt; 1e6; x++) { Math.sqrt(Math.log(x)); } 以模拟计算。这应该会增加池必须创建额外工人的机会。保持外循环的范围较高。内循环的范围可以保持原样。
  • 这并没有改变这种情况——我仍然总是得到 15 getPoolSize()。我在 JDK-14.0.1 上——也许他们改变了最近 JDK 中的逻辑?最终他们将为此使用 Project Jacquard 虚拟线程(这不需要产生额外的物理线程——新的虚拟线程可以以轻量级的方式产生),但我认为 JDK 14 中尚未启用,除非你启用预览功能。
  • 实际上比这更简单。调用线程根本不等待。相反,它正在参与元素的处理。这就是为什么将公共池配置为具有“CPU 核心数减一”工作线程的原因。与启动的非工作线程一起,您可以获得与 CPU 内核一样多的工作线程。
  • @LukeHutchison Holger 指出了正确答案。很抱歉我最初的错误答案。希望它不会太罗嗦,并提供清晰的解释。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多