【问题标题】:How can I process a Java stream with more than the default number of threads?如何处理超过默认线程数的 Java 流?
【发布时间】:2016-02-23 15:40:50
【问题描述】:

默认情况下,Java 流由common thread pool 处理,该common thread pool 使用默认参数构造。正如another question 中已回答的那样,可以通过指定自定义池或设置java.util.concurrent.ForkJoinPool.common.parallelism 系统参数来调整这些默认值。

但是,我无法通过这两种方法中的任何一种来增加分配给流处理的线程数。例如,考虑下面的程序,它处理包含在其第一个参数中指定的文件中的 IP 地址列表并输出解析的地址。在具有大约 13000 个唯一 IP 地址的文件上运行此程序,我发现使用 Oracle Java Mission Control 的线程少至 16 个。其中,只有五个是ForkJoinPool 工人。然而,这个特定的任务会从更多的线程中受益,因为线程大部分时间都在等待 DNS 响应。所以我的问题是,我怎样才能真正增加使用的线程数?

我已经在三个环境中尝试过该程序;这些是操作系统报告的线程数。

  • Java SE 运行时环境在运行 Windows 7 的 8 核机器上构建 1.8.0_73-b02:17 个线程
  • Java SE 运行时环境在运行 OS X Darwin 15.2.0 的 2 核机器上构建 1.8.0_66-b17:23 个线程
  • 在运行 FreeBSD 11.0 的 24 核机器上的 openjdk 版本 1.8.0_72:44 个线程

import java.io.IOException;
import java.net.InetAddress;
import java.net.UnknownHostException;
import java.nio.file.Files;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.concurrent.ForkJoinPool;

/** Resolve IP addresses in file args[0] using 100 threads */
public class Resolve100 {
    /** Resolve the passed IP address into a name */
    static String addressName(String ipAddress) {
        try {
            return InetAddress.getByName(ipAddress).getHostName();
        } catch (UnknownHostException e) {
            return ipAddress;
        }
    }

    public static void main(String[] args) {
        Path path = Paths.get(args[0]);
        ForkJoinPool fjp = new ForkJoinPool(100);
        try {
            fjp.submit(() -> {
                try {
                    Files.lines(path)
                    .parallel()
                    .map(line -> addressName(line))
                    .forEach(System.out::println);
                } catch (IOException e) {
                    System.err.println("Failed: " + e);
                }
            }).get();
        } catch (Exception e) {
            System.err.println("Failed: " + e);
        }
    }
}

【问题讨论】:

  • 您应该将Files.lines() 括在try-with-resources 语句中!
  • 我建议您在尝试对其进行并行处理之前将这些行添加到列表中。当它预先知道有多少条目时,它会做得更好。

标签: java multithreading java-stream forkjoinpool


【解决方案1】:

您的方法存在两个问题。首先是使用自定义 FJP 不会改变流 API 创建的单个任务的最大数量,因为这是在 in the following way 中定义的:

static final int LEAF_TARGET = ForkJoinPool.getCommonPoolParallelism() << 2;

因此,即使您使用自定义池,并行任务的数量也会受到commonPoolParallelism * 4 的限制。 (实际上并不是硬性限制,而是一个目标,但在很多情况下任务数等于这个数)。

上述问题可以通过使用java.util.concurrent.ForkJoinPool.common.parallelism 系统属性来解决,但是在这里你遇到了另一个问题:Files.lines 并行化非常糟糕。有关详细信息,请参阅this question。特别是,对于 13000 条输入线,最大可能的加速是 3.17 倍(假设每条线的处理时间大致相同),即使您有 100 个 CPU。我的StreamEx 库为此提供了解决方法(使用StreamEx.ofLines(path).parallel() 创建流)。另一种可能的解决方案是将文件行顺序读取到List,然后从中创建一个并行流:

Files.readAllLines(path).parallelStream()...

这将与系统属性一起使用。然而,一般来说,当任务涉及 I/O 时,Stream API 并不适合并行处理。更灵活的解决方案是对每一行使用CompletableFuture

ForkJoinPool fjp = new ForkJoinPool(100);
List<CompletableFuture<String>> list = Files.lines(path)
    .map(line -> CompletableFuture.supplyAsync(() -> addressName(line), fjp))
    .collect(Collectors.toList());
list.stream().map(CompletableFuture::join)
    .forEach(System.out::println);

这样您就不需要调整系统属性,并且可以将单独的池用于单独的任务。

【讨论】:

  • 并且不应该提到这种改变线程数量的技术完全依赖于实现,未指定的行为,开发人员应该依赖于任何东西
  • @Holger,我猜你的意思是 .submit 方法,对吧?
  • 谢谢! CompletableFuture 方法确实产生了 100 个线程并提供了一个数量级的加速。这是数字。原件:48m40.036s; CompletableFuture:0m37.465s。 (请注意,原始版本也在热 DNS 缓存上运行。)
  • @Diomidis Spinellis:对,submit 方法,它改变了流的行为,因为流使用了 Fork/Join,这是一个实现细节。
  • @Tagir Valeev:您的方法的一个限制是它需要与项目数量成比例的内存。是否可以在单个流上使用指定数量的线程运行处理?我尝试了简单的方法(消除 collect() 和 stream.list()),但线程数再次减少到 16。我不确定发生了什么。
猜你喜欢
  • 1970-01-01
  • 2013-07-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-04-29
相关资源
最近更新 更多