【问题标题】:Flexible CountDownLatch can't use Phaser because of limit灵活的 CountDownLatch 由于限制不能使用 Phaser
【发布时间】:2016-01-28 18:54:14
【问题描述】:

我收到一个包含 N 个条目的大文件。对于每个条目,我正在创建一个新线程。我需要等待所有 N 个线程都被终止。

一开始我使用的是 Phaser,但它的实现仅限于 65K 方。所以,是因为 N 可能像 100K 一样爆炸。

然后,我尝试了 CountDownLatch。这很好用,非常简单的概念和非常简单的实现。但我不知道N的数量。

Phaser 是我的解决方案,但有这个限制。

有什么想法吗?

这篇文章是相关的: Flexible CountDownLatch?

【问题讨论】:

    标签: java multithreading synchronization phaser


    【解决方案1】:

    听起来您要解决的问题是尽快处理大量任务并等待处理完成。

    同时处理大量任务的问题在于,它可能会导致过多的上下文切换,并且基本上会削弱您的机器并减慢处理速度,使其超过一定数量(取决于硬件)的并发线程。这意味着您需要对正在执行的并发工作线程设置上限。

    Phaser 和 CountDownLatch 都是同步原语,它们的目的是提供对关键代码块的访问控制,而不是管理并行执行。

    在这种情况下,我会使用Executor service。支持添加任务(多种形式,包括Runnable)。

    您可以使用Executors 类轻松创建ExecutorService。我建议为此使用 fixed size thread pool,最大线程数为 20-100 - 取决于您的任务的 CPU 密集程度。任务所需的计算能力越多,可以处理的并行线程数就越少,而不会严重降低性能。

    有多种方法可以等待所有任务完成:

    • 收集由submit 方法返回的所有Future 实例,并简单地对所有这些实例调用get。这可确保在您的循环完成时执行每个任务。
    • Shut down 执行器服务和wait for all the submitted tasks to finish。此方法的缺点是您必须指定等待任务完成的最长时间。此外,它不太优雅,您并不总是想关闭Executor,这取决于您是在编写单次应用程序还是在之后继续运行的服务器 - 如果是服务器应用程序,您将肯定要沿用以前的方法。

    最后,这里有一段代码 sn-p 说明了这一切:

    List<TaskFromFile> tasks = loadFileAndCreateTasks();
    ExecutorService executor = Executors.newFixedThreadPool(50);
    
    for(TaskFromFile task : tasks) {
        // createRunnable is not necessary in case your task implements Runnable
        executor.submit(createRunnable(task));
    }
    
    // assuming single-shot batch job
    executor.shutdown();
    executor.awaitTermination(MAX_WAIT_TIME_MILLIS, TimeUnit.MILLISECONDS);
    

    【讨论】:

    • 我已经实现了一些。我正在使用带有线程池和信号量的 ExecutorService,以避免 CPU 崩溃。但是,我将分析“awaitTermination”。谢谢!
    【解决方案2】:

    ReusableCountLatchCountDownLatch 的替代方案,也允许增量。

    用法如下:

    ReusableCountLatch latch = new ReusableCountLatch(); // creates latch with initial count 0
    ReusableCountLatch latch = new ReusableCountLatch(10); // creates latch with initial count 10
    
    latch.increment(); // increments counter
    
    latch.decrement(); // decrement counter
    
    latch.waitTillZero(); // blocks until counts falls to zero
    
    boolean succeeded = latch.waitTillZero(200, MILLISECONDS); // waits for up to 200 milliseconds until count falls to zero
    
    int count = latch.getCount(); // gets actual count
    

    要使用它,只需将此 gradle/maven 依赖项添加到您的项目中:'com.github.matejtymes:javafixes:1.3.1'

    更多详情可以在这里找到:https://github.com/MatejTymes/JavaFixes

    【讨论】:

      【解决方案3】:

      使用 AtomicInteger,您可以轻松实现相同的目标。用 1 初始化并随着每个新线程递增。一旦在工人和生产者中完成,减量并获得。如果为零,则运行您的整理 Runnable。

      【讨论】:

      • 并非如此。我没有工人和消费者。我有 N 个线程,我需要等待所有 N 个线程完成。假设我将原子整数初始化为零,然后触发 N 个线程,并且在每个线程的开头,有一个 +1。如果由于某种调度原因,一个线程 A 连续增加 +1 和减少 -1,并且这两个操作中间没有其他线程 B,我会完成,而实际上我没有。
      • 我将启动线程称为生产者。您必须在 启动线程之前加入。启动所有线程后,您仍然需要从 main 开始 1 和 dec,否则可能会发生类似上述情况。
      • 从主要减少可能会在相同的情况下运行。主线程将计数器放入 1。线程 A 放入 2,然后线程 A 放入 1。然后主线程将其减少到零。线程 B、C... 可以执行,但 MAIN 线程不会等待它们。
      • 同样,inc() 必须在 thread.start() 之前完成。所以在 MAIN 线程(而不是其他线程)中,在调用 MAIN dec() 之前,它可能不会归零。
      • 抱歉没看懂,但是,如何让主线程等待其他N个线程?
      猜你喜欢
      • 2012-02-14
      • 2018-01-17
      • 1970-01-01
      • 2017-05-27
      • 2021-11-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多