【问题标题】:What is ParallelStream Queue Behavior?什么是 ParallelStream 队列行为?
【发布时间】:2018-09-05 03:06:25
【问题描述】:

我正在使用parallelStream 并行上传一些文件,有些是大文件,有些是小文件。我注意到并非所有工人都被使用。

起初一切正常,所有线程都在使用(我将并行度选项设置为 16)。然后在某个点(一旦它获得更大的文件),它只使用一个线程

简化代码:

files.parallelStream().forEach((file) -> {
    try (FileInputStream fileInputStream = new FileInputStream(file)) {
                IDocumentStorageAdaptor uploader = null;

                try {
                    logger.debug("Adaptors before taking: " + uploaderPool.size());
                    uploader = uploaderPool.take();
                    logger.debug("Took an adaptor!");
                    logger.debug("Adaptors after taking: " + uploaderPool.size());
                    uploader.addNewFile(file);
                } finally {
                    if (uploader != null) {
                        logger.debug("Adding one back!");
                        uploaderPool.put(uploader);
                        logger.debug("Adaptors after putting: " + uploaderPool.size());
                    }
                }
            } catch (InterruptedException | IOException e) {
                throw new UploadException(e);
            }
});

uploaderPool 是一个 ArrayBlockingQueue。 日志:

[ForkJoinPool.commonPool-worker-8] - Adaptors before taking: 0
[ForkJoinPool.commonPool-worker-15] - Adding one back!
[ForkJoinPool.commonPool-worker-8] - Took an adaptor!
[ForkJoinPool.commonPool-worker-15] - Adaptors after putting: 0
...
...
...
[ForkJoinPool.commonPool-worker-10] - Adding one back!
[ForkJoinPool.commonPool-worker-10] - Adaptors after putting: 16
[ForkJoinPool.commonPool-worker-10] - Adaptors before taking: 16
[ForkJoinPool.commonPool-worker-10] - Took an adaptor!
[ForkJoinPool.commonPool-worker-10] - Adaptors after taking: 15
[ForkJoinPool.commonPool-worker-10] - Adding one back!
[ForkJoinPool.commonPool-worker-10] - Adaptors after putting: 16
[ForkJoinPool.commonPool-worker-10] - Adaptors before taking: 16
[ForkJoinPool.commonPool-worker-10] - Took an adaptor!
[ForkJoinPool.commonPool-worker-10] - Adaptors after taking: 15

似乎所有工作(列表中的项目)都分布在 16 个线程中,并且委派给一个线程的事情只会等到线程空闲工作而不是使用可用线程。有没有办法改变 parallelStream 的工作排队方式?我阅读了 forkjoinpool 文档,其中提到了工作窃取,但仅适用于衍生的子任务。

我的另一个计划可能是随机排列我正在使用 parallelStream 的列表,这可能会平衡一些事情。

谢谢!

【问题讨论】:

    标签: java multithreading forkjoinpool


    【解决方案1】:

    并行流的拆分与计算启发式算法针对数据并行操作进行了调整,而不是针对 IO 并行操作进行了调整。 (换句话说,它们被调整为保持 CPU 忙碌,但不会产生比 CPU 更多的任务。)因此,它们偏向于计算而不是分叉。目前没有覆盖这些选择的选项。

    【讨论】:

    • 你对 io-parallel 操作有什么建议吗?
    • @kyl 这在很大程度上取决于您的问题。对于纯粹的 IO-bound 问题,普通的ThreadPoolExecutor 就可以解决问题,而且很容易设置。如果混合使用计算、内存缓冲和 IO,它会变得更加棘手,因为您需要一种机制来确保它们中的每一个都有适当的界限。并行流非常适合它的计算绑定方面。如果您只是枚举文件,那么每个任务的内存很少(只是文件名),因此通过循环或流操作将它们提交给 TPE 应该可以正常工作。
    • 听起来不错,我的问题是上传文件,这似乎主要是 i/o,可能还有一点内存缓冲。最初使用并行流是为了便于设置,但会尝试使用 ThreadPoolExecuter。谢谢!
    • @kyl 这里的挑战是如何优雅地关闭,如果这对你的情况很重要的话。有关示例,请参见 Java 并发实践 (amzn.to/2nzZnkl),第 7 章。
    猜你喜欢
    • 2012-03-09
    • 2011-08-25
    • 2022-07-21
    • 2020-04-19
    • 2015-06-25
    • 1970-01-01
    • 2010-11-06
    • 2010-11-19
    相关资源
    最近更新 更多