【问题标题】:Performance of executorService multithreading poolexecutorService多线程池的性能
【发布时间】:2019-01-07 13:49:05
【问题描述】:

我正在使用 Java 的并发库 ExecutorService 来运行我的任务。写入数据库的阈值是 200 QPS,但是这个程序在 15 个线程的情况下只能达到 20 QPS。我尝试了 5、10、20、30 个线程,它们甚至比 15 个线程还慢。代码如下:

ExecutorService executor = Executors.newFixedThreadPool(15);
List<Callable<Object>> todos = new ArrayList<>();

for (final int id : ids) {
    todos.add(Executors.callable(() -> {
        try {
            TestObject test = testServiceClient.callRemoteService();
    SaveToDatabase();
        } catch (Exception ex) {}
    }));
}
try {
    executor.invokeAll(todos);
} catch (InterruptedException ex) {} 
executor.shutdown();

1)我检查了运行这个程序的linux服务器的CPU使用率,使用率分别为90%和60%(它有4个CPU)。内存使用率仅为 20%。所以CPU和内存还是不错的。数据库服务器的 CPU 使用率很低(大约 20%)。什么可以阻止速度达到 200 QPS?也许这个服务电话:testServiceClient.callRemoteService()?我检查了该调用的服务器配置,它允许每秒进行大量调用。

2) 如果 ids 中的 id 计数超过 50000,使用 invokeAll 是不是一个好主意?我们是否应该将其拆分为更小的批次,例如每批次 5000 个?

【问题讨论】:

  • 是否可以发布正在执行的查询类型? IE。他们都在访问或修改同一行或表吗?它们都是读取、写入还是两者的组合?
  • @M.Deinum "线程数受限于你拥有的核心数...如果远程调用需要很多时间,核心在等待响应,不会被可用于其他任何事情。” - 这至少具有误导性。当然,您可以在任何时间点拥有比 CPU cores 更多的 threads尤其是当一个线程需要等待 I/O(网络、磁盘、...)(“阻塞”)时,另一个(“可运行”)线程将被换入以使用 CPU与此同时。因此,特别是如果您在一个线程中有相对较多的等待,您希望使用 更多 个线程来保持 CPU 忙碌。
  • 这取决于您的 CPU 架构。此外,更多的线程最终将导致换出线程,这在 java 中将导致相当大的性能损失,因为堆栈和其他所有内容都需要在暂停和取消暂停线程时进行序列化/反序列化。此外,为 IO 繁重的操作使用大量线程也不是一个好的解决方案(它适用于 CPU 密集型操作)。因此,它确实不受内核数量的限制,但在只有 16 个内核的情况下创建 200 个线程可能会使您的程序比仅使用 20 个线程时更慢。
  • @M.Deinum “堆栈的性能受到影响,其他所有内容都需要在暂停和取消暂停线程时进行序列化/反序列化。” - 我不知道你在说什么。 Java 线程 1:1 映射到操作系统线程,根本不存在诸如“序列化/反序列化”堆栈或“暂停”线程之类的事情。正如我所说,如果每个线程都可以使用 100% 的 CPU,那么您将无法通过使用多个线程来获得性能。多线程只是屏蔽 I/O 延迟的方法。
  • @M.Deinum 现代 CPU 上的上下文切换会产生非常小的损失。更大的问题通常是内存访问的局部性丢失,这可能或多或少地使数据缓存无效。但这与Java无关,在大多数情况下也不是什么大问题。例如。即使惩罚总共是 1ms,等待 5ms 等待网络响应仍然是切换到另一个线程的好时机。

标签: java multithreading executorservice


【解决方案1】:

除了重复创建和销毁线程池非常昂贵之外,此代码中没有任何内容可以阻止此查询率。我建议使用 Streams API,它不仅更简单,而且可以重用内置线程池

int[] ids = ....
IntStream.of(ids).parallel()
                 .forEach(id -> testServiceClient.callRemoteService(id));

这是一个使用普通服务的基准。主要开销是创建连接的延迟。

public static void main(String[] args) throws IOException {
    ServerSocket ss = new ServerSocket(0);
    Thread service = new Thread(() -> {
        try {
            for (; ; ) {
                try (Socket s = ss.accept()) {
                    s.getOutputStream().write(s.getInputStream().read());
                }
            }
        } catch (Throwable t) {
            t.printStackTrace();
        }
    });
    service.setDaemon(true);
    service.start();

    for (int t = 0; t < 5; t++) {
        long start = System.nanoTime();
        int[] ids = new int[5000];
        IntStream.of(ids).parallel().forEach(id -> {
            try {
                Socket s = new Socket("localhost", ss.getLocalPort());
                s.getOutputStream().write(id);
                s.getInputStream().read();
            } catch (IOException e) {
                e.printStackTrace();
            }
        });
        long time = System.nanoTime() - start;
        System.out.println("Throughput " + (int) (ids.length * 1e9 / time) + " connects/sec");
    }
}

打印

Throughput 12491 connects/sec
Throughput 13138 connects/sec
Throughput 15148 connects/sec
Throughput 14602 connects/sec
Throughput 15807 connects/sec

@grzegorz-piwowarek 提到,使用 ExecutorService 会更好。

    ExecutorService es = Executors.newFixedThreadPool(8);
    for (int t = 0; t < 5; t++) {
        long start = System.nanoTime();
        int[] ids = new int[5000];
        List<Future> futures = new ArrayList<>(ids.length);
        for (int id : ids) {
            futures.add(es.submit(() -> {
                try {
                    Socket s = new Socket("localhost", ss.getLocalPort());
                    s.getOutputStream().write(id);
                    s.getInputStream().read();
                } catch (IOException e) {
                    e.printStackTrace();
                }
            }));
        }
        for (Future future : futures) {
            future.get();
        }
        long time = System.nanoTime() - start;
        System.out.println("Throughput " + (int) (ids.length * 1e9 / time) + " connects/sec");
    }
    es.shutdown();

在这种情况下产生几乎相同的结果。

【讨论】:

  • 这只是 ForkJoinPool 不应该用于 IO,我认为是这样
  • @GrzegorzPiwowarek 可能不是,但它使用起来很简单。完成这项工作后,使用 ExecutorService 很可能会更好。
  • 可能值得注意的是,您在示例中使用了“localhost”,这可能比远程网络快几个数量级,基本上使任务完全受 CPU 限制。如果每个网络连接只是导致每个线程必须等待几毫秒的响应,那么更多线程的缩放效果会更明显。
  • @JimmyB 和 DB 请求可能需要更长的时间,从而为更大的池大小增加更多好处。 +1
【解决方案2】:

你为什么把自己限制在这么少的线程数上?

这样你就错失了表现机会。看来您的任务确实受 CPU 限制。网络操作(远程服务+数据库查询)可能会占用大部分时间来完成每个任务。在这些时间里,单个任务/线程需要等待某个事件(网络,...),另一个线程可以使用 CPU。您为系统提供的线程越多,等待其网络 I/O 完成同时仍有一些线程同时使用 CPU 的线程可能就越多。

我建议您大幅增加执行程序的线程数。正如您所说,两个远程服务器都没有得到充分利用,我认为您的程序运行的主机目前是瓶颈。尝试增加(加倍?)线程数,直到您的 CPU 利用率接近 100% 或内存或远程端成为瓶颈。

顺便说一句,你shutdown执行者,但你真的在等待任务终止吗?你如何衡量“QPS”?

我又想到一件事:如何处理数据库连接? IE。 SaveToDatabase()s 是如何同步的?所有线程是否共享(并竞争)单个连接?或者,更糟糕的是,每个线程会创建一个到数据库的新连接,做它的事情,然后再次关闭连接吗?这可能是一个严重的瓶颈,因为建立 TCP 连接和进行身份验证握手可能会占用与运行简单 SQL 语句一样多的时间。

如果 ids 中的 id 计数超过 50000,使用 调用所有?我们是否应该将其拆分为更小的批次,例如每批 5000 个 批量?

正如@Vaclav Stengl 已经写的那样,Executors 有内部队列,它们在其中排队并从中处理任务。所以不用担心那个。您也可以在创建每个任务后立即为每个任务调用submit。这允许第一个任务在您仍在创建/准备后续任务时已经开始执行,这很有意义,尤其是当每个任务 creation 花费相对较长的时间时,但在所有其他情况下都不会受到伤害。将invokeAll 视为您已经拥有任务集合的情况的便捷方法。如果您自己连续创建任务并且您已经可以访问 ExecutorService 来运行它们,只需 submit() 他们 a.s.a.p.

【讨论】:

    【解决方案3】:

    关于批量拆分: ExecutorService 有用于存储任务的内部队列。在您的情况下,ExecutorService executor = Executors.newFixedThreadPool(15); 有 15 个线程,因此最多 15 个任务将同时运行,其他任务将存储在队列中。队列的大小可以参数化。默认情况下,大小将放大到最大 int。在方法execute 内部调用InvokeAll,当所有线程都在工作时,此方法会将任务放入队列中。

    恕我直言,CPU 未达到 100% 的可能有两种情况:

    1. 尝试扩大线程池
    2. 线程正在等待testServiceClient.callRemoteService() 完成,同时 CPU 正在运行

    【讨论】:

      【解决方案4】:

      QPS 的问题可能是带宽限制或事务执行(它会锁定表或行)。所以你只是增加池大小是行不通的。另外,您可以尝试使用生产者-消费者模式。

      【讨论】:

        猜你喜欢
        • 2011-11-01
        • 2014-08-01
        • 1970-01-01
        • 1970-01-01
        • 2014-11-22
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2023-04-04
        相关资源
        最近更新 更多