【问题标题】:ParallelStream Collecting to Set?ParallelStream 收集到设置?
【发布时间】:2021-06-03 06:00:20
【问题描述】:

我有一个程序可以从使用parallelStream() 中受益匪浅(大型数据集,其中涉及一些映射和过滤方案,但不依赖于外部变量/同步),但必须作为一个集合收集。

我对并行流有些陌生(这是我第一次使用它们),并尝试使用以下代码却发现这导致了非并发后端的并发修改和幕后的死锁情况.

此映射尝试使用 Linux 本机命令 sudo blockdev --getsize64 unmounted_device_here 获取未安装磁盘的文件大小(我不知道 Java 是否可以在 Linux 上获取未安装磁盘的完整大小,所以我'我只是使用本机方法,因为无论如何这只会在 Linux 系统上发布)

映射方法(死锁):

var mountPath = Paths.get("/dev");
            //Do NVME Drives First
            var list = new ArrayList<Path>(10);
            //Looks like nvme1n1
            //For reasons beyond my understanding replacing [0-9] with \\d does not work here
            try (var directoryStream = Files.newDirectoryStream(mountPath, "nvme[0-9]n[0-9]")) {
                for (Path path : directoryStream) {
                    list.add(path);
                }
            }
//Map to DrivePacket (path, long), note that blockdev return bytes -> GB
var nvmePackets = list.parallelStream().map((drive) -> new DrivePacket(drive,
                    (Long.parseLong(runCommand("sudo", "blockdev", "--getsize64", drive.toAbsolutePath().toString())) / (1024 * 1024 * 1024))))
                    .collect(Collectors.toSet());

IOUtils 来自 Apache 实用程序类:

        <dependency>
            <groupId>org.apache.commons</groupId>
            <artifactId>commons-lang3</artifactId>
            <version>3.12.0</version>
        </dependency>

runCommand(执行本机调用):

public static String runCommand(String... command) {
        try {
            if (DEBUG_MODE) {
                systemMessage("Running Command: " + Arrays.asList(command).stream().collect(Collectors.joining(" ")));
            }
            var builder = new ProcessBuilder(command);
            var result = IOUtils.toString(builder.start().getInputStream(), StandardCharsets.UTF_8).replaceAll("\n", "");
            if (DEBUG_MODE) {
                System.out.println("Result: " + result);
            }
            return result;
        } catch (IOException ex) {
            throw new IllegalStateException(ex);

        }
    }

DrivePacket 类:

   /**
     * A record of the relevant information for a drive
     *
     * Path is the fully qualified /dev/DRIVE path
     */
    public record DrivePacket(Path drivePath, long driveSize) {}

既然操作受益于并发,有没有办法使用parallelStream来做到这一点?还是我必须使用其他技术?

它总是挂在执行这行代码的停止处,并在我使用调试器时在ForkJoinTask.javaexternalAwaitDone(); 处永远等待。

很遗憾,我找不到与 toConcurrentMap() 类似的 Set

我怎样才能避免这种死锁,同时仍然获得计算的并行性并让最终结果成为一个集合?

系统:JDK 16

编辑 0:更新了可重现性代码

鉴于映射代码调用不共享数据的子例程,我不确定为什么这会导致死锁情况。

【问题讨论】:

  • 死锁是否发生在complicated_mapping_here后面的某个地方?
  • 请提供minimal reproducible exampleWith a quick example (Ideone),我无法重现该问题。
  • 明确地说,所有收集器都应该使用并行流,并且永远不应该发生“对非并发后端的并发修改”。有关非并发和并发收集器之间的区别,请参阅this answer。 JDK 16 中的 bug 并非不可能,但我还是先分析一下“complicated_mapping_here”,然后再在 JDK 代码中查找原因。
  • 令人惊讶的是,您使用IOUtils.toString(builder.start().getInputStream(), …),要求进程终止,读取完整输出,但如果进程使用其他通道之一,这可能会死锁。您应该使用Process p = builder.redirectError(Redirect.INHERIT) .start(); p.getOutputStream().close(); var result = IOUtils.toString( p.getInputStream(), … 之类的东西来确保打印错误并立即结束子进程的读取尝试。
  • 旁注:\\d 不起作用的原因是该模式不是正则表达式,而是getPathMatcher 中描述的全局模式。当你使用Files.newDirectoryStream(…)时,会自动添加glob:前缀,所以你不能切换到正则表达式。

标签: java concurrency java-stream java-16


【解决方案1】:

也许您可以将其收集到ConcurrentMap,然后获取keySet。 (假设为映射对象定义了equalshashcode方法):

list.parallelStream()
    .map(x -> complicated_mapping_here)
    .collect(Collectors.toConcurrentMap(Function.identity(),
                                        x -> Boolean.TRUE /*dummy value*/ ))
    .keySet();

【讨论】:

  • 也许,但根据@Holger 对原始帖子的评论,这无关紧要,因为转换为集合是以串行方式进行的(尽管比真正的并发集合慢一点) 一些更微妙的东西在起作用:)
  • @SarahSzabo java.lang.ProcessImpl(将从 builder.start() 调用)使用一些共享资源。但是对于您的代码,似乎没有访问这些对象。您能否尝试使用不同的简单命令而不是sudo blockdev。即使我不确定,请检查使用blockdev --getSize 命令同时获取磁盘大小是否有任何限制。因为您在“列表”上使用parallelStream,所以它可能有重复的条目(尽管从代码看来,重复的驱动器名称似乎不存在)。你可以试试“设置”。
猜你喜欢
  • 1970-01-01
  • 2021-04-06
  • 1970-01-01
  • 2010-11-01
  • 1970-01-01
  • 1970-01-01
  • 2015-11-13
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多