【问题标题】:How to manage M threads (1 per task) ensuring only N threads at the same time. With N < M. In Java如何管理 M 个线程(每个任务 1 个)确保同时只有 N 个线程。 N < M。在Java中
【发布时间】:2009-09-13 05:31:54
【问题描述】:

我在 java 中有一个任务队列。此队列位于数据库中的一个表中。

我需要:

  • 每个任务仅 1 个线程
  • 同时运行的线程不超过 N 个。这是因为线程有数据库交互,我不想打开一堆数据库连接。

我想我可以这样做:

final Semaphore semaphore = new Semaphore(N);
while (isOnJob) {
    List<JobTask> tasks = getJobTasks();
    if (!tasks.isEmpty()) {
        final CountDownLatch cdl = new CountDownLatch(tasks.size());
        for (final JobTask task : tasks) {
            Thread tr = new Thread(new Runnable() {

                @Override
                public void run() {
                    semaphore.acquire();
                    task.doWork();
                    semaphore.release();
                    cdl.countDown();
                }

            });
        }
        cdl.await();
    }
}

我知道存在一个 ExecutorService 类,但我不确定是否可以使用它。

那么,您认为这是最好的方法吗?或者您能否澄清一下 ExecutorService 是如何工作的以解决这个问题?

最终解决方案:

我认为最好的解决方案是:

while (isOnJob) {
    ExecutorService executor = Executors.newFixedThreadPool(N);
    List<JobTask> tasks = getJobTasks();
    if (!tasks.isEmpty()) {
        for (final JobTask task : tasks) {
            executor.submit(new Runnable() {

                @Override
                public void run() {
                    task.doWork();
                }

            });
        }
    }
    executor.shutdown();
    executor.awaitTermination(Long.MAX_VALUE, TimeUnit.HOURS);
}

非常感谢遮阳篷。顺便说一句,我正在使用连接池,但是对数据库的查询非常繁重,我不想同时拥有不受控制的任务数量。

【问题讨论】:

    标签: java concurrency multithreading


    【解决方案1】:

    您确实可以使用ExecutorService。例如,使用newFixedThreadPool 方法创建一个新的固定线程池。这样,除了缓存线程之外,您还可以保证不超过n 个线程同时运行。

    类似的东西:

    private static final ExecutorService executor = Executors.newFixedThreadPool(N);
    // ...
    while (isOnJob) {
        List<JobTask> tasks = getJobTasks();
        if (!tasks.isEmpty()) {
            List<Future<?>> futures = new ArrayList<Future<?>>();
            for (final JobTask task : tasks) {
                    Future<?> future = executor.submit(new Runnable() {    
                            @Override
                            public void run() {
                                    task.doWork();
                            }
                    });
                    futures.add(future);
            }
            // you no longer need to use await
            for (Future<?> fut : futures) {
              fut.get();
            }
        }
    }
    

    请注意,您不再需要使用锁存器,因为get 将在必要时等待计算完成。

    【讨论】:

    • 看来我也不需要信号量,是吗?
    【解决方案2】:

    我同意 JG 的观点,即 ExecutorService 是要走的路……但我认为你们都让它变得比需要的更复杂。

    与其创建大量线程(每个任务1个),为什么不创建一个固定大小的线程池(使用Executors.newFixedThreadPool(N))并将所有任务提交给它?不需要信号量或类似的东西 - 只需在获得作业时将作业提交到线程池,线程池将一次处理最多 N 个线程。

    如果您不打算一次使用超过 N 个线程,为什么要创建它们?

    【讨论】:

      【解决方案3】:

      使用具有未绑定队列和固定最大线程大小的 ThreadPoolExecutor 实例,例如Executors.newFixedThreadPool(N)。这将接受大量任务,但只会同时执行 N 个任务。

      如果您选择有界队列(容量为 N),Executor 将拒绝任务的执行(具体取决于您可以配置的策略当直接使用 ThreadPoolExecutor 而不是使用 Executors 工厂时 - 请参阅 RejectedExecutionHandler)。

      如果你需要“真正的”拥塞控制,你应该设置一个容量为N的绑定BlockingQueue。从数据库中获取您想要完成的任务并将它们放入到队列中 - 如果队列已满,则调用线程将阻塞。在另一个线程中(可能也开始使用 Executor API),您BlockingQueue 中获取 任务并将它们提交给 Executor时间>。如果 BlockingQueue 为空,则调用线程也将阻塞。要表示您已完成,请使用“特殊”对象(例如,标记队列中最后一项/最后一项的单例)。

      【讨论】:

        【解决方案4】:

        实现良好的性能还取决于线程中需要完成的工作类型。如果您的数据库是处理中的瓶颈,我会开始关注您的线程如何访问数据库。使用连接池可能是正确的。这可能会帮助您实现更高的吞吐量,因为工作线程可以重用池中的数据库连接。

        【讨论】:

          猜你喜欢
          • 2015-01-29
          • 2022-01-19
          • 2021-12-18
          • 1970-01-01
          • 2021-06-03
          • 2012-02-20
          • 1970-01-01
          • 1970-01-01
          • 2018-10-21
          相关资源
          最近更新 更多