Javaの同時実行ユーティリティ: Semaphore, CountDownLatch, CyclicBarrierの応用と原理

Semaphoreの使用と原理

概要

適用シナリオ:同時に共有リソースにアクセスできるスレッドの上限を制限するために使用されます。

例:各時点で最大3つのスレッドがリソースにアクセスする

package concurrency.example;

import lombok.extern.slf4j.Slf4j;

import java.util.concurrent.Semaphore;

@Slf4j(topic = "c.ResourceLimiterExample")
public class ResourceLimiterExample {
    public static void main(String[] args) {
        // 1. Semaphoreオブジェクトの作成
        // permitsパラメータはアクセスを許可するスレッド数を制限し、fairパラメータは fairnessを制御します。
        Semaphore resourceLimiter = new Semaphore(3);
        
        // 2. 10個のスレッドを同時に実行しますが、同一時点でリソースを取得できるスレッドは最大3つです。
        // semaphoreはリソースへのアクセスを制限します。
        for (int i = 0; i < 10; i++) {
            new Thread(() -> {
                try {
                    resourceLimiter.acquire();
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                try {
                    log.debug("リソースを使用中...");
                    try {
                        Thread.sleep(1000);
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                    log.debug("リソースの使用を終了しました。");
                } finally {
                    resourceLimiter.release();
                }
            }).start();
        }
    }
}

実行結果:同一時点で最大3つのスレッドがリソースを取得できます。

22:05:10.123 [Thread-2] DEBUG c.ResourceLimiterExample - リソースを使用中...
22:05:10.123 [Thread-1] DEBUG c.ResourceLimiterExample - リソースを使用中...
22:05:10.123 [Thread-0] DEBUG c.ResourceLimiterExample - リソースを使用中...
22:05:11.125 [Thread-2] DEBUG c.ResourceLimiterExample - リソースの使用を終了しました。
22:05:11.125 [Thread-1] DEBUG c.ResourceLimiterExample - リソースの使用を終了しました。
22:05:11.125 [Thread-0] DEBUG c.ResourceLimiterExample - リソースの使用を終了しました。
22:05:11.125 [Thread-3] DEBUG c.ResourceLimiterExample - リソースを使用中...
22:05:11.125 [Thread-4] DEBUG c.ResourceLimiterExample - リソースを使用中...
22:05:11.125 [Thread-5] DEBUG c.ResourceLimiterExample - リソースを使用中...
22:05:12.126 [Thread-4] DEBUG c.ResourceLimiterExample - リソースの使用を終了しました。
22:05:12.126 [Thread-5] DEBUG c.ResourceLimiterExample - リソースの使用を終了しました。
22:05:12.126 [Thread-6] DEBUG c.ResourceLimiterExample - リソースを使用中...
22:05:12.127 [Thread-7] DEBUG c.ResourceLimiterExample - リソースを使用中...
22:05:12.127 [Thread-3] DEBUG c.ResourceLimiterExample - リソースの使用を終了しました。
22:05:12.127 [Thread-8] DEBUG c.ResourceLimiterExample - リソースを使用中...
22:05:13.128 [Thread-8] DEBUG c.ResourceLimiterExample - リソースの使用を終了しました。
22:05:13.128 [Thread-9] DEBUG c.ResourceLimiterExample - リソースを使用中...
22:05:13.128 [Thread-7] DEBUG c.ResourceLimiterExample - リソースの使用を終了しました。
22:05:13.128 [Thread-6] DEBUG c.ResourceLimiterExample - リソースの使用を終了しました。
22:05:14.129 [Thread-9] DEBUG c.ResourceLimiterExample - リソースの使用を終了しました。

Semaphoreの応用:共有リソースへのアクセスを制限する

応用シナリオ

  • 1) Semaphoreによるレートリミット:アクセスのピーク時にはリクエストスレッドをブロックし、ピークが過ぎたら許可を解放します。もちろん、これは単一マシンのスレッド数を制限するのに適しており、リソース数ではなくスレッド数を制限するだけです(リソースの数とリクエストスレッドの数は別の概念です)。
  • 2) Semaphoreを使用してシンプルな接続プールを実装します。『フライウェイトモード』での実装(wait/notifyを使用)と比較すると、パフォーマンスと可読性が明らかに優れています。

- 接続がリソースに対応する場合、たとえばデータベース接続プールはスレッドがデータベース接続に対応するため、この2)のシナリオはSemaphoreの使用に非常に適しています。

Semaphoreを使用してカスタムデータベース接続プールを最適化する

wait/notifyとアトミック配列を組み合わせて、フライウェイトモードに基づいた接続プールを実装する

- 以前のwait/notifyは主に、スレッドプールのリソースがすべて割り当てられた後にwaitを呼び出してスレッドをwaitsetに入れ、スレッドがリソースを解放したときにnotifyを呼び出してスレッドを唤醒するために使用されていました。

最適化後のコード
package concurrency.example;

import lombok.extern.slf4j.Slf4j;

import java.sql.Connection;
import java.util.Random;
import java.util.concurrent.Semaphore;
import java.util.concurrent.atomic.AtomicIntegerArray;

@Slf4j(topic = "c.ConnectionPool")
class ConnectionPool {
    // 01 接続プールのサイズ
    private final int maxConnections;
    // 02 接続オブジェクトの配列
    private Connection[] connections;
    // 03 接続ステータスの配列、0はアイドル、1はビジー
    private AtomicIntegerArray connectionStatus;
    // 04 semaphoreを使用してスレッドプールの容量を管理する
    private Semaphore semaphore;

    ConnectionPool(int maxConnections) {
        this.maxConnections = maxConnections;
        this.semaphore = new Semaphore(maxConnections);
        this.connectionStatus = new AtomicIntegerArray(new int[maxConnections]);
        this.connections = new Connection[maxConnections];

        for (int i = 0; i < maxConnections; ++i) {
            this.connections[i] = new FakeConnection();
        }
    }

    public Connection borrowConnection() {
        try {
            semaphore.acquire(); // リソースの許可を取得し、リソースがない場合はスレッドをブロックする
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        for (int i = 0; i < maxConnections; ++i) {
            if (connectionStatus.get(i) == 0) {
                // 注意:ここではCAS操作を使用して、複数のスレッドによるステータス変数の安全性を確保する必要があります
                connectionStatus.compareAndSet(i, 0, 1);
                log.warn("接続を借用しました: {}", i);
                return connections[i];
            }
        }
        return null;
    }

    public void returnConnection(Connection conn) {
        for (int i = 0; i < maxConnections; ++i) {
            if (connections[i] == conn) {
                connectionStatus.set(i, 0);
                log.warn("接続を解放しました: {}", i);
                semaphore.release(); // 許可を解放する
                break;
            }
        }
    }
}

public class ConnectionPoolExample {
    public static void main(String[] args) {
        ConnectionPool pool = new ConnectionPool(2);
        for (int i = 0; i < 5; ++i) {
            new Thread(() -> {
                Connection tmp = pool.borrowConnection();
                try {
                    Thread.sleep(new Random().nextInt(1000));
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                pool.returnConnection(tmp);
            }).start();
        }
    }
}

// mock:偽の
// ここでは偽の接続プールオブジェクトを実装しています
class FakeConnection implements Connection {
    // Connectionインターフェースのメソッドを実装
    // ...
}

Semaphoreの原理(AQS待ち行列のノードモードはShared)

概要:その基本的な考え方は、読み取りロックの取得と解放と基本的に一致しています。

Semaphoreのコンストラクタ

    public Semaphore(int permits) {
        sync = new NonfairSync(permits);
    }
    public Semaphore(int permits, boolean fair) {
        sync = fair ? new FairSync(permits) : new NonfairSync(permits);
    }
  • Semaphoreにも2種類の同期器があることがわかります:公平な同期器と非公平な同期器
static final class FairSync extends Sync {
        private static final long serialVersionUID = 2014338818796000944L;

        FairSync(int permits) {
            super(permits);
        }

        protected int tryAcquireShared(int acquires) {
            for (;;) {
                if (hasQueuedPredecessors())
                    return -1;
                int available = getState();
                int remaining = available - acquires;
                if (remaining < 0 ||
                    compareAndSetState(available, remaining))
                    return remaining;
            }
        }
    }
========================================================================================
abstract static class Sync extends AbstractQueuedSynchronizer {
        private static final long serialVersionUID = 1192457210091910933L;
        Sync(int permits) {      // 許可の数はAQSのstateに保存されます
            setState(permits);
        }

        final int getPermits() {
            return getState();
        }
    ......
    }
  • 許可の数がAQSのstateに保存されていることがわかります。

Semaphoreのacquireメソッドの原理(リソースの数を取得する)

基本的な考え方

  • 2つのケースに分けて説明します:
    • ケース1:現在のリソースが十分な場合、リソースの数を変更します。
    • ケース2:リソースが不足している場合、スレッドをAQSの待ち行列に入れ、parkを使用して実行を停止させます。
  • いくつかの詳細:
    • 1) stateのメンテナンス(CASメカニズム)
    • 2) 待ち行列のメンテナンス(CASメカニズム)
    • 3) 最初の失敗後、リソースを取得するためのいくつかの試行

CountDownLatch(カウントダウンラッチ)の使用と応用

概要

適用シナリオ:スレッド間の同期と協力を行うために使用されます。すべてのスレッドがカウントダウンを完了するのを待ちます。

  • コンストラクタパラメータは待機カウント値を初期化するために使用されます
  • await() はカウントがゼロになるのを待つために使用されます
  • countDown() はカウントを1減らすために使用されます

joinとの違い:joinは比較的低レベルのAPIであり、使用が面倒で、スレッドの終了を待たなければならないのに対し、CountDownLatchはより柔軟なスレッド間の同期方法を提供します(スレッドは特定の条件を満たした後にcountDown()メソッドを呼び出すことができます)。

簡単な例

package concurrency.example;

import lombok.extern.slf4j.Slf4j;

import java.util.concurrent.CountDownLatch;

@Slf4j(topic = "c.TaskCoordinatorExample")
public class TaskCoordinatorExample {
    public static void main(String[] args) throws InterruptedException {
        /*本質的には、CountDownLatchはstate変数をカウントとして利用し、join関数と同様の機能を実現しています*/
        CountDownLatch taskCoordinator = new CountDownLatch(3);

        new Thread(() -> {
            log.warn("タスク1開始");
            try {
                Thread.sleep(1000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            log.warn("タスク1終了...");
            taskCoordinator.countDown(); // スレッドの実行が完了したら、stateを1減らす
        }, "t1").start();

        new Thread(() -> {
            log.warn("タスク2開始");
            try {
                Thread.sleep(1000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            log.warn("タスク2終了...");
            taskCoordinator.countDown(); // スレッドの実行が完了したら、stateを1減らす
        }, "t2").start();

        new Thread(() -> {
            log.warn("タスク3開始");
            try {
                Thread.sleep(3000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            log.warn("タスク3終了...");
            taskCoordinator.countDown(); // スレッドの実行が完了したら、stateを1減らす
        }, "t3").start();

        log.warn("待機開始");
        taskCoordinator.await(); // stateが0になるまで、以下の処理は実行されず、AQS待ち行列で待機します
        log.warn("待機終了");
    }
}

CountDownLatchとスレッドプールの併用

package concurrency.example;

import lombok.extern.slf4j.Slf4j;

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

@Slf4j(topic = "c.ThreadPoolCoordinatorExample")
public class ThreadPoolCoordinatorExample {
    public static void main(String[] args) throws InterruptedException {
        /*本質的には、CountDownLatchはstate変数をカウントとして利用し、AQSから継承したロックがstateが0でない場合にロックを取得できないという特性を利用して、join関数と同様の機能を実現しています*/
        CountDownLatch taskCoordinator = new CountDownLatch(3);
        ExecutorService pool = Executors.newFixedThreadPool(4);

        pool.submit(() -> {
            log.warn("タスク1開始");
            try {
                Thread.sleep(1000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            log.warn("タスク1終了...");
            taskCoordinator.countDown(); // スレッドの実行が完了したら、stateを1減らす
        });

        pool.submit(() -> {
            log.warn("タスク2開始");
            try {
                Thread.sleep(1000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            log.warn("タスク2終了...");
            taskCoordinator.countDown(); // スレッドの実行が完了したら、stateを1減らす
        });

        pool.submit(() -> {
            log.warn("タスク3開始");
            try {
                Thread.sleep(1000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            log.warn("タスク3終了...");
            taskCoordinator.countDown(); // スレッドの実行が完了したら、stateを1減らす
        });

        pool.submit(() -> {
            log.warn("待機開始");
            try {
                taskCoordinator.await(); // stateが0になるまで待機
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            log.warn("待機終了");
        });

        pool.shutdown();
    }
}

CountDownLatchの典型的な応用シナリオの紹介

シナリオ1:複数のゲームリソースの同時ロード完了通知

このシナリオでは、あるスレッドが次のプロセスに入るには、複数のリソースがロードされる必要があります。ここでは10個と仮定します。スレッドプールを使用して、このスレッドが10個のタスクをそれぞれ1つのリソースをロードするように提出し、countDownLatch.awaitですべてのリソースのロードが完了するのを待ち、各タスクがリソースを準備できたらcountDown()を呼び出します。

シナリオ2:マイクロサービスシナリオでの複数のRPCリモート呼び出しの同時実行

このシナリオでは、現在のユーザーのリクエストには、商品情報、注文情報、配送情報の3つの情報など、複数のサーバーリソースが必要です。順次実行する場合、明らかに効率が低いです。この場合、スレッドプールとCountDownLatchを併用して、スレッドのfutureオブジェクトを使用してスレッドの実行結果を返すことで、リソースを並行して取得し(効率を向上させます)。

CyclicBarrier(サイクリックバリア)の使用と注意点(バッチタスクの循環実行)

概要

役割:スレッド間の協力を行うために使用されます。スレッドが特定のカウントを満たすのを待ちます。コンストラクタで『カウント数』を設定し、各スレッドが「同期」が必要な時点でawait()メソッドを呼び出して待機し、待機しているスレッド数が『カウント数』を満たした場合、実行を続けます。

  • 各回のawait()呼び出しはstateを1減らし、stateが0になると、コンストラクタの2番目のパラメータrunableインターフェースが呼び出されます。

適用シナリオ:繰り返し実行されるバッチタスクに適用されます。CountDownLatchの循環版と見なすことができます。

package concurrency.example;

import lombok.extern.slf4j.Slf4j;

import java.util.concurrent.*;

@Slf4j(topic = "c.PhaseSynchronizerExample")
public class PhaseSynchronizerExample {

    public static void main(String[] args) {
        ExecutorService service = Executors.newFixedThreadPool(2);
        /*パラメータ1:スレッドプールのサイズ、パラメータ2:stateが0になったときに実行されるrunableインターフェース*/
        CyclicBarrier phaseSynchronizer = new CyclicBarrier(2, () -> {
            log.debug("フェーズ1とフェーズ2が完了しました...");
        });
        /*スレッドプール内の2つのタスクを3回循環実行し、CyclicBarrierのカウントが0になると自動的に初期値に戻るため、
           このクラスは再利用可能です。*/
        for (int i = 0; i < 3; i++) { // フェーズ1、フェーズ2、フェーズ1
            service.submit(() -> {
                log.debug("フェーズ1開始...");
                try {
                    Thread.sleep(1000);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                try {
                    phaseSynchronizer.await(); // 2-1=1
                } catch (InterruptedException | BrokenBarrierException e) {
                    e.printStackTrace();
                }
            });
            service.submit(() -> {
                log.debug("フェーズ2開始...");
                try {
                    Thread.sleep(2000);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                try {
                    phaseSynchronizer.await(); // 1-1=0
                } catch (InterruptedException | BrokenBarrierException e) {
                    e.printStackTrace();
                }
            });
        }
        service.shutdown();
    }
}

注意点

固定数のタスクを循環的に実行する場合、スレッドプールのサイズとCyclicBarrierの最初のパラメータが一致していることを必ず確認してください。

Executors.newFixedThreadPool(2);
CyclicBarrier phaseSynchronizer = new CyclicBarrier(2, () -> {
            log.debug("フェーズ1とフェーズ2が完了しました...");
});

スレッドプールのサイズと初期カウントが一致しない場合:

Executors.newFixedThreadPool(3);
CyclicBarrier phaseSynchronizer = new CyclicBarrier(2, () -> {
            log.debug("フェーズ1とフェーズ2が完了しました...");
});

この不一致により、最初の循環中にタスクが完了する前に、余分なアイドルプロセスが次の循環の同じタスクを実行してしまいます。これは、スレッドプールのサイズとCyclicBarrierのカウントが一致していない場合に発生する典型的な問題です。

タグ: Java 並行処理 java.util.concurrent Semaphore CountDownLatch

7月28日 16:37 投稿