File.lines() 并行流的内存使用情况

Aar*_*ron 3 java-8 java-stream

我正在使用 Files.lines() 从大文件(8GB+)读取行。如果按顺序处理,它会很好地工作,并且内存占用非常低。一旦我将parallel()添加到流中,它似乎就会永远挂在它正在处理的数据上,最终导致内存不足异常。我相信这是 Spliterator 在尝试拆分时缓存数据的结果,但我不确定。我剩下的唯一想法是编写一个带有 trySplit 方法的自定义 Spliterator,该方法剥离少量数据进行拆分,而不是尝试将文件拆分为一半或更多。有人遇到过这种情况么?

dka*_*zel 5

跟踪代码我的猜测是isSpliterator使用的。谁的方法有这样的注释:Files.lines()Spliterators.IteratorSpliteratortrySplit()

        /*
         * Split into arrays of arithmetically increasing batch
         * sizes.  This will only improve parallel performance if
         * per-element Consumer actions are more costly than
         * transferring them into an array.  The use of an
         * arithmetic progression in split sizes provides overhead
         * vs parallelism bounds that do not particularly favor or
         * penalize cases of lightweight vs heavyweight element
         * operations, across combinations of #elements vs #cores,
         * whether or not either are known.  We generate
         * O(sqrt(#elements)) splits, allowing O(sqrt(#cores))
         * potential speedup.
         */
Run Code Online (Sandbox Code Playgroud)

然后,代码看起来像是分成了 1024 条记录(行)的倍数的批次。因此第一个分割将读取 1024 行,然后下一个分割将读取 2048 行,依此类推。每次分割都会读取越来越大的批量大小。

如果您的文件确实很大,它最终将达到最大批量大小 33,554,432,即1<<25. 请记住,这是行而不是字节,这可能会导致内存不足错误,特别是当您开始让多个线程读取这么多数据时。

这也解释了速度放缓的原因。在线程可以处理这些行之前,会提前读取这些行。

所以我要么根本不使用parallel(),要么如果你必须使用,因为你所做的计算每行都很昂贵,请编写你自己的 Spliterator,它不会像这样分割。也许总是使用一批 1024 就可以了。