ThreadPoolExecutorとSynchronousQueueの動作メカニズムの詳細解説

ThreadPoolExecutorのインスタンスを生成し、特定の条件下でタスクを投入した場合の動作を分析します。ここでは、コアスレッド数を1、最大スレッド数を2、そしてキューにSynchronousQueueを使用する設定を想定します。そして、3つのタスクを順次投入します。

以下のコード例では、タスク1とタスク2は正常に実行されますが、タスク3は拒否されます。

import java.util.concurrent.*;

public class ThreadPoolExample {
    public static void main(String[] args) {
        ExecutorService executor = new ThreadPoolExecutor(
                1, 2, 10, TimeUnit.SECONDS,
                new SynchronousQueue<>(), new ThreadPoolExecutor.AbortPolicy());

        executor.execute(() -> {
            pause(1000L);
            System.out.println("タスク1の実行完了");
        });
        pause(100L);
        executor.execute(() -> {
            pause(1000L);
            System.out.println("タスク2の実行完了");
        });
        pause(100L);
        executor.execute(() -> {
            pause(1000L);
            System.out.println("タスク3の実行完了");
        });
    }

    private static void pause(long millis) {
        try {
            Thread.sleep(millis);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

実行結果は以下の通りです。タスク3がRejectedExecutionExceptionをスローして拒否されていることがわかります。

Exception in thread "main" java.util.concurrent.RejectedExecutionException: Task java.util.concurrent.FutureTask@... rejected from java.util.concurrent.ThreadPoolExecutor@...[Running, pool size = 2, active threads = 2, queued tasks = 0, completed tasks = 0]
	at java.util.concurrent.ThreadPoolExecutor$AbortPolicy.rejectedExecution(ThreadPoolExecutor.java:2063)
	at java.util.concurrent.ThreadPoolExecutor.reject(ThreadPoolExecutor.java:830)
	at java.util.concurrent.ThreadPoolExecutor.execute(ThreadPoolExecutor.java:1379)
	at ThreadPoolExample.main(ThreadPoolExample.java:15)
タスク1の実行完了
タスク2の実行完了

この動作は、ThreadPoolExecutorとSynchronousQueueのソースコードを分析することで理解できます。

  1. ThreadPoolExecutorのexecuteメソッドのロジック: executeメソッドは、まずworkQueue.offer(command)を呼び出してタスクをキューに投入しようとします。この操作が失敗した場合(falseが返された場合)、addWorkerメソッドを呼び出して新しいスレッドを生成し、タスクの実行を試みます。さらに、addWorkerも失敗した場合には、最終的にrejectメソッドが呼び出され、指定された拒否ポリシー(この例ではAbortPolicy)が適用されます。

    public void execute(Runnable command) {
        if (command == null)
            throw new NullPointerException();
        int c = ctl.get();
        if (workerCountOf(c) < corePoolSize) {
            if (addWorker(command, true))
                return;
            c = ctl.get();
        }
        if (isRunning(c) && workQueue.offer(command)) {
            int recheck = ctl.get();
            if (! isRunning(recheck) && remove(command))
                reject(command);
            else if (workerCountOf(recheck) == 0)
                addWorker(null, false);
        }
        else if (!addWorker(command, false))
            reject(command);
    }
    
  2. SynchronousQueueのofferメソッドの特性: SynchronousQueueは、要素を一時的に保持せず、要素の受け渡しが直接行われる特殊なブロッキングキューです。そのofferメソッドは、要素を取り出すスレッド(コンシューマ)が待機している場合にのみ成功(trueを返す)し、待機するスレッドがいない場合には失敗(falseを返す)します。

    public boolean offer(E e) {
        if (e == null) throw new NullPointerException();
        return transferer.transfer(e, true, 0) != null;
    }
    

したがって、このシナリオの流れは以下のようになります。

  • タスク1の投入: スレッドプールに空きスレッドがあるため、addWorkerが成功し、タスク1が実行されます。
  • タスク2の投入: タスク1を実行中のため、スレッドプールに空きスレッドがありません。workQueue.offer(command)が呼び出されますが、SynchronousQueueにはコンシューマがいないためofferは失敗し(falseを返す)、addWorkerが呼び出されます。最大スレッド数に達していないため、新しいスレッドが生成され、タスク2が実行されます。
  • タスク3の投入: タスク1とタスク2を実行中のため、スレッドプールに空きスレッドがありません。offerは再び失敗します。そして、addWorkerが呼び出されますが、最大スレッド数(2)に達しているため、addWorkerも失敗します。その結果、rejectメソッドが呼び出され、タスク3は拒否されます。

この挙動は、SynchronousQueueを「容量が0のキュー」と捉えると、他のブロッキングキュー(例: LinkedBlockingQueue)に対しても同様のロジックで説明できます。SynchronousQueueは、常に満杯(容量0)であるため、offerは常に失敗し、addWorkerが呼び出されるからです。

タグ: Java ThreadPoolExecutor SynchronousQueue

8月5日 19:58 投稿