【问题标题】:Concatenating parallel streams连接并行流
【发布时间】:2015-05-26 16:45:43
【问题描述】:

假设我有两个int[] 数组input1input2。我只想从第一个中获取正数,从第二个中获取不同的数字,将它们合并在一起,排序并存储到结果数组中。这可以使用流来执行:

int[] result = IntStream.concat(Arrays.stream(input1).filter(x -> x > 0), 
                   Arrays.stream(input2).distinct()).sorted().toArray();

我想加快任务,所以我考虑让流并行。通常这只是意味着我可以在流构造和终端操作之间的任何地方插入.parallel(),结果是一样的。 IntStream.concat 的 JavaDoc 表示,如果任何输入流是并行的,则生成的流将是并行的。所以我认为将parallel() 制作成input1 流或input2 流或连接流会产生相同的结果。

其实我错了:如果我将.parallel() 添加到结果流中,输入流似乎保持顺序。此外,我可以将输入流(其中一个或两者)标记为.parallel(),然后将结果流转换为.sequential(),但输入保持并行。所以实际上有 8 种可能性:input1、input2 和串联流中的任何一个可以并行或不并行:

int[] sss = IntStream.concat(Arrays.stream(input1).filter(x -> x > 0),
                Arrays.stream(input2).distinct()).sorted().toArray();
int[] ssp = IntStream.concat(Arrays.stream(input1).filter(x -> x > 0),
                Arrays.stream(input2).distinct()).parallel().sorted().toArray();
int[] sps = IntStream.concat(Arrays.stream(input1).filter(x -> x > 0), 
                Arrays.stream(input2).parallel().distinct()).sequential().sorted().toArray();
int[] spp = IntStream.concat(Arrays.stream(input1).filter(x -> x > 0), 
                Arrays.stream(input2).parallel().distinct()).sorted().toArray();
int[] pss = IntStream.concat(Arrays.stream(input1).parallel().filter(x -> x > 0),
                Arrays.stream(input2).distinct()).sequential().sorted().toArray();
int[] psp = IntStream.concat(Arrays.stream(input1).parallel().filter(x -> x > 0),
                Arrays.stream(input2).distinct()).sorted().toArray();
int[] pps = IntStream.concat(Arrays.stream(input1).parallel().filter(x -> x > 0),
                Arrays.stream(input2).parallel().distinct()).sequential().sorted().toArray();
int[] ppp = IntStream.concat(Arrays.stream(input1).parallel().filter(x -> x > 0),
                Arrays.stream(input2).parallel().distinct()).sorted().toArray();

benchmarked 针对不同输入大小的所有版本(在 Core i5 4xCPU、Win7 上使用 JDK 8u45 64bit)并在每种情况下得到不同的结果:

Benchmark           (n)  Mode  Cnt       Score       Error  Units
ConcatTest.SSS      100  avgt   20       7.094 ±     0.069  us/op
ConcatTest.SSS    10000  avgt   20    1542.820 ±    22.194  us/op
ConcatTest.SSS  1000000  avgt   20  350173.723 ±  7140.406  us/op
ConcatTest.SSP      100  avgt   20       6.176 ±     0.043  us/op
ConcatTest.SSP    10000  avgt   20     907.855 ±     8.448  us/op
ConcatTest.SSP  1000000  avgt   20  264193.679 ±  6744.169  us/op
ConcatTest.SPS      100  avgt   20      16.548 ±     0.175  us/op
ConcatTest.SPS    10000  avgt   20    1831.569 ±    13.582  us/op
ConcatTest.SPS  1000000  avgt   20  500736.204 ± 37932.197  us/op
ConcatTest.SPP      100  avgt   20      23.871 ±     0.285  us/op
ConcatTest.SPP    10000  avgt   20    1141.273 ±     9.310  us/op
ConcatTest.SPP  1000000  avgt   20  400582.847 ± 27330.492  us/op
ConcatTest.PSS      100  avgt   20       7.162 ±     0.241  us/op
ConcatTest.PSS    10000  avgt   20    1593.332 ±     7.961  us/op
ConcatTest.PSS  1000000  avgt   20  383920.286 ±  6650.890  us/op
ConcatTest.PSP      100  avgt   20       9.877 ±     0.382  us/op
ConcatTest.PSP    10000  avgt   20     883.639 ±    13.596  us/op
ConcatTest.PSP  1000000  avgt   20  257921.422 ±  7649.434  us/op
ConcatTest.PPS      100  avgt   20      16.412 ±     0.129  us/op
ConcatTest.PPS    10000  avgt   20    1816.782 ±    10.875  us/op
ConcatTest.PPS  1000000  avgt   20  476311.713 ± 19154.558  us/op
ConcatTest.PPP      100  avgt   20      23.078 ±     0.622  us/op
ConcatTest.PPP    10000  avgt   20    1128.889 ±     7.964  us/op
ConcatTest.PPP  1000000  avgt   20  393699.222 ± 56397.445  us/op

从这些结果我只能得出结论,distinct() 步骤的并行化会降低整体性能(至少在我的测试中)。

所以我有以下问题:

  1. 是否有关于如何更好地使用串联流的并行化的官方指南?测试所有可能的组合并不总是可行的(尤其是在连接两个以上的流时),所以有一些“经验法则”会很好。
  2. 似乎如果我连接直接从集合/数组创建的流(在连接之前不执行中间操作),那么结果不太依赖于 parallel() 的位置。这是真的吗?
  3. 除了串联之外,是否还有其他情况,结果取决于流管道的并行化点?

【问题讨论】:

    标签: java parallel-processing java-8 java-stream


    【解决方案1】:

    规范准确地描述了您得到的结果——当您考虑到这一点时,与其他操作不同,我们谈论的不是单个管道,而是三个不同的 Streams,它们保留了独立于其他的属性。

    规范说:“如果任一输入流是并行的,则结果流是 [...] 并行的。”这就是你得到的;如果任一 input 流是并行的,则 resulting 流是并行的(但您可以在之后将其转换为顺序)。但是将 resulting 流更改为并行或顺序不会改变 input 流的性质,也不会将并行和顺序流馈送到concat

    关于性能后果,请咨询documentation, paragraph “Stream operations and pipelines”

    中间操作进一步分为无状态有状态操作。无状态操作,例如filtermap,在处理新元素时不会保留先前看到的元素的状态——每个元素都可以独立于其他元素的操作进行处理。有状态的操作,例如 distinctsorted,在处理新元素时可能会合并来自以前看到的元素的状态。

    有状态操作可能需要在产生结果之前处理整个输入。例如,在查看流的所有元素之前,无法通过对流进行排序产生任何结果。因此,在并行计算下,一些包含有状态中间操作的管道可能需要对数据进行多次传递,或者可能需要缓冲重要数据。仅包含无状态中间操作的管道可以一次性处理,无论是顺序的还是并行的,数据缓冲最少。

    您已经选择了两个命名的 有状态 操作并将它们组合起来。因此,结果流的.sorted() 操作需要缓冲整个内容,然后才能开始排序,这意味着distinct 操作的完成。不同的操作显然很难并行化,因为线程必须同步已经看到的值。

    所以要回答您的第一个问题,这与 concat 无关,而只是 distinct 无法从并行执行中受益。

    这也会使您的第二个问题过时,因为您在两个串联的流中执行完全不同的操作,因此您不能对预先串联的集合/数组执行相同的操作。连接数组并在结果数组上运行 distinct 不太可能产生更好的结果。

    关于你的第三个问题,flatMapparallel 流的行为可能会让人感到意外……

    【讨论】:

      猜你喜欢
      • 2010-09-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-07-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多