Gestion des ressources concurrentes avec Semaphore en Java

Exemple concret : un pool de 5 connexions de base de données. Chaque thread doit obtenir un permis avant d'utiliser une connexion. Si aucun permis n'est disponible, le thread peut soit attendre (via acquire()) soit échouer immédiatement (via tryAcquire()).

public class ConnectionPoolManager {
    public static void main(String[] args) throws InterruptedException {
        Semaphore pool = new Semaphore(5);
        CountDownLatch latch = new CountDownLatch(8);
        for (int i = 0; i < 8; i++) {
            int id = i;
            new Thread(() -> {
                handleConnection(pool, id);
                latch.countDown();
            }).start();
        }
        latch.await();
        System.out.println("Permis disponibles : " + pool.availablePermits());
    }

    private static void handleConnection(Semaphore pool, int id) {
        try {
            if (pool.tryAcquire(200, TimeUnit.MILLISECONDS)) {
                System.out.println("Connexion " + id + " établie");
                Thread.sleep(3000);
                pool.release();
                System.out.println("Connexion " + id + " libérée");
            } else {
                System.out.println("Connexion " + id + " refusée");
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

Résultat typique : les 5 premières connexions sont acceptées, les autres échouent après délai d'attente. Cela montre le comportement non bloquant de tryAcquire() comparé à acquire() qui bloque jusqu'à l'obtention d'un permis.

Dans Dubbo, le paramètre executes limite le nombre de threads pour une méthode spécifique. Par exemple :

<dubbo:service interface="com.example.OrderService" executes="150"/>
<dubbo:service interface="com.example.CancelService" executes="25"/>

La mise en œuvre interne utilise un Semaphore pour éviter les problèmes de concurrence :

public class ThrottlingFilter {
    public Result invoke(Invoker<?> invoker, Invocation invocation) throws RpcException {
        int maxThreads = invoker.getUrl().getMethodParameter(
            invocation.getMethodName(), "executes", 0
        );
        if (maxThreads <= 0) return invoker.invoke(invocation);
        
        Semaphore semaphore = RpcStatus.getSemaphore(
            invoker.getUrl(), invocation.getMethodName(), maxThreads
        );
        if (!semaphore.tryAcquire()) {
            throw new RpcException("Trop de threads pour " + invocation.getMethodName());
        }
        try {
            return invoker.invoke(invocation);
        } finally {
            semaphore.release();
        }
    }
}

Semaphore est basé sur AQS (AbstractQueuedSynchronizer). En mode partagé, lorsqu'un thread obtient un permis, il réveille immédiatement les threads suivants dans la file d'attente, car plusieurs threads peuvent coexister. Ce comportement diffère des verrous exclusifs où le réveil n'arrive qu'après libération.

Un risque courant est le phénomène de "burst" : toutes les ressources sont consommées rapidement en début de période, laissant la ressource inutilisée ensuite. Par exemple, une limite de 100 requêtes/seconde peut être atteinte en 50ms, rendant le système inactif pendant 950ms.

Étiquettes: Semaphore java-concurrency AQS dubbo resource-limiting

Publié le 1 septembre à 21h35