Comprendre les politiques de rejet des pools de threads et la gestion des pertes de données

Dans l'écosystème Java, un pool de threads (ThreadPool) est un outil puissant pour gérer l'exécution asynchrone des tâches. Cependant, lorsque la charge dépasse les capacités de traitement définies par la configuration du pool, un mécanisme de protection intervient : la politique de rejet (Rejection Policy).

Qu'est-ce qu'une politique de rejet ?

Une politique de rejet est une stratégie appliquée par l'exécuteur lorsqu'il ne peut plus accepter de nouvelles tâches. Cela se produit généralement dans deux situations combinées :

  • Le nombre de threads actifs a atteint sa limite maximale (maximumPoolSize).
  • La file d'attente (Work Queue) utilisée pour stocker les tâches en attente est totalement saturée.

Les types de politiques standards en Java

L'interface RejectedExecutionHandler du package java.util.concurrent définit le comportement à adopter. Voici les quatre implémentations fournies par le JDK :

1. AbortPolicy (Comportement par défaut)

Cette stratégie lève immédiatement une exception de type RejectedExecutionException au moment de la soumission. Elle informe l'appelant que le système est saturé.

// Exemple d'affectation
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy());

Impact : Arrête le flux de soumission. C'est idéal pour détecter rapidement un problème de dimensionnement en environnement de test.

2. CallerRunsPolicy

Au lieu de rejeter la tâche ou de lever une exception, cette politique force le thread qui a soumis la tâche (le thread appelant) à l'exécuter lui-même.

// Exemple d'affectation
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());

Impact : Crée un effet de contre-pression (backpressure). En occupant le thread appelant, la vitesse de soumission ralentit naturellement, laissant au pool le temps de traiter les tâches en attente.

3. DiscardPolicy

Cette politique supprime simplement la nouvelle tâche sans aucune notification ni exception.

// Exemple d'affectation
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.DiscardPolicy());

Impact : Très risqué, car l'application n'a aucun moyen de savoir qu'une tâche a été ignorée. À n'utiliser que pour des données non critiques (ex: télémétrie facultative).

4. DiscardOldestPolicy

L'exécuteur supprime la tâche la plus anicenne en tête de file d'attente pour tenter de soumetre la nouvelle tâche à nouveau.

// Exemple d'affectation
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.DiscardOldestPolicy());

Impact : Privilégie les données récentes. Utile pour les systèmes temps réel où l'information périmée n'a plus de valeur.

La ploitique de rejet provoque-t-elle des pertes de données ?

La réponse dépend directement de la stratégie choisie :

  • Perte certaine : DiscardPolicy et DiscardOldestPolicy entraînent une perte irrémédiable de données car les tâches sont jetées sans traitement.
  • Perte potentielle : AbortPolicy peut entraîner une perte si l'exception n'est pas correctement capturée et gérée par l'application appelante.
  • Sécurité maximale : CallerRunsPolicy est la seule qui garantit que la tâche sera exécutée, évitant ainsi la perte de données, au prix d'un ralentissement global de l'application.

Implémentation personnalisée

Pour des besoins spécifiques (journalisation, persistance en base de données de secours), il est possible de créer son propre gestionnaire :

public class LoggingRejectionHandler implements RejectedExecutionHandler {
    @Override
    public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
        System.err.println("Alerte : Tâche rejetée détectée. ID: " + r.hashCode());
        // Logique de sauvegarde ou d'alerte ici
    }
}

Configuration pratique : Exemple Java SE

Voici comment configurer un pool robuste avec une politique d'exécution par l'appelant :

import java.util.concurrent.*;

public class PoolManager {
    public static void main(String[] args) {
        BlockingQueue<Runnable> queue = new ArrayBlockingQueue<>(10);
        
        ThreadPoolExecutor workerPool = new ThreadPoolExecutor(
            3,                       // Cœur du pool
            6,                       // Maximum de threads
            30L, TimeUnit.SECONDS,   // Temps de garde
            queue,
            new ThreadPoolExecutor.CallerRunsPolicy()
        );

        for (int i = 0; i < 25; i++) {
            final int taskId = i;
            workerPool.execute(() -> {
                try {
                    Thread.sleep(500);
                    System.out.println("Tâche " + taskId + " terminée par " + Thread.currentThread().getName());
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                }
            });
        }
        workerPool.shutdown();
    }
}

Configuration dans Spring avec ThreadPoolTaskExecutor

Dans un environnement Spring, la configuration s'effectue via le bean ThreadPoolTaskExecutor :

@Configuration
public class AsyncConfig {

    @Bean(name = "primaryExecutor")
    public ThreadPoolTaskExecutor primaryExecutor() {
        ThreadPoolTaskExecutor taskExec = new ThreadPoolTaskExecutor();
        
        taskExec.setCorePoolSize(5);
        taskExec.setMaxPoolSize(10);
        taskExec.setQueueCapacity(50);
        taskExec.setThreadNamePrefix("AsyncThread-");
        
        // Définition de la stratégie de rejet
        taskExec.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
        
        taskExec.initialize();
        return taskExec;
    }
}

Stratégies pour éviter les pertes de données

Pour garantir l'intégrité des données lors d'une surcharge, considérez les approches suivantes :

  1. File d'attente persistante : Utiliser un système de messagerie externe (RabbitMQ, Kafka) pour découpler la réception de l'exécution.
  2. Surveillance active : Monitorer la taille de la file d'attente et ajuster dynamiquement les paramètres du pool via JMX ou des métriques Actuator.
  3. Dégradation gracieuse : En cas de rejet, basculer vers un mode de traitement simplifié ou retourner une réponse de surcharge au client (HTTP 503).

Étiquettes: Java multithreading Concurrency spring-framework ThreadPool

Publié le 30 août à 10h48