L'API CompletableFuture en Java représente un outil puissant pour la programmation asynchrone et concurrente, offrant une manière élégante de gérer les opérations non bloquantes. Cependant, sa puissance s'accompagne de pièges qui, s'ils ne sont pas identifiés et gérés correctement, peuvent entraîner des problèmes de performance, des fuites de mémoire ou des bugs difficiles à diagnostiquer. Cet article explore les écueils les plus fréquents lors de l'utilisation de CompletableFuture et propose des stratégies pour les éviter.
Introduciton à CompletableFuture
Pour ceux qui découvrent CompletableFuture, ses capacités peuvent être impressionnantes. Il facilite l'écriture de code asynchrone de manière plus lisible et composable qu'avec les interfaces Future traditionnelles. Voici un exemple d'utilisation simple :
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
public class DemoBasiqueCompletableFuture {
public static void main(String[] args) throws InterruptedException, ExecutionException {
// Exécution d'un calcul asynchrone
CompletableFuture<String> resultatFutur = CompletableFuture.supplyAsync(() -> {
// Simulation d'une opération longue
try {
TimeUnit.SECONDS.sleep(1);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
System.err.println("Opération interrompue.");
}
return "Opération asynchrone terminée avec succès !";
});
// Récupération du résultat (bloquant)
String message = resultatFutur.get();
System.out.println(message);
}
}
Cette simplicité apparente masque des complexités sous-jacentes. Abordons maintenant les pièges principaux.
Gestion inappropriée des pools de threads
L'un des pièges les plus courants et pourtant les plus négligés est la mauvaise gestion des pools de threads avec CompletableFuture.
Le piège du pool de threads par défaut
Lorsqu'aucune instance d'Executor n'est spécifiée, CompletableFuture utilise le pool de threads commun de ForkJoinPool (ForkJoinPool.commonPool()). Ce pool est optimisé pour les tâches intensives en CPU et a une taille généralement limitée au nombre de cœurs CPU disponibles moins un. Utiliser ce pool pour des tâches intensives en I/O ou bloquantes peut entraîner des performances médiocres, voire des blocages complets de l'application.
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ForkJoinPool;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
public class UtilisationDefautDangereuse {
// Danger: Utilisation massive du pool de threads par défaut pour des tâches I/O
public List<String> traiterLotsDeDonnees(List<String> entrees) {
List<CompletableFuture<String>> futures = entrees.stream()
.map(donnee -> CompletableFuture.supplyAsync(() -> {
// Cette tâche simule une opération I/O bloquante
// qui utilise le ForkJoinPool.commonPool() par défaut.
return recupererDonneeExterne(donnee);
}))
.collect(Collectors.toList());
// Attendre que toutes les tâches se terminent
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
return futures.stream()
.map(CompletableFuture::join) // Bloque pour obtenir le résultat
.collect(Collectors.toList());
}
private String recupererDonneeExterne(String identifiant) {
try {
// Simulation d'un appel réseau ou BD
TimeUnit.MILLISECONDS.sleep(150);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
System.err.println("Récupération de donnée interrompue.");
}
return "DonnéeTraitée(" + identifiant + ")";
}
public static void main(String[] args) {
UtilisationDefautDangereuse demo = new UtilisationDefautDangereuse();
List<String> testData = new ArrayList<>();
for (int i = 0; i < 100; i++) {
testData.add("Item-" + i);
}
long debut = System.currentTimeMillis();
List<String> resultats = demo.traiterLotsDeDonnees(testData);
long fin = System.currentTimeMillis();
System.out.println("Traitement de " + resultats.size() + " éléments en " + (fin - debut) + " ms");
System.out.println("Taille du pool commun: " + ForkJoinPool.commonPool().getPoolSize());
}
}
Analyse du problème :
- Le pool de threads par défaut est dimensionné pour les calculs, pas pour l'I/O.
- De nombreuses tâches I/O bloquantes peuvent saturer le pool par défaut, causant des goulots d'étranglement et des ralentissements.
- Si le taux de soumisison des tâches dépasse la capacité de traitement, cela peut mener à des problèmes de mémoire (
OutOfMemoryError).
Utilisation correcte des pools de threads
La solution consiste à fournir des pools de threads personnalisés et adaptés au type de tâche (I/O intensive ou CPU intensive).
import java.util.concurrent.*;
import java.util.function.Function;
public class GestionnairePoolsCorrecte {
private final ExecutorService executeurIO;
private final ExecutorService executeurCPU;
public GestionnairePoolsCorrecte() {
// Pool pour les tâches intensives en I/O - Plus grand pool
this.operateurIO = new ThreadPoolExecutor(
50, // Nombre de threads principaux
100, // Nombre maximum de threads
60L, TimeUnit.SECONDS, // Temps de survie des threads inactifs
new LinkedBlockingQueue<>(2000), // File d'attente de travail
new ThreadFactory() { // Fabrique de threads personnalisée
private int compteur = 0;
public Thread newThread(Runnable r) {
return new Thread(r, "pool-io-" + compteur++);
}
},
new ThreadPoolExecutor.CallerRunsPolicy() // Politique de rejet
);
// Pool pour les tâches intensives en CPU - Plus petit pool
this.operateurCPU = new ThreadPoolExecutor(
Runtime.getRuntime().availableProcessors(), // Nb de cœurs CPU
Runtime.getRuntime().availableProcessors() * 2, // Max threads (à adapter)
60L, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(100),
new ThreadFactory() {
private int compteur = 0;
public Thread newThread(Runnable r) {
return new Thread(r, "pool-cpu-" + compteur++);
}
},
new ThreadPoolExecutor.AbortPolicy()
);
}
public CompletableFuture<String> traiterAvecPoolIO(String identifiant) {
return CompletableFuture.supplyAsync(() -> {
System.out.println("Tâche IO dans thread: " + Thread.currentThread().getName());
return recupererInfosDuSysteme(identifiant);
}, operateurIO);
}
public CompletableFuture<Integer> calculerAvecPoolCPU(int nombre) {
return CompletableFuture.supplyAsync(() -> {
System.out.println("Tâche CPU dans thread: " + Thread.currentThread().getName());
return effectuerCalculLourd(nombre);
}, operateurCPU);
}
private String recupererInfosDuSysteme(String id) {
try {
TimeUnit.MILLISECONDS.sleep(200); // Simule une latence I/O
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
return "Donnée_pour_" + id;
}
private Integer effectuerCalculLourd(int valeur) {
long sum = 0;
for (int i = 0; i < 1000000; i++) {
sum += Math.sqrt(i * valeur);
}
return (int) (sum % 1000); // Retourne un résultat arbitraire
}
// Méthode de nettoyage des ressources
public void arreterPools() {
operateurIO.shutdown();
operateurCPU.shutdown();
try {
if (!operateurIO.awaitTermination(5, TimeUnit.SECONDS)) {
operateurIO.shutdownNow();
}
if (!operateurCPU.awaitTermination(5, TimeUnit.SECONDS)) {
operateurCPU.shutdownNow();
}
} catch (InterruptedException e) {
operateurIO.shutdownNow();
operateurCPU.shutdownNow();
Thread.currentThread().interrupt();
}
}
public static void main(String[] args) throws ExecutionException, InterruptedException {
GestionnairePoolsCorrecte gestionnaire = new GestionnairePoolsCorrecte();
// Exemple d'utilisation
CompletableFuture<String> futureIO = gestionnaire.traiterAvecPoolIO("utilisateurA");
CompletableFuture<Integer> futureCPU = gestionnaire.calculerAvecPoolCPU(123);
System.out.println("Résultat IO: " + futureIO.get());
System.out.println("Résultat CPU: " + futureCPU.get());
gestionnaire.arreterPools();
}
}
Gestion des exceptions : Pourquoi les exceptions disparaissent-elles ?
Un autre défi avec CompletableFuture est la gestion des exceptions, qui peuvent sembler « disparaître » si elles ne sont pas traitées explicitement.
Scénario typique de perte d'exception
Dans les chaînes de CompletableFuture, une expection non gérée dans une étape peut empêcher les étapes suivantes de s'exécuter, sans pour autant remonter l'exception immédiatement au point d'appel, surtout si la méthode get() n'est pas appelée ou si elle est appelée sans gestion de l'ExecutionException.
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
public class DisparitionExceptions {
public void demontrerPerteExceptionSimple() {
CompletableFuture<String> futur = CompletableFuture.supplyAsync(() -> {
System.out.println("Exécution de l'opération risquée...");
if (Math.random() > 0.5) {
throw new IllegalStateException("Opération métier échouée !");
}
return "Opération réussie";
});
// La chaîne de transformation ne sera pas exécutée en cas d'exception
CompletableFuture<String> resultatChaine = futur.thenApply(res -> {
System.out.println("Résultat traité: " + res);
return res.toUpperCase();
});
try {
// L'exception sera encapsulée dans une ExecutionException ici.
// Si .get() n'est pas appelé, l'exception ne remonte pas.
String finalResult = resultatChaine.get();
System.out.println("Résultat final: " + finalResult);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} catch (ExecutionException e) {
System.err.println("Exception capturée: " + e.getClass().getName());
System.err.println("Cause racine: " + e.getCause().getMessage());
}
}
// Une perte d'exception encore plus discrète
public void demontrerPerteExceptionSilencieuse() {
CompletableFuture.supplyAsync(() -> {
System.out.println("Démarrage d'une tâche critique...");
throw new CustomBusinessException("Erreur critique non gérée.");
}).thenAccept(donnee -> {
// Cette étape ne sera jamais atteinte si une exception est levée en amont
System.out.println("Traitement de la donnée: " + donnee);
});
// Le programme continue, l'exception est ignorée car .get() n'est pas appelé
System.out.println("Le programme principal se termine, mais une exception a pu être perdue.");
try {
TimeUnit.MILLISECONDS.sleep(100); // Laisse le temps au thread de s'exécuter
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
static class CustomBusinessException extends RuntimeException {
public CustomBusinessException(String message) {
super(message);
}
}
public static void main(String[] args) {
DisparitionExceptions demo = new DisparitionExceptions();
System.out.println("--- Démonstration de la perte simple ---");
demo.demontrerPerteExceptionSimple();
System.out.println("\n--- Démonstration de la perte silencieuse ---");
demo.demontrerPerteExceptionSilencieuse();
}
}
Mécanismes de gestion des exceptions de CompletableFuture
Les exceptions dans CompletableFuture sont propagées le long de la chaîne de calcul et sont encapsulées dans un CompletionException ou une ExecutionException si l'on appelle get().
exceptionally(Function<Throwable, T> fn): Permet de récupérer après une exception et de fournir une valeur de secours.handle(BiFunction<T, Throwable, R> fn): Similaire àexceptionallymais est toujours appelé, que le calcul ait réussi ou échoué, et permet de transformer le résultat ou l'exception.whenComplete(BiConsumer<T, Throwable> action): Exécute une action lorsque laCompletableFutureest terminée (réussie ou échouée), sans modifier le résultat ou l'exception. Utile pour la journalisation ou les effets de bord.
Méthodes correctes de gestion des exceptions
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.function.BiConsumer;
import java.util.function.Function;
public class CorrecteurExceptions {
// Simule une opération qui peut échouer
private String effectuerOperationRisquee() {
if (Math.random() < 0.6) {
throw new RuntimeException("Erreur lors de l'opération simulée !");
}
return "Donnée traitée";
}
// Méthode 1: Utilisation de exceptionally pour la récupération
public CompletableFuture<String> gererAvecRecuperation() {
return CompletableFuture.supplyAsync(this::effectuerOperationRisquee)
.exceptionally(erreur -> {
System.err.println("Opération échouée, récupération avec valeur par défaut: " + erreur.getMessage());
return "VALEUR_DEFAUT_EXCEPTION";
});
}
// Méthode 2: Utilisation de handle pour un traitement unifié (résultat ou exception)
public CompletableFuture<String> gererAvecTraitementUnifie() {
return CompletableFuture.supplyAsync(this::effectuerOperationRisquee)
.handle((resultat, erreur) -> {
if (erreur != null) {
System.err.println("Gestion unifiée: Une erreur est survenue: " + erreur.getMessage());
return "VALEUR_ERREUR_UNIFIEE";
}
return resultat + " (traité unifié)";
});
}
// Méthode 3: Utilisation de whenComplete pour les effets de bord (journalisation, alerte)
public CompletableFuture<Void> gererAvecEffetsDeBord() {
return CompletableFuture.supplyAsync(this::effectuerOperationRisquee)
.whenComplete((resultat, erreur) -> {
if (erreur != null) {
journaliserErreur(erreur);
notifierAdmin(erreur);
} else {
System.out.println("Opération réussie: " + resultat);
}
})
.thenApply(res -> res + " (après effet de bord)") // Continue la chaîne même après whenComplete
.thenAccept(System.out::println); // Consomme le résultat
}
// Méthode 4: Gestion des exceptions dans des opérations composées
public CompletableFuture<String> gererEnComposition() {
CompletableFuture<String> etape1 = CompletableFuture.supplyAsync(() -> {
// throw new CustomBusinessException("Erreur dans l'étape 1");
return "ResultatEtape1";
});
CompletableFuture<String> etape2 = etape1.thenCompose(res1 -> {
// throw new TimeoutException("Délai dépassé dans l'étape 2");
return CompletableFuture.supplyAsync(() -> res1 + "ResultatEtape2");
});
// Gérer l'exception à la fin de la chaîne pour capturer toutes les erreurs en amont
return etape2.exceptionally(erreur -> {
Throwable causeRacine = obtenirCauseRacine(erreur);
if (causeRacine instanceof CustomBusinessException) {
System.err.println("Exception métier détectée: " + causeRacine.getMessage());
return "Repli_Metier";
} else if (causeRacine instanceof TimeoutException) {
System.err.println("Délai dépassé: " + causeRacine.getMessage());
return "Repli_Delai";
} else {
System.err.println("Erreur inconnue: " + causeRacine.getMessage());
return "Repli_General";
}
});
}
private void journaliserErreur(Throwable t) {
System.err.println("LOG ERREUR: " + t.getMessage());
}
private void notifierAdmin(Throwable t) {
System.out.println("ALERTE ADMIN: " + t.getMessage());
}
private Throwable obtenirCauseRacine(Throwable throwable) {
Throwable cause = throwable;
while (cause.getCause() != null && cause.getCause() != cause) {
cause = cause.getCause();
}
return cause;
}
static class CustomBusinessException extends RuntimeException {
public CustomBusinessException(String message) { super(message); }
}
public static void main(String[] args) throws ExecutionException, InterruptedException {
CorrecteurExceptions demo = new CorrecteurExceptions();
System.out.println("\n--- Démo recuperation ---");
System.out.println(demo.gererAvecRecuperation().get());
System.out.println("\n--- Démo unifie ---");
System.out.println(demo.gererAvecTraitementUnifie().get());
System.out.println("\n--- Démo effets de bord ---");
demo.gererAvecEffetsDeBord().join(); // Utilise join pour attendre la fin
System.out.println("\n--- Démo composition ---");
// Pour tester l'exception, décommenter les throws dans gererEnComposition
System.out.println(demo.gererEnComposition().get());
}
}
L'enfer des callbacks : Quand l'asynchrone devient une source de douleur
L'utilisation intensive de thenCompose ou thenApply, surtout dans des logiques séquentielles complexes, peut rapidement mener à un "callback hell", rendant le code difficile à lire, à comprendre et à maintenir.
Exemple d'enfer des callbacks
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
public class CalamiteCallbacks {
// Classes de données simplifiées
record Client(String id, String nom) {}
record Article(String id, String nom, double prix) {}
record Commande(String id, String clientId, List<Article> articles) {}
record Reduction(double pourcentage) {}
// Services simulés
CompletableFuture<Client> obtenirInfosClient(String idClient) {
return CompletableFuture.supplyAsync(() -> {
System.out.println("Récupération infos client: " + idClient);
simulerLatence(50);
return new Client(idClient, "NomClient_" + idClient);
});
}
CompletableFuture<List<Commande>> obtenirHistoriqueCommandes(String idClient) {
return CompletableFuture.supplyAsync(() -> {
System.out.println("Récupération historique commandes pour: " + idClient);
simulerLatence(80);
return List.of(new Commande("CMD1", idClient, List.of(new Article("ART1", "ProduitA", 100.0))));
});
}
CompletableFuture<Reduction> calculerReductionClient(Client client, List<Commande> historique) {
return CompletableFuture.supplyAsync(() -> {
System.out.println("Calcul réduction pour: " + client.nom());
simulerLatence(30);
return new Reduction(historique.size() > 0 ? 0.10 : 0.05);
});
}
CompletableFuture<Commande> creerNouvelleCommande(Client client, Reduction reduction, List<Article> articles) {
return CompletableFuture.supplyAsync(() -> {
System.out.println("Création nouvelle commande pour: " + client.nom());
simulerLatence(70);
return new Commande("NOUVELLE_CMD", client.id(), articles);
});
}
CompletableFuture<String> envoyerConfirmation(Commande commande) {
return CompletableFuture.supplyAsync(() -> {
System.out.println("Envoi confirmation pour commande: " + commande.id());
simulerLatence(40);
return "Confirmation envoyée pour " + commande.id();
});
}
private void simulerLatence(int ms) {
try {
TimeUnit.MILLISECONDS.sleep(ms);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
// Le scénario de "callback hell"
public CompletableFuture<String> orchestrerProcessusCommande(String idClient, List<Article> articles) {
return obtenirInfosClient(idClient)
.thenCompose(client ->
obtenirHistoriqueCommandes(client.id())
.thenCompose(historique ->
calculerReductionClient(client, historique)
.thenCompose(reduction ->
creerNouvelleCommande(client, reduction, articles)
.thenCompose(nouvelleCommande ->
envoyerConfirmation(nouvelleCommande)
)
)
)
);
}
public static void main(String[] args) throws Exception {
CalamiteCallbacks demo = new CalamiteCallbacks();
List<Article> articlesPanier = List.of(new Article("P1", "Produit Test", 25.0));
System.out.println("Début de l'orchestration...");
String resultat = demo.orchestrerProcessusCommande("client123", articlesPanier).get();
System.out.println("Processus terminé: " + resultat);
}
}
Solutions pour une programmation asynchrone structurée
Pour éviter l'enfer des callbacks, on peut adopter plusieurs stratégies :
- **Utiliser un objet de contexte mutable :** Passer un objet de contexte à travers les étapes de la chaîne permet de conserver les résultats intermédiaires et d'éviter de redemander des informations ou de les passer en paramètres à chaque étape.
- **Combiner des tâches parallèles :** Utiliser
thenCombine,thenAcceptBoth,thenRunpour les tâches qui peuvent s'exécuter en parallèle. - **Attendre plusieurs tâches indépendantes :** Utiliser
allOfouanyOfpour synchroniser plusieurs futures.
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.TimeUnit;
public class ProgrammationAsynchroneStructurée {
// Classes de données simplifiées (utilisées aussi dans l'exemple de CalamiteCallbacks)
record Client(String id, String nom) {}
record Article(String id, String nom, double prix) {}
record Commande(String id, String clientId, List<Article> articles) {}
record Reduction(double pourcentage) {}
// Nouvelles classes de données
record Adresse(String rue, String ville) {}
record Preferences(String langue) {}
record ProfilUtilisateur(Client client, List<Commande> historique, List<Adresse> adresses, Preferences preferences) {}
// Services simulés (réutilisés de CalamiteCallbacks pour la cohérence)
CompletableFuture<Client> obtenirInfosClient(String idClient) { /* ... */ return CompletableFuture.supplyAsync(() -> { simulerLatence(50); return new Client(idClient, "NomClient_" + idClient); }); }
CompletableFuture<List<Commande>> obtenirHistoriqueCommandes(String idClient) { /* ... */ return CompletableFuture.supplyAsync(() -> { simulerLatence(80); return List.of(new Commande("CMD1", idClient, List.of(new Article("ART1", "ProduitA", 100.0)))); });}
CompletableFuture<Reduction> calculerReductionClient(Client client, List<Commande> historique) { /* ... */ return CompletableFuture.supplyAsync(() -> { simulerLatence(30); return new Reduction(historique.size() > 0 ? 0.10 : 0.05); });}
CompletableFuture<Commande> creerNouvelleCommande(Client client, Reduction reduction, List<Article> articles) { /* ... */ return CompletableFuture.supplyAsync(() -> { simulerLatence(70); return new Commande("NOUVELLE_CMD", client.id(), articles); });}
CompletableFuture<String> envoyerConfirmation(Commande commande) { /* ... */ return CompletableFuture.supplyAsync(() -> { simulerLatence(40); return "Confirmation envoyée pour " + commande.id(); });}
CompletableFuture<List<Adresse>> obtenirAdressesClient(String idClient) { /* ... */ return CompletableFuture.supplyAsync(() -> { simulerLatence(60); return List.of(new Adresse("10 Rue Fictive", "Ville Test")); });}
CompletableFuture<Preferences> obtenirPreferencesClient(String idClient) { /* ... */ return CompletableFuture.supplyAsync(() -> { simulerLatence(25); return new Preferences("fr"); });}
CompletableFuture<List<String>> obtenirNotifications(String idClient) { /* ... */ return CompletableFuture.supplyAsync(() -> { simulerLatence(45); return List.of("Nouvelle promotion !", "Votre commande est en route."); });}
private void simulerLatence(int ms) {
try {
TimeUnit.MILLISECONDS.sleep(ms);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
// Objet de contexte pour un workflow structuré
static class ContexteWorkflowCommande {
String idClient;
Client client;
List<Commande> historiqueCommandes;
Reduction reduction;
Commande commandeNouvelle;
List<Article> articlesPanier; // Ajouté pour le panier
public ContexteWorkflowCommande(String idClient, List<Article> articlesPanier) {
this.idClient = idClient;
this.articlesPanier = articlesPanier;
}
}
// Processus de commande structuré
public CompletableFuture<String> orchestrerCommandeStructurée(String idClient, List<Article> articles) {
ContexteWorkflowCommande ctx = new ContexteWorkflowCommande(idClient, articles);
return obtenirInfosClient(ctx.idClient)
.thenCompose(client -> {
ctx.client = client;
return obtenirHistoriqueCommandes(client.id());
})
.thenCompose(historique -> {
ctx.historiqueCommandes = historique;
return calculerReductionClient(ctx.client, historique);
})
.thenCompose(reduction -> {
ctx.reduction = reduction;
return creerNouvelleCommande(ctx.client, reduction, ctx.articlesPanier);
})
.thenCompose(nouvelleCommande -> {
ctx.commandeNouvelle = nouvelleCommande;
return envoyerConfirmation(nouvelleCommande);
})
.exceptionally(ex -> {
System.err.println("Erreur globale dans le processus de commande: " + ex.getMessage());
return "Échec du processus de commande.";
});
}
// Utilisation de thenCombine pour les tâches parallèles
public CompletableFuture<ProfilUtilisateur> obtenirProfilUtilisateur(String idClient) {
CompletableFuture<Client> clientFuture = obtenirInfosClient(idClient);
CompletableFuture<List<Commande>> historiqueFuture = obtenirHistoriqueCommandes(idClient);
CompletableFuture<List<Adresse>> adressesFuture = obtenirAdressesClient(idClient);
CompletableFuture<Preferences> prefsFuture = obtenirPreferencesClient(idClient);
// Combine les résultats de quatre futures en parallèle
return CompletableFuture.allOf(clientFuture, historiqueFuture, adressesFuture, prefsFuture)
.thenApply(v -> { // 'v' est Void car allOf ne renvoie pas de résultat
try {
Client client = clientFuture.get();
List<Commande> historique = historiqueFuture.get();
List<Adresse> adresses = adressesFuture.get();
Preferences prefs = prefsFuture.get();
return new ProfilUtilisateur(client, historique, adresses, prefs);
} catch (Exception e) {
throw new CompletionException("Erreur lors de la combinaison des données du profil", e);
}
});
}
// Utilisation de allOf pour exécuter plusieurs tâches indépendantes et attendre leur complétion
public CompletableFuture<Map<String, Object>> obtenirDonneesTableauDeBord(String idClient) {
CompletableFuture<Client> infosClientFuture = obtenirInfosClient(idClient);
CompletableFuture<List<Commande>> commandesFuture = obtenirHistoriqueCommandes(idClient);
CompletableFuture<List<String>> notificationsFuture = obtenirNotifications(idClient);
// Attendre que toutes les futures soient terminées
return CompletableFuture.allOf(infosClientFuture, commandesFuture, notificationsFuture)
.thenApply(v -> {
Map<String, Object> tableauDeBord = new HashMap<>();
try {
tableauDeBord.put("client", infosClientFuture.get());
tableauDeBord.put("commandes", commandesFuture.get());
tableauDeBord.put("notifications", notificationsFuture.get());
} catch (Exception e) {
throw new CompletionException("Échec de la récupération des données du tableau de bord", e);
}
return tableauDeBord;
});
}
public static void main(String[] args) throws Exception {
ProgrammationAsynchroneStructurée demo = new ProgrammationAsynchroneStructurée();
List<Article> panier = List.of(new Article("itemA", "Article A", 50.0));
System.out.println("--- Démonstration de l'orchestration structurée ---");
String resultatCommande = demo.orchestrerCommandeStructurée("user456", panier).get();
System.out.println("Résultat final de la commande: " + resultatCommande);
System.out.println("\n--- Démonstration de la récupération de profil avec thenCombine et allOf ---");
ProfilUtilisateur profil = demo.obtenirProfilUtilisateur("user789").get();
System.out.println("Profil de l'utilisateur: " + profil);
System.out.println("\n--- Démonstration de la récupération des données du tableau de bord ---");
Map<String, Object> tableau = demo.obtenirDonneesTableauDeBord("user001").get();
System.out.println("Données du tableau de bord: " + tableau);
}
}
Fuites de mémoire : Les consommateurs de ressources cachés
Une utilisation négligente de CompletableFuture peut entraîner des fuites de mémoire, en particulier dans les applications à longue durée de vie.
Scénarios courants de fuites de mémoire
- **Cache à croissance illimitée :** Si des
CompletableFuturesont stockés dans un cache sans mécanisme d'expiration, ils peuvent s'accumuler indéfiniment. - **Futures non complétés :** Des tâches qui ne se terminent jamais (ex: appels réseau bloqués indéfiniment) peuvent faire que les
CompletableFutureassociés restent en mémoire, avec leurs références et celles des objets qu'ils contiennent. - **Références circulaires :** Un
CompletableFutureet l'objet qui le gère se référant mutuellement peuvent empêcher la collecte de garbage.
import java.lang.ref.WeakReference;
import java.util.Collections;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
public class DemonstreFuiteMemoire {
private final Map<String, CompletableFuture<String>> cacheFuturesBloquants = new ConcurrentHashMap<>();
private final ExecutorService longRunningExecutor = Executors.newFixedThreadPool(5);
// Scenario 1: Cache non limité de CompletableFuture
public CompletableFuture<String> recupererDonneeAvecFuite(String cle) {
// computeIfAbsent peut créer une référence forte persistante si la future n'est jamais nettoyée
return cacheFuturesBloquants.computeIfAbsent(cle, k -> {
System.out.println("Chargement de " + k + " dans le cache.");
return CompletableFuture.supplyAsync(() -> simulerChargementDonnee(k));
});
}
private String simulerChargementDonnee(String cle) {
try {
TimeUnit.MILLISECONDS.sleep(100); // Tâche rapide, mais la future reste si non invalidée
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
return "DonneeDe_" + cle;
}
// Scenario 2: Accumulation de futures jamais complétées
public void accumulerFuturesInachevees() {
for (int i = 0; i < 100_000; i++) {
CompletableFuture<String> futureLongueDuree = CompletableFuture.supplyAsync(() -> {
try {
// Simule une tâche bloquante infinie
TimeUnit.DAYS.sleep(365);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
return "Résultat_long";
}, longRunningExecutor);
// Ces futures ne se termineront jamais et retiendront de la mémoire.
// On pourrait les stocker implicitement ou explicitement dans une collection.
}
System.out.println("Créé 100,000 futures potentiellement bloquées.");
}
// Scenario 3: Référence circulaire (simplifié)
public class GestionnaireTacheDangereux {
private CompletableFuture<String> tacheEnCours;
private String etat = "INIT"; // État interne du gestionnaire
public void demarrerTache() {
tacheEnCours = CompletableFuture.supplyAsync(() -> {
// La lambda capture 'this' (GestionnaireTacheDangereux), créant une dépendance
while (!"TERMINE".equals(etat)) {
System.out.println("Tâche en cours... état: " + etat);
try {
TimeUnit.MILLISECONDS.sleep(50);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
return "Tâche complétée.";
});
}
public CompletableFuture<String> obtenirTacheEnCours() {
return tacheEnCours;
}
public void setEtat(String nouvelEtat) {
this.etat = nouvelEtat;
}
}
public static void main(String[] args) throws Exception {
DemonstreFuiteMemoire demo = new DemonstreFuiteMemoire();
// Démo Scenario 1
System.out.println("\n--- Démo Fuite Cache ---");
for (int i = 0; i < 5; i++) {
demo.recupererDonneeAvecFuite("cle" + i).join();
}
System.out.println("Taille du cache après 5 clés: " + demo.cacheFuturesBloquants.size());
// Sans nettoyage, ces 5 futures restent en mémoire.
// Démo Scenario 2
System.out.println("\n--- Démo Futures Inachevees ---");
// demo.accumulerFuturesInachevees(); // Attention: ceci va consommer beaucoup de mémoire!
// Pour une démo réelle, il faudrait surveiller l'usage mémoire.
// Démo Scenario 3
System.out.println("\n--- Démo Référence Circulaire ---");
GestionnaireTacheDangereux gestionnaire = demo.new GestionnaireTacheDangereux();
gestionnaire.demarrerTache();
TimeUnit.MILLISECONDS.sleep(200);
gestionnaire.setEtat("TERMINE");
System.out.println("Résultat tâche: " + gestionnaire.obtenirTacheEnCours().get());
// Après la fin de la tâche, si le gestionnaire n'est plus utilisé ailleurs,
// la référence circulaire pourrait potentiellement prolonger la vie de l'objet.
demo.longRunningExecutor.shutdownNow(); // Arrêter le pool pour la démo 2
}
}
Détection et prévention des fuites de mémoire
Pour prévenir les fuites de mémoire :
- **Utiliser des caches avec expiration :** Des bibliothèques comme Caffeine ou Guava Cache offrent des mécanismes d'expiration et de taille maximale.
- **Appliquer des délais d'attente :** Utiliser
orTimeoutoucompleteOnTimeoutpour les futures, afin d'éviter qu'elles ne restent bloquées indéfiniment. - **Nettoyer les références :** S'assurer que les références aux
CompletableFuturene sont pas conservées inutilement après leur achèvement. Pour les références circulaires, envisager desWeakReferencesi le cycle de vie est géré par la GC. - **Outils de monitoring :** Utiliser des profilers Java (VisualVM, JProfiler, YourKit) pour identifier les objets qui occupent la mémoire et leur chemin de référence.
import com.github.benmanes.caffeine.cache.Cache;
import com.github.benmanes.caffeine.cache.Caffeine;
import com.github.benmanes.caffeine.cache.RemovalCause;
import java.lang.ref.WeakReference;
import java.util.concurrent.*;
import java.util.function.Function;
public class ProtectionFuiteMemoire {
// Utilisation de Caffeine Cache avec expiration et gestion des futures
private final Cache<String, CompletableFuture<String>> cacheSecurise;
private final ScheduledExecutorService moniteurFutures;
private final ExecutorService tacheExecutor = Executors.newCachedThreadPool();
public ProtectionFuiteMemoire() {
this.cacheSecurise = Caffeine.newBuilder()
.maximumSize(500) // Taille maximale du cache
.expireAfterAccess(10, TimeUnit.MINUTES) // Expiration après un certain temps d'inactivité
.removalListener((String key, CompletableFuture<String> future, RemovalCause cause) -> {
if (future != null && !future.isDone()) {
System.out.println("Cache: Annulation de la future non terminée pour clé " + key + ", cause: " + cause);
future.cancel(true); // Tente d'annuler la tâche sous-jacente
}
})
.build();
this.moniteurFutures = Executors.newSingleThreadScheduledExecutor();
demarrerMoniteurFutures();
}
// Méthode sécurisée pour récupérer des données avec cache
public CompletableFuture<String> recupererDonneeSecurisee(String cle) {
try {
// Utilise get(key, mappingFunction) de Caffeine pour charger si absent
return cacheSecurise.get(cle, k -> {
CompletableFuture<String> nouvelleFuture = CompletableFuture.supplyAsync(() -> simulerChargementDonneeLong(k), tacheExecutor);
// Ajouter un délai d'attente à la future elle-même
return nouvelleFuture.orTimeout(30, TimeUnit.SECONDS)
.exceptionally(ex -> {
System.err.println("Erreur ou timeout pour la clé " + k + ": " + ex.getMessage());
cacheSecurise.invalidate(k); // Invalider l'entrée du cache en cas d'erreur ou timeout
return "Donnée de secours pour " + k;
});
});
} catch (Exception e) { // ExecutionException de get, par exemple si mappingFunction échoue
throw new RuntimeException("Échec de récupération pour " + cle, e);
}
}
private String simulerChargementDonneeLong(String cle) {
try {
TimeUnit.SECONDS.sleep(new Random().nextInt(10) + 1); // Simule une tâche longue (1-10s)
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return "Interrompu_" + cle;
}
return "Chargé_" + cle;
}
// Utilisation de WeakReference pour éviter les cycles de référence explicites
public static class GestionnaireTacheSecurise {
private WeakReference<CompletableFuture<String>> refTacheActuelle;
public void demarrerTacheSecurisee() {
CompletableFuture<String> tache = CompletableFuture.supplyAsync(() -> {
System.out.println("Exécution d'une tâche sécurisée.");
return "Tâche sécurisée terminée.";
});
refTacheActuelle = new WeakReference<>(tache);
// Lorsque la tâche est complète, la référence faible peut être nettoyée par la GC
tache.whenComplete((result, ex) -> {
if (ex == null) {
System.out.println("Tâche sécurisée terminée. Résultat: " + result);
} else {
System.err.println("Tâche sécurisée échouée. Erreur: " + ex.getMessage());
}
// Optionnel: nettoyer explicitement la WeakReference si nécessaire
refTacheActuelle = null;
});
}
public CompletableFuture<String> obtenirTacheSecurisee() {
return refTacheActuelle != null ? refTacheActuelle.get() : null;
}
}
// Moniteur périodique des futures non complétées (pour diagnostic)
private void demarrerMoniteurFutures() {
moniteurFutures.scheduleAtFixedRate(() -> {
long futuresNonTerminees = cacheSecurise.asMap().values().stream()
.filter(f -> !f.isDone())
.count();
if (futuresNonTerminees > 0) {
System.out.println("MONITEUR: " + futuresNonTerminees + " futures non terminées dans le cache.");
}
}, 0, 1, TimeUnit.MINUTES); // Vérifie toutes les minutes
}
// Nettoyage des ressources
public void arreter() {
moniteurFutures.shutdownNow();
tacheExecutor.shutdownNow();
}
public static void main(String[] args) throws Exception {
ProtectionFuiteMemoire demo = new ProtectionFuiteMemoire();
System.out.println("--- Démo Cache Sécurisé ---");
demo.recupererDonneeSecurisee("prod1").join();
demo.recupererDonneeSecurisee("prod2").join();
System.out.println("Tentative de récupération d'une donnée qui prend du temps (potentiel timeout)...");
demo.recupererDonneeSecurisee("long_task_prod").join();
System.out.println("Cache size: " + demo.cacheSecurise.estimatedSize());
System.out.println("\n--- Démo Gestionnaire Tache Sécurisé ---");
GestionnaireTacheSecurise gestionnaireSecurise = new GestionnaireTacheSecurise();
gestionnaireSecurise.demarrerTacheSecurisee();
CompletableFuture<String> tache = gestionnaireSecurise.obtenirTacheSecurisee();
if (tache != null) {
System.out.println("Résultat de la tâche sécurisée: " + tache.get());
}
demo.arreter();
}
}
Absence de contrôle de délai d'attente
L'omission de contrôles de délai d'attente (timeout) dans les opérations asynchrones est une source fréquente de blocage de threads, d'épuisement des ressources et d'expériences utilisateur dégradées.
La gravité des problèmes de délai d'attente
Sans délai, une opération bloquante (comme un appel réseau qui ne répond pas) peut retenir un thread indéfiniment. Dans un pool de threads, cela signifie qu'un thread est perdu pour d'autres tâches, pouvant mener à la famine des threads et au blocage complet de l'application.
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
public class RisquesSansDelai {
private final ExecutorService execExterne = Executors.newFixedThreadPool(2);
// Code dangereux: pas de contrôle de délai d'attente sur .get()
public String recuperationDangereuse() {
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
System.out.println("Appel à un service externe potentiellement bloquant...");
return appelServiceBloquant();
}, execExterne);
try {
// Si appelServiceBloquant() ne se termine jamais, ce .get() bloquera indéfiniment.
System.out.println("Tentative de récupération du résultat...");
return future.get();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return "interrompu";
} catch (ExecutionException e) {
System.err.println("Erreur lors de l'exécution: " + e.getCause().getMessage());
return "erreur";
}
}
private String appelServiceBloquant() {
try {
TimeUnit.DAYS.sleep(365); // Simule un blocage quasi infini
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
System.out.println("Appel de service externe interrompu.");
}
return "réponse";
}
// Exemple de fuite de ressources due à des tâches bloquées dans un pool limité
public void exempleFuiteRessources() {
ExecutorService poolLimité = Executors.newFixedThreadPool(5);
AtomicInteger tachesSoumises = new AtomicInteger(0);
for (int i = 0; i < 10; i++) { // Soumet plus de tâches que le pool ne peut gérer simultanément
CompletableFuture.runAsync(() -> {
int numeroTache = tachesSoumises.incrementAndGet();
System.out.println("Tâche " + numeroTache + " démarrée dans thread " + Thread.currentThread().getName());
try {
TimeUnit.SECONDS.sleep(30); // Tâche longue et bloquante
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
System.out.println("Tâche " + numeroTache + " interrompue.");
}
System.out.println("Tâche " + numeroTache + " terminée (ou interrompue).");
}, poolLimité);
}
System.out.println("10 tâches soumises. Le pool de 5 threads sera rapidement saturé.");
poolLimité.shutdown(); // Tente d'arrêter le pool, mais les tâches bloquées empêcheront la terminaison.
try {
if (!poolLimité.awaitTermination(60, TimeUnit.SECONDS)) {
System.err.println("Le pool n'a pas pu s'arrêter dans le délai imparti. Certaines tâches sont probablement bloquées.");
poolLimité.shutdownNow(); // Force l'arrêt, interrompant les threads
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
poolLimité.shutdownNow();
}
}
public void arreterExecuteurs() {
execExterne.shutdownNow();
}
public static void main(String[] args) throws Exception {
RisquesSansDelai demo = new RisquesSansDelai();
System.out.println("--- Démonstration de récupération dangereuse (peut bloquer) ---");
// String resultat = demo.recuperationDangereuse(); // Décommenter pour voir le blocage
// System.out.println("Résultat: " + resultat);
System.out.println("\n--- Démonstration de fuite de ressources ---");
demo.exempleFuiteRessources();
demo.arreterExecuteurs();
}
}
Solutions complètes de contrôle de délai d'attente
Java 9 a introduit des méthodes spécifiques pour la gestion des délais d'attente : orTimeout et completeOnTimeout. Pour Java 8, une approche manuelle est nécessaire.
import java.util.concurrent.*;
import java.util.function.Function;
import java.util.Random;
public class SolutionCompleteDelai {
private final ScheduledExecutorService planificateurDelai;
private final ExecutorService tachesExecutor = Executors.newCachedThreadPool();
public SolutionCompleteDelai() {
this.planificateurDelai = Executors.newScheduledThreadPool(2);
}
private String appelServiceExterne(String identifiant) {
try {
int duree = new Random().nextInt(6) + 1; // 1-6 secondes
System.out.println("Service externe " + identifiant + " va durer " + duree + "s");
TimeUnit.SECONDS.sleep(duree);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return "interrompu";
}
return "réponse de " + identifiant;
}
// Méthode 1: Utilisation de orTimeout (Java 9+)
public CompletableFuture<String> avecOuDelai() {
return CompletableFuture.supplyAsync(() -> appelServiceExterne("avecOuDelai"), tachesExecutor)
.orTimeout(3, TimeUnit.SECONDS) // Délai de 3 secondes
.exceptionally(throwable -> {
if (throwable instanceof TimeoutException) {
System.err.println("Opération 'avecOuDelai' a dépassé le délai!");
return "repli-timeout";
}
System.err.println("Opération 'avecOuDelai' a échoué: " + throwable.getMessage());
return "repli-erreur";
});
}
// Méthode 2: Utilisation de completeOnTimeout (Java 9+)
public CompletableFuture<String> avecCompleterSurDelai() {
return CompletableFuture.supplyAsync(() -> appelServiceExterne("avecCompleterSurDelai"), tachesExecutor)
.completeOnTimeout("valeur-defaut-timeout", 2, TimeUnit.SECONDS); // Fournit une valeur par défaut en cas de timeout
}
// Méthode 3: Contrôle de délai d'attente manuel (compatible Java 8)
public CompletableFuture<String> avecDelaiManuel() {
CompletableFuture<String> futurTache = CompletableFuture.supplyAsync(() -> appelServiceExterne("avecDelaiManuel"), tachesExecutor);
CompletableFuture<String> futurDelai = new CompletableFuture<>();
planificateurDelai.schedule(() -> {
futurDelai.completeExceptionally(new TimeoutException("Délai d'attente manuel dépassé"));
}, 4, TimeUnit.SECONDS);
// Applique l'un ou l'autre: le premier qui se termine gagne
return futurTache.applyToEither(futurDelai, Function.identity())
.exceptionally(throwable -> {
if (throwable instanceof TimeoutException) {
System.err.println("Opération 'avecDelaiManuel' a dépassé le délai!");
return "repli-manuel-timeout";
}
System.err.println("Opération 'avecDelaiManuel' a échoué: " + throwable.getMessage());
return "repli-manuel-erreur";
});
}
// Méthode 4: Contrôle de délai d'attente par couches (pour les workflows)
public CompletableFuture<String> avecDelaiEtage() {
return CompletableFuture.supplyAsync(() -> appelServiceExterne("phase1"), tachesExecutor)
.orTimeout(2, TimeUnit.SECONDS) // Délai pour la phase 1
.thenCompose(res1 -> CompletableFuture.supplyAsync(() -> appelServiceExterne("phase2_" + res1), tachesExecutor)
.orTimeout(3, TimeUnit.SECONDS) // Délai pour la phase 2
)
.thenCompose(res2 -> CompletableFuture.supplyAsync(() -> appelServiceExterne("phase3_" + res2), tachesExecutor)
.orTimeout(4, TimeUnit.SECONDS) // Délai pour la phase 3
)
.exceptionally(throwable -> {
Throwable causeRacine = obtenirCauseRacine(throwable);
if (causeRacine instanceof TimeoutException) {
System.err.println("Délai dépassé dans une des phases.");
return "repli-phase-timeout";
}
System.err.println("Erreur générale dans les phases: " + causeRacine.getMessage());
return "repli-general";
});
}
// Méthode 5: Stratégie de délai d'attente configurable
public CompletableFuture<String> avecDelaiConfigurable(String typeOperation) {
ConfigurationDelai config = obtenirConfigurationDelai(typeOperation);
return CompletableFuture.supplyAsync(() -> appelServiceExterne(typeOperation), tachesExecutor)
.orTimeout(config.getDelai(), config.getUniteTemps())
.exceptionally(config.getStrategieRepli());
}
private ConfigurationDelai obtenirConfigurationDelai(String typeOperation) {
return switch (typeOperation) {
case "rapide" -> new ConfigurationDelai(1, TimeUnit.SECONDS, t -> "repli-rapide");
case "normal" -> new ConfigurationDelai(5, TimeUnit.SECONDS, t -> "repli-normal");
case "lent" -> new ConfigurationDelai(15, TimeUnit.SECONDS, t -> "repli-lent");
default -> new ConfigurationDelai(8, TimeUnit.SECONDS, t -> "repli-defaut");
};
}
private Throwable obtenirCauseRacine(Throwable throwable) {
Throwable cause = throwable;
while (cause.getCause() != null && cause.getCause() != cause) {
cause = cause.getCause();
}
return cause;
}
// Classe pour la configuration du délai
public static class ConfigurationDelai {
private final long delai;
private final TimeUnit uniteTemps;
private final Function<Throwable, String> strategieRepli;
public ConfigurationDelai(long delai, TimeUnit uniteTemps, Function<Throwable, String> strategieRepli) {
this.delai = delai;
this.uniteTemps = uniteTemps;
this.strategieRepli = strategieRepli;
}
public long getDelai() { return delai; }
public TimeUnit getUniteTemps() { return uniteTemps; }
public Function<Throwable, String> getStrategieRepli() { return strategieRepli; }
}
// Méthode de nettoyage des ressources
public void arreter() {
planificateurDelai.shutdownNow();
tachesExecutor.shutdownNow();
}
public static void main(String[] args) throws ExecutionException, InterruptedException {
SolutionCompleteDelai demo = new SolutionCompleteDelai();
System.out.println("--- Démo orTimeout ---");
System.out.println(demo.avecOuDelai().get());
System.out.println("\n--- Démo completeOnTimeout ---");
System.out.println(demo.avecCompleterSurDelai().get());
System.out.println("\n--- Démo manuelTimeout ---");
System.out.println(demo.avecDelaiManuel().get());
System.out.println("\n--- Démo layeredTimeout ---");
System.out.println(demo.avecDelaiEtage().get());
System.out.println("\n--- Démo configurableTimeout (rapide) ---");
System.out.println(demo.avecDelaiConfigurable("rapide").get());
System.out.println("\n--- Démo configurableTimeout (lent) ---");
System.out.println(demo.avecDelaiConfigurable("lent").get());
demo.arreter();
}
}
Le contrôle des délais d'attente est essentiel pour la résilience des applications asynchrones.