Javaにおけるスレッドプールの仕組みと実装

スレッドプールの基本概念

スレッドプールとは、複数のスレッドを管理するコンテナであり、一度作成されたスレッドは再利用可能で、頻繁な生成・破棄のコストを削減します。

スレッドプールの主な利点

  • リソース消費の低減:スレッドの生成・終了回数が減少し、既存のスレッドを複数のタスクに再利用可能。
  • 応答速度の向上:タスク到着時に空きスレッドが存在すれば即座に処理可能。システムのフリーズリスクを低減。
  • スレッド管理の容易さ:無制限のスレッド生成はシステム負荷を増大させ、安定性を損ねるため、一元的な制御が可能。

核心となる設計思想

スレッドの再利用」—— 同じスレッドが複数のタスクを順次処理する仕組み。

プーリング技術(Pool)

リソースの再利用を目的としたプログラミング手法。大量のリクエストに対してもパフォーマンスを最適化し、接続確立のコストを軽減。

ExecutorとThreadPoolExecutor

Javaの并发ライブラリでは、java.util.concurrent.Executorインタフェースがスレッドプールの抽象化を提供しています。実際の実装はThreadPoolExecutorです。

private final HashSet<Worker> workers = new HashSet<Worker>();

コンストラクタの主要パラメータ:

public ThreadPoolExecutor(int corePoolSize,
                          int maximumPoolSize,
                          long keepAliveTime,
                          TimeUnit unit,
                          BlockingQueue<Runnable> workQueue,
                          ThreadFactory threadFactory,
                          RejectedExecutionHandler handler)

具体的なインスタンス例:

public static ThreadPoolExecutor executor = new ThreadPoolExecutor(
    2,                    // カーネルスレッド数
    5,                    // 最大スレッド数
    1000,                 // アイドル時間(ミリ秒)
    TimeUnit.MINUTES,
    new LinkedBlockingQueue<Runnable>(10), // キュー容量10
    Executors.defaultThreadFactory(),
    new ThreadPoolExecutor.AbortPolicy()
);

パラメータ解説

パラメータ意味
corePoolSize常時稼働する最小スレッド数。
maximumPoolSize最大同時実行可能なスレッド数。キューが満杯になった場合に追加される。
keepAliveTime非コアスレッドがアイドル状態で保持される時間。
unitkeepAliveTimeの単位(秒、分など)。
workQueue未実行のタスクを一時的に格納するブロッキングキュー。
threadFactory新規スレッドの生成時に使用。名前付与などカスタマイズ可能。
handlerキューもスレッドも満杯時の拒否ポリシー。

拒否ポリシー(RejectedExecutionHandler)

  • AbortPolicy:例外を投げて呼び出し側に通知(デフォルト)。
  • CallerRunsPolicy:呼び出しスレッド自身がタスクを実行(負荷軽減)。
  • DiscardPolicy:タスクを無視して破棄。
  • DiscardOldestPolicy:キュー先頭のタスクを破棄し、現在のタスクを再試行。

動作プロセス

  1. 初期状態:スレッドは未生成(遅延初期化)。
  2. execute()呼び出しでタスク受信。
  3. 判断手順:
    1. 実行中スレッド数 < corePoolSize → 直ちにスレッド生成。
    2. 実行中スレッド数 ≥ corePoolSize → キューに格納。
    3. キュー満杯かつ実行中スレッド数 < maximumPoolSize → 非コアスレッド生成。
    4. キュー満杯かつ実行中スレッド数 ≥ maximumPoolSize → 拒否ポリシー発動。
  4. スレッドがタスク完了後、キューから次のタスクを取得。
  5. アイドル時間がkeepAliveTimeを超えると、コアスレッド数を超えるスレッドは終了。

タスクの送信方法

ExecutorServiceインターフェースの主なメソッド:

メソッド説明
void execute(Runnable command)Runnableタスクを実行(戻り値なし)。
Future<?> submit(Runnable task)タスクを登録。結果はFutureで取得。
Future<T> submit(Callable<T> task)Callableタスクを登録。戻り値付き。
List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks)すべてのタスクを実行し、結果リストを返す。
T invokeAny(Collection<? extends Callable<T>> tasks)最初に完了したタスクの結果を返す。

execute vs submit

  • executeRunnableのみ対応。異常は直ちにスロー。
  • submitRunnableおよびCallable対応。内部でFutureTaskにラップされ、executeを呼ぶ。
  • 異常処理:submitは例外を隠蔽。get()で再発生させる必要あり。

シャットダウン操作

メソッド説明
void shutdown()新しいタスクを受け付けないが、キュー内タスクは処理済み。
List<Runnable> shutdownNow()直ちに停止。実行中のタスク中断。キュー内のタスクを返却。
boolean isShutdown()SHUTDOWN/STOP状態ならtrue
boolean isTerminated()全タスク終了後にtrue
boolean awaitTermination(long timeout, TimeUnit unit)終了待ち。タイムアウトで例外発生。

例外処理の方法

異常はsubmitではキャッチされず、Future.get()で再発生。

// メソッド1:try-catchで直接処理
executor.submit(() -> {
    try {
        System.out.println("task1");
        int i = 1 / 0;
    } catch (Exception e) {
        e.printStackTrace();
    }
});

// メソッド2:Futureで取得
Future<?> future = executor.submit(() -> {
    System.out.println("task1");
    int i = 1 / 0;
    return true;
});
System.out.println(future.get()); // 異常がスローされる

Executorsファクトリメソッド

4種類のプリセットスレッドプールを提供:

1. newFixedThreadPool(n)

public static ExecutorService newFixedThreadPool(int nThreads) {
    return new ThreadPoolExecutor(nThreads, nThreads, 0L, TimeUnit.MILLISECONDS,
            new LinkedBlockingQueue<Runnable>());
}
  • コア数=最大数 → 救急スレッドなし。
  • キューは無限サイズ(Integer.MAX_VALUE)→ OOMリスクあり。
  • 長時間実行タスクに適している。

2. newCachedThreadPool()

public static ExecutorService newCachedThreadPool() {
    return new ThreadPoolExecutor(0, Integer.MAX_VALUE, 60L, TimeUnit.SECONDS,
            new SynchronousQueue<Runnable>>());
}
  • コアスレッド数0、最大スレッド数無限 → 多量のスレッド生成リスク。
  • SynchronousQueue:要素を保持せず、puttakeが同期。
  • 短時間タスクが多く、高密度の処理に適している。

3. newSingleThreadExecutor()

public static ExecutorService newSingleThreadExecutor() {
    return new FinalizableDelegatedExecutorService(
        new ThreadPoolExecutor(1, 1, 0L, TimeUnit.MILLISECONDS,
                new LinkedBlockingQueue<Runnable>()));
}
  • 常に1スレッドで順番に実行。
  • 失敗時は自動再起動(内部で新スレッド生成)。
  • 外部からThreadPoolExecutorのメソッドは呼び出せない(デコレータパターン)。

4. newScheduledThreadPool()

定期実行や遅延実行に対応。スケジューリング機能を持つ。

ブロッキングキューの種類

線形データ構造の阻塞キュー。スレッド安全かつ非同期操作をサポート。

クラス特徴
ArrayBlockingQueue固定サイズの配列ベースのキュー(有界)。
LinkedBlockingQueue連結リストベース。デフォルト無限サイズ(有界/無限可設定)。
PriorityBlockingQueue優先度順に並べ替えられる無限キュー。
DelayQueue指定時間後の実行を保証するキュー。
SynchronousQueue要素を保持しない。1:1のペアリングが必要。
LinkedTransferQueue非同期トランスファー可能な無限キュー。
LinkedBlockingDeque双方向の連結リストキュー。

ブロッキング操作の仕様

  • put(E item):キューが満杯ならブロック。空くまで待機。
  • take():キューが空ならブロック。要素出現まで待機。

タグ: Java Thread Pool ExecutorService ThreadPoolExecutor BlockingQueue

9月15日 09:05 投稿