【问题标题】:Processing random numbers in parallel Java stream在并行 Java 流中处理随机数
【发布时间】:2016-04-15 11:27:21
【问题描述】:

我想从 0 - 50 范围内生成 5 个不同的随机数,然后对它们并行执行一些操作。当我写这篇文章时,程序从未结束:

new Random().ints(0, 50)
            .distinct()
            .limit(5)
            .parallel()
            .forEach(d -> System.out.println("s: " + d));

我尝试使用 peek 对其进行调试。我有无限数量的 c: 行,50 行 d: 行,但零 l:s: 行:

new Random().ints(0, 50)
            .peek(d -> System.out.println("c: " + d))
            .distinct()
            .peek(d -> System.out.println("d: " + d))
            .limit(5)
            .peek(d -> System.out.println("l: " + d))
            .parallel()
            .forEach(d -> System.out.println("s: " + d));

我的实现有什么问题?

【问题讨论】:

  • IntStream.iterate(…) 这样的无限流和随机数流之间的一个显着区别是随机数流并不是真正无限的,而是具有Long.MAX_VALUE 的大小,甚至报告说,这可能有有趣的效果……
  • 不是this question的重复,请看我的回答。

标签: java parallel-processing java-stream


【解决方案1】:

首先,请注意.parallel()改变了整个流水线的并行状态,所以它会影响所有的操作,而不仅仅是后续的操作。你的情况

new Random().ints(0, 50)
            .distinct()
            .limit(5)
            .parallel()
            .forEach(d -> System.out.println("s: " + d));

是一样的

new Random().ints(0, 50)
            .parallel()
            .distinct()
            .limit(5)
            .forEach(d -> System.out.println("s: " + d));

您不能只并行化部分管道。它是平行的还是不平行的。

现在回到你的问题。由于Random.ints 是无序流,因此选择了distinctlimit 的无序实现,因此它不是this question 的重复项(问题出在有序的不同实现中)。这里的问题在于无序的limit() 实现。为了减少可能的争用,它不会检查在不同线程中找到的元素的总数,直到每个子任务至少获得 128 个元素或上游耗尽(请参阅implementation1 << 7 = 128)。在您的情况下,上游distinct() 仅找到 50 个不同的元素并拼命遍历输入以希望找到更多,但下游limit() 不会发出停止处理的信号,因为它想在检查是否之前收集至少 128 个元素达到了限制(这不是很聪明,因为限制小于 128)。因此,要使这个东西正常工作,您应该至少选择(128 * 数量的 CPU)不同的元素。在我的 4 核机器上使用 new Random().ints(0, 512) 成功,而 new Random().ints(0, 511) 卡住了。

为了解决这个问题,我建议按顺序收集随机数并在那里创建一个新流:

int[] ints = new Random().ints(0, 50).distinct().limit(5).toArray();
Arrays.stream(ints).parallel()
      .forEach(d -> System.out.println("s: " + d));

我假设您想要执行一些昂贵的下游处理。在这种情况下,并行生成 5 个随机数并不是很有用。这部分按顺序执行会更快。

更新:提交了bug report 并提交了patch

【讨论】:

    【解决方案2】:

    您致电ints(0, 50)

    返回一个有效无限的伪随机 int 值流, 每个都符合给定的原点(包括)和绑定(不包括)。

    我原本以为问题出在未终止的IntStream,但我重复了这个问题。

    new Random().ints(0, 50)
                .distinct().limit(5)
                .parallel().forEach(a -> System.out.println(a));
    

    进入无限循环,而

    new Random().ints(0, 50)
                .distinct().limit(5)
                .forEach(a -> System.out.println(a));
    

    正确完成。

    我的 Stream 知识不太好,我无法解释它,但显然并行化并不能很好地发挥作用(可能是由于无限流)。

    【讨论】:

    • 但是当我删除 .parallel() 时,程序将正确打印 5 个不同的数字并退出。为什么在限制后添加.parallel()会使执行无限?
    • @janinko 当你添加parallel() 时,整个流被并行化() 这意味着distinct() 必须使用多个线程来计算。您可能希望收集结果并仅在此基础上进行并行处理。
    • @PeterLawrey 所以我在什么地方使用.parallel() 并不重要?我的想法是只有在.parallel() 之后指定的事情才会并行完成。
    • @janinko 恕我直言,但事实并非如此。
    【解决方案3】:

    与您尝试做的最接近的选择可能是使用iterateunordered

    Random ran = new Random();
    IntStream.iterate(ran.nextInt(50), i -> ran.nextInt(50))
        .unordered()
        .distinct()
        .limit(5)
        .parallel()
        .forEach(System.out::println);
    

    将无限流与distinctparallel 一起使用可能会很昂贵或导致无响应。请参阅API Notethis question 了解更多信息。

    【讨论】:

    • 这似乎有效,但是当我用 ran.ints(0,50) 替换 IntStream.iterate 时,它会循环。为什么IntStream 来自IntStream.iterate 方法的行为与IntStream 来自Random.ints 方法的行为不同?
    • @janinko 你说得对。行为差异看起来很奇怪。我不知道确切的答案,但我怀疑这可能是因为元素是 split 进行并行化的方式。
    猜你喜欢
    • 1970-01-01
    • 2016-05-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-09-25
    相关资源
    最近更新 更多