Maîtriser les Opérations Stream en Java 8

Pourquoi le Concept de Stream est Essentiel en Java 8

Introduits avec Java 8, les Streams représentent une amélioration majeure pour le traitement des collections, se distinguant radicalement des flux d'entrée/sortie (java.io.InputStream, OutputStream) ou des approches XML comme StAX. Loin d'être des structures de données, les Streams sont des outils puissants dédiés aux opérations d'agrégation et de traitement massif de données sur des collections d'objets. Grâce aux expressions Lambda, également nouvelles en Java 8, l'API Stream améliore considérablement l'efficacité du développement et la lisibilité du code.

Un autre atout majeur est la prise en charge native des modes séquentiels et parallèles pour les opérations de collcete. Le mode parallèle tire parti des architectures multi-cœurs en utilisant le framework Fork/Join pour décomposer et accélérer le traitement des tâches. Écrire du code parallèle est souvent complexe et sujet aux erreurs, mais l'API Stream permet de créer des programmes concurrents performants sans écrire une seule ligne de code de gestion de threads. Le package java.util.stream est donc le résultat d'une convergence entre la programmation fonctionnelle et l'ère des processeurs multi-cœurs.

Qu'est-ce qu'une Opération d'Agrégation ?

Historiquement, les applications Java EE s'appuyaient souvent sur les bases de données relationnelles pour des calculs agrégés tels que le montant moyen dépensé par client mensuellement, l'article le plus cher en vente, ou le top 10 des produits les plus vendus. Cependant, à l'ère du big data, avec des sources de données variées et des volumes massifs, il est souvent nécessaire de réaliser des agrégations directement au niveau de l'application, en s'affranchissant des RDBMS ou en complétant leurs capacités.

L'API de collections traditionnelle de Java offrait peu de méthodes d'assistance pour ces scénarios, obligeant les développeurs à itérer manuellement sur les collections pour implémenter la logique d'agrégation. Cette approche est fastidieuse et peu efficace. Par exemple, pour trouver toutes les transactions d'un certain type et obtenir leurs identifiants triés par valeur décroissante en Java 7, il fallait écrire un code similaire à ceci :

Exemple 1 : Tri et filtrage avec Java 7

import java.util.ArrayList;
import java.util.Collections;
import java.util.Comparator;
import java.util.List;

class Payment {
    enum Category { FOOD, DRINKS, ELECTRONICS }
    private int id;
    private Category type;
    private double amount;

    public Payment(int id, Category type, double amount) {
        this.id = id;
        this.type = type;
        this.amount = amount;
    }

    public int getId() { return id; }
    public Category getType() { return type; }
    public double getAmount() { return amount; }
}

public class LegacyProcessing {
    public static void main(String[] args) {
        List<Payment> allPayments = new ArrayList<>();
        allPayments.add(new Payment(1, Payment.Category.FOOD, 150.75));
        allPayments.add(new Payment(2, Payment.Category.DRINKS, 25.50));
        allPayments.add(new Payment(3, Payment.Category.FOOD, 300.00));
        allPayments.add(new Payment(4, Payment.Category.ELECTRONICS, 500.00));
        allPayments.add(new Payment(5, Payment.Category.FOOD, 75.20));

        List<Payment> foodPayments = new ArrayList<>();
        for (Payment p : allPayments) {
            if (p.getType() == Payment.Category.FOOD) {
                foodPayments.add(p);
            }
        }

        Collections.sort(foodPayments, new Comparator<Payment>() {
            @Override
            public int compare(Payment p1, Payment p2) {
                return Double.compare(p2.getAmount(), p1.getAmount()); // Tri décroissant
            }
        });

        List<Integer> paymentIds = new ArrayList<>();
        for (Payment p : foodPayments) {
            paymentIds.add(p.getId());
        }
        System.out.println("IDs de paiements alimentaires (Java 7) : " + paymentIds);
    }
}

En Java 8, l'utilisation de l'API Stream permet un code plus concis et lisible, et l'exécution est plus rapide en mode parallèle.

Exemple 2 : Tri et filtrage avec Java 8 (API Stream)

import java.util.Arrays;
import java.util.List;
import java.util.Comparator;
import java.util.stream.Collectors;

// Payment class as defined above would be needed
// (Omitted for brevity in this snippet but assumed to be available)

public class StreamProcessing {
    public static void main(String[] args) {
        List<Payment> allPayments = Arrays.asList(
            new Payment(1, Payment.Category.FOOD, 150.75),
            new Payment(2, Payment.Category.DRINKS, 25.50),
            new Payment(3, Payment.Category.FOOD, 300.00),
            new Payment(4, Payment.Category.ELECTRONICS, 500.00),
            new Payment(5, Payment.Category.FOOD, 75.20)
        );

        List<Integer> foodPaymentIds = allPayments.parallelStream()
            .filter(p -> p.getType() == Payment.Category.FOOD)
            .sorted(Comparator.comparingDouble(Payment::getAmount).reversed()) // Tri par montant décroissant
            .map(Payment::getId)
            .collect(Collectors.toList());

        System.out.println("IDs de paiements alimentaires (Java 8 Stream) : " + foodPaymentIds);
    }
}

Introduction aux Flux (Streams)

Qu'est-ce qu'un Flux ?

Un flux Java n'est pas une collection de données ; il ne stocke aucune information. Il s'agit plutôt d'une séquence d'éléments sur laquelle des opérations peuvent être effectuées. C'est une sorte d'itérateur avancé : tandis qu'un itérateur classique requiert de traverser explicitement les éléments un par un, un flux permet de déclarer les opérations souhaitées (par exemple, "filtrer les chaînes de plus de 10 caractères", "obtenir la première lettre de chaque chaîne") et gère l'itération et la transformation des données de manière interne.

Un flux est unidirectionnel et non réutilisable : une fois qu'il a été parcouru, il est "consommé" et ne peut plus être utilisé, à l'image de l'eau qui s'écoule. Cependant, contrairement aux itérateurs impératifs et séquentiels, les flux peuvent être traités en parallèle. En mode séquentiel, chaque élément est traité l'un après l'autre. En mode parallèle, les données sont divisées en segments, chacun traité par un thread distinct, avant que les résultats ne soient combinés. Cette capacité de parallélisation repose sur le framework Fork/Join (JSR166y) introduit en Java 7, simplifiant considérablement l'écriture de code concurrent.

Une autre caractéristique notable des flux est qu'ils peuvent être issus de sources de données potentiellement infinies.

Composition d'un Pipeline de Flux

L'utilisation typique d'un flux se déroule en trois étapes :

  1. Acquisition d'une source de données (source).
  2. Transformation des données (opérations intermédiaires).
  3. Exécution d'une opération finale pour obtenir le résultat (opération terminale).

Chaque transformation produit un nouveau flux sans modifier l'original, ce qui permet d'enchaîner les opérations dans un "pipeline" de flux.

Diagramme d'un pipeline de flux Java 8Figure 1 : Structure d'un pipeline de flux

Les sources de flux peuvent être variées :

  • Collections et tableaux : Collection.stream(), Collection.parallelStream(), Arrays.stream(T[] array) ou Stream.of().
  • BufferedReader : java.io.BufferedReader.lines().
  • Factories statiques : java.util.stream.IntStream.range(), java.nio.file.Files.walk().
  • Construction manuelle : via java.util.Spliterator.
  • Autres : Random.ints(), BitSet.stream(), Pattern.splitAsStream(java.lang.CharSequence), JarFile.stream().

Les opérations sur les flux se classent en deux catégories :

  • Opérations Intermédiaires (Intermediate) : Elles peuvent être zéro ou plusieurs. Leur but est d'ouvrir le flux, d'effectuer des mappages ou des filtrages, puis de retourner un nouveau flux pour l'opération suivante. Ces opérations sont "paresseuses" (lazy) ; elles ne déclenchent pas le parcours réel du flux tant qu'une opération termianle n'est pas appelée.
  • Opération Terminale (Terminal) : Il ne peut y en avoir qu'une seule par flux. Une fois exécutée, le flux est "consommé". C'est cette opération qui initie le parcours du flux et produit un résultat ou un effet de bord.

La nature paresseuse des opérations intermédiaires signifie que plusieurs transformations sont combinées et exécutées en un seul passage lors de l'opération terminale, optimisant ainsi les performances. Les fonctions de transformation sont ajoutées à une "liste" et appliquées séquentiellement à chaque élément du flux lors de l'appel de l'opération terminale.

Certaines opérations sont dites "à court-circuit" (short-circuiting) :

  • Une opération intermédiaire est à court-circuit si elle accepte un flux potentiellement infini mais retourne un flux fini.
  • Une opération terminale est à court-circuit si elle accepte un flux potentiellement infini mais peut produire un résultat en temps fini.

Une opération à court-circuit est nécessaire (mais pas suffisante) pour traiter un flux infini en un temps limité.

Exemple 3 : Calcul du poids total d'articles rouges

import java.util.Arrays;
import java.util.List;

class Item {
    enum Color { RED, BLUE, GREEN }
    private String name;
    private Color color;
    private int weight;

    public Item(String name, Color color, int weight) {
        this.name = name;
        this.color = color;
        this.weight = weight;
    }

    public Color getColor() { return color; }
    public int getWeight() { return weight; }
}

public class ItemProcessing {
    public static void main(String[] args) {
        List<Item> items = Arrays.asList(
            new Item("Stapler", Item.Color.BLUE, 100),
            new Item("Pen", Item.Color.RED, 20),
            new Item("Notebook", Item.Color.GREEN, 300),
            new Item("Book", Item.Color.RED, 500)
        );

        int totalRedWeight = items.stream()
            .filter(i -> i.getColor() == Item.Color.RED) // Opération intermédiaire : filtrage
            .mapToInt(Item::getWeight)                   // Opération intermédiaire : transformation en IntStream
            .sum();                                      // Opération terminale : somme
        
        System.out.println("Poids total des articles rouges : " + totalRedWeight + "g");
    }
}

Dans cet exemple, stream() obtient la source, filter et mapToInt sont des opérations intermédiaires pour la sélection et la transformation, et sum() est l'opération terminale qui calcule le total du poids des articles rouges conformes.

Détails des Opérations de Flux

L'utilisation des flux peut être conceptualisée comme un processus de type "filtrer-mapper-réduire" pour produire un résultat final ou un effet de bord.

Construction et Conversion de Flux

Voici quelques méthodes courantes pour créer des flux :

Exemple 4 : Méthodes de construction de flux

import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import java.util.stream.Stream;

public class StreamCreation {
    public static void main(String[] args) {
        // 1. À partir de valeurs individuelles
        Stream<String> streamFromValues = Stream.of("pomme", "banane", "orange");
        streamFromValues.forEach(System.out::println);

        // 2. À partir d'un tableau
        String[] fruitArray = new String[] {"kiwi", "mangue", "poire"};
        Stream<String> streamFromArray1 = Stream.of(fruitArray);
        Stream<String> streamFromArray2 = Arrays.stream(fruitArray);
        streamFromArray1.forEach(System.out::println);
        streamFromArray2.forEach(System.out::println);

        // 3. À partir d'une collection
        List<String> fruitList = Arrays.asList(fruitArray);
        Stream<String> streamFromList = fruitList.stream();
        streamFromList.forEach(System.out::println);
    }
}

Il existe des versions spécialisées de flux pour les types primitifs : IntStream, LongStream et DoubleStream. Leur utilisation est recommandée pour éviter les coûts de boxing et unboxing associés aux flux de types enveloppés (Stream<Integer>, etc.).

Exemple 5 : Construction de flux numériques

import java.util.stream.IntStream;

public class PrimitiveStreams {
    public static void main(String[] args) {
        System.out.println("IntStream.of(array):");
        IntStream.of(new int[]{10, 20, 30}).forEach(System.out::println); // À partir d'un tableau d'int

        System.out.println("IntStream.range(inclusive, exclusive):");
        IntStream.range(5, 8).forEach(System.out::println); // De 5 à 7

        System.out.println("IntStream.rangeClosed(inclusive, inclusive):");
        IntStream.rangeClosed(10, 12).forEach(System.out::println); // De 10 à 12
    }
}

Les flux peuvent également être convertis en d'autres structures de données. Chaque flux ne peut être utilisé qu'une seule fois. Pour des raisons de concision, les exemples ci-dessous réutilisent des variables de flux, ce qui nécessiterait de recréer le flux à chaque fois en production.

Exemple 6 : Conversion de flux vers d'autres structures

import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import java.util.Set;
import java.util.Stack;
import java.util.stream.Collectors;
import java.util.stream.Stream;

public class StreamConversion {
    public static void main(String[] args) {
        Stream<String> currentStream = Stream.of("alpha", "beta", "gamma");

        // 1. En tableau
        String[] stringArray = currentStream.toArray(String[]::new);
        System.out.println("Array: " + Arrays.toString(stringArray));
        currentStream = Stream.of("alpha", "beta", "gamma"); // Recréer le flux

        // 2. En collection
        List<String> stringList = currentStream.collect(Collectors.toList());
        System.out.println("List: " + stringList);
        currentStream = Stream.of("alpha", "beta", "gamma");

        Set<String> stringSet = currentStream.collect(Collectors.toSet());
        System.out.println("Set: " + stringSet);
        currentStream = Stream.of("alpha", "beta", "gamma");

        Stack<String> stringStack = currentStream.collect(Collectors.toCollection(Stack::new));
        System.out.println("Stack: " + stringStack);
        currentStream = Stream.of("alpha", "beta", "gamma");

        // 3. En chaîne de caractères
        String combinedString = currentStream.collect(Collectors.joining(", "));
        System.out.println("Combined String: " + combinedString);
    }
}

Opérations sur les Flux

Une fois qu'une structure de données est enveloppée dans un flux, diverses opérations peuvent être appliquées. Elles sont généralement classées comme suit :

  • Intermédiaires : map (mapToInt, flatMap, etc.), filter, distinct, sorted, peek, limit, skip, parallel, sequential, unordered.
  • Terminales : forEach, forEachOrdered, toArray, reduce, collect, min, max, count, anyMatch, allMatch, noneMatch, findFirst, findAny, iterator.
  • À court-circuit : anyMatch, allMatch, noneMatch, findFirst, findAny, limit.

map et flatMap

L'opération map transforme chaque élément d'un flux d'entrée en un élément différent dans le flux de sortie, établissant une relation 1:1. Si vous êtes familier avec les langages fonctionnels comme Scala, cette fonction vous sera familière.

Exemple 7 : Conversion en majuscules

import java.util.Arrays;
import java.util.List;
import java.util.stream.Collectors;

public class MapExample {
    public static void main(String[] args) {
        List<String> words = Arrays.asList("hello", "world", "java");
        List<String> upperCaseWords = words.stream()
            .map(String::toUpperCase)
            .collect(Collectors.toList());
        System.out.println("Mots en majuscules : " + upperCaseWords); // [HELLO, WORLD, JAVA]
    }
}

Exemple 8 : Calcul des carrés

import java.util.Arrays;
import java.util.List;
import java.util.stream.Collectors;

public class MapSquares {
    public static void main(String[] args) {
        List<Integer> numbers = Arrays.asList(2, 3, 5, 7);
        List<Integer> squares = numbers.stream()
            .map(n -> n * n)
            .collect(Collectors.toList());
        System.out.println("Carrés des nombres : " + squares); // [4, 9, 25, 49]
    }
}

Pour les scénarios où un élément d'entrée doit correspondre à zéro ou plusieurs éléments de sortie (relation 1:N), flatMap est utilisé. Il "aplatit" une structure hiérarchique, extrayant les éléments de niveau inférieur et les regroupant dans un nouveau flux sans les structures intermédiaires.

Exemple 9 : Aplatir une liste de listes

import java.util.Arrays;
import java.util.List;
import java.util.stream.Collectors;
import java.util.stream.Stream;

public class FlatMapExample {
    public static void main(String[] args) {
        Stream<List<String>> nestedLists = Stream.of(
            Arrays.asList("A"),
            Arrays.asList("B", "C"),
            Arrays.asList("D", "E", "F")
        );
        Stream<String> flattenedStream = nestedLists
            .flatMap(List::stream); // Chaque List est transformée en un Stream, puis combinée
        flattenedStream.forEach(System.out::print); // ABCDEF
        System.out.println();
    }
}

filter

L'opération filter évalue chaque élément du flux par rapport à un prédicat et ne conserve que ceux qui satisfont la condition, formant ainsi un nouveau flux.

Exemple 10 : Filtrage des nombres pairs

import java.util.Arrays;
import java.util.stream.Stream;

public class FilterEvens {
    public static void main(String[] args) {
        Integer[] initialNumbers = {10, 11, 12, 13, 14, 15};
        Integer[] evenNumbers = Stream.of(initialNumbers)
            .filter(n -> n % 2 == 0) // Garde seulement les nombres pairs
            .toArray(Integer[]::new);
        System.out.println("Nombres pairs : " + Arrays.toString(evenNumbers)); // [10, 12, 14]
    }
}

Exemple 11 : Extraction de mots non vides d'un texte

import java.util.Arrays;
import java.util.List;
import java.util.stream.Collectors;
import java.util.stream.Stream;

public class FilterWords {
    public static void main(String[] args) {
        String textBlock = "Ceci est un exemple de texte. Il contient plusieurs mots.";
        List<String> words = Arrays.stream(textBlock.split("[\\s.,!?]+")) // Divise le texte en mots
            .filter(word -> word.length() > 0)                                // Élimine les chaînes vides
            .collect(Collectors.toList());
        System.out.println("Mots extraits : " + words);
    }
}

forEach et peek

La méthode forEach prend une expression Lambda et l'exécute sur chaque élément du flux. C'est une opération terminale.

Exemple 12 : Affichage des prénoms (Java 8 vs. Pré-Java 8)

import java.util.ArrayList;
import java.util.List;

class User {
    enum Gender { MALE, FEMALE }
    private String name;
    private Gender gender;

    public User(String name, Gender gender) {
        this.name = name;
        this.gender = gender;
    }
    public String getName() { return name; }
    public Gender getGender() { return gender; }
}

public class ForEachExample {
    public static void main(String[] args) {
        List<User> userList = new ArrayList<>();
        userList.add(new User("Alice", User.Gender.FEMALE));
        userList.add(new User("Bob", User.Gender.MALE));
        userList.add(new User("Charlie", User.Gender.MALE));
        userList.add(new User("Diana", User.Gender.FEMALE));

        System.out.println("Noms des hommes (Java 8 forEach):");
        userList.stream()
            .filter(u -> u.getGender() == User.Gender.MALE)
            .forEach(u -> System.out.println(u.getName()));

        System.out.println("\nNoms des hommes (Pré-Java 8 for-loop):");
        for (User u : userList) {
            if (u.getGender() == User.Gender.MALE) {
                System.out.println(u.getName());
            }
        }
    }
}

forEach est conçu pour les expressions Lambda et un style compact. Pour une optimisation multi-cœurs, parallelStream().forEach() peut être utilisé, bien que l'ordre des éléments ne soit pas garanti. Notez que forEach est une opération terminale et consomme le flux. Tenter d'utiliser un flux une seconde fois après un forEach entraînera une erreur.

Pour effectuer une opération sur chaque élément sans consommer le flux (par exemple, pour le débogage), l'opération intermédiaire peek est appropriée. Elle exécute une action sur chaque élément et retourne un nouveau flux.

Exemple 13 : Utilisation de peek pour le débogage

import java.util.List;
import java.util.stream.Collectors;
import java.util.stream.Stream;

public class PeekExample {
    public static void main(String[] args) {
        List<String> processedData = Stream.of("un", "deux", "trois", "quatre")
            .filter(e -> e.length() > 3)                                // Filtrage
            .peek(e -> System.out.println("Valeur filtrée: " + e))      // Effet de bord (debug)
            .map(String::toUpperCase)                                   // Transformation
            .peek(e -> System.out.println("Valeur mappée: " + e))       // Effet de bord (debug)
            .collect(Collectors.toList());
        System.out.println("Résultat final : " + processedData);
    }
}

Il est important de noter que les expressions Lambda dans forEach ne peuvent pas modifier les variables locales externes et ne peuvent pas utiliser break ou return pour interrompre la boucle.

findFirst et Optional

findFirst est une opération terminale et à court-circuit qui retourne le premier élément du flux, ou un Optional vide si le flux est vide. Optional est un conteneur qui peut contenir une valeur ou être vide, conçu pour éviter les exceptions NullPointerException.

Exemple 14 : Gestion des valeurs nulles avec Optional

import java.util.Optional;

public class OptionalUsage {
    // Méthode pour imprimer une chaîne en toute sécurité
    public static void safePrint(String text) {
        Optional.ofNullable(text).ifPresent(s -> System.out.println("Texte: " + s));
    }

    // Méthode pour obtenir la longueur d'une chaîne, retournant -1 si nulle
    public static int getSafeLength(String text) {
        return Optional.ofNullable(text)
                       .map(String::length)
                       .orElse(-1);
    }

    public static void main(String[] args) {
        String data1 = "  Java Stream  ";
        String data2 = null;
        String data3 = "";

        System.out.println("--- safePrint ---");
        safePrint(data1);
        safePrint(data2);
        safePrint(data3);

        System.out.println("\n--- getSafeLength ---");
        System.out.println("Longueur de '" + data1 + "': " + getSafeLength(data1));
        System.out.println("Longueur de '" + data2 + "': " + getSafeLength(data2));
        System.out.println("Longueur de '" + data3 + "': " + getSafeLength(data3));
    }
}

L'utilisation d'Optional améliore la lisibilité du code pour la gestion des valeurs nulles et permet des vérifications au moment de la compilatoin, réduisant ainsi les NPE (NullPointerException) à l'exécution. Des méthodes comme findAny, max/min, reduce, et IntStream.average() (retournant OptionalDouble) retournent des valeurs Optional.

reduce

Cette méthode combine les éléments d'un flux en un seul résultat. Elle prend un "élément initial" (identité) et une règle de combinaison (un BinaryOperator) pour agréger les éléments. Les opérations de concaténation de chaînes, de sommation, de recherche du minimum ou du maximum sont des cas spécifiques de reduce. Par exemple, la somme d'entiers peut être exprimée comme : integers.reduce(0, Integer::sum);.

Si aucun élément initial n'est fourni, la méthode combine les deux premiers éléments et retourne un Optional, car le flux pourrait être vide.

Exemple 15 : Cas d'utilisation de reduce

import java.util.Arrays;
import java.util.Optional;
import java.util.stream.Stream;

public class ReduceExample {
    public static void main(String[] args) {
        // Concaténation de chaînes : "Hello World !"
        String sentence = Stream.of("Hello", "World", "!").reduce("", (acc, str) -> acc + " " + str).trim();
        System.out.println("Phrase : '" + sentence + "'");

        // Recherche du minimum : min = -5.0
        double minTemperature = Stream.of(10.5, 1.2, -5.0, 3.8)
            .reduce(Double.MAX_VALUE, Double::min);
        System.out.println("Température minimale : " + minTemperature);

        // Somme des entiers avec valeur initiale : total = 15
        int totalSum = Stream.of(1, 2, 3, 4, 5).reduce(0, Integer::sum);
        System.out.println("Somme totale (avec identité) : " + totalSum);

        // Somme des entiers sans valeur initiale : total = Optional[15]
        Optional<Integer> optionalSum = Stream.of(1, 2, 3, 4, 5).reduce(Integer::sum);
        optionalSum.ifPresent(sum -> System.out.println("Somme totale (sans identité) : " + sum));

        // Filtrage et concaténation : "aei" (voyelles minuscules)
        String vowels = Stream.of("A", "b", "E", "f", "I", "j")
            .filter(s -> s.matches("[aeiouAEIOU]")) // Filtrer les voyelles
            .map(String::toLowerCase)               // Convertir en minuscule
            .reduce("", String::concat);
        System.out.println("Voyelles concaténées : " + vowels);
    }
}

Les méthodes reduce() avec une valeur initiale retournent directement l'objet agrégé, tandis que celles sans valeur initiale retournent un Optional.

limit et skip

limit(n) retourne un nouveau flux contenant les n premiers éléments du flux original. skip(n), quant à lui, ignore les n premiers éléments et retourne un flux avec les éléments restants.

Exemple 16 : Impact de limit et skip sur les opérations

import java.util.ArrayList;
import java.util.List;
import java.util.stream.Collectors;

class Product {
    public int id;
    private String name;

    public Product(int id, String name) {
        this.id = id;
        this.name = name;
    }

    public String getName() {
        System.out.println("Appel de getName() pour le produit : " + name); // Pour observer l'exécution
        return name;
    }

    @Override
    public String toString() {
        return "Product{id=" + id + ", name='" + name + "'}";
    }
}

public class LimitSkipExample {
    public static void main(String[] args) {
        List<Product> productCatalog = new ArrayList<>();
        for (int i = 1; i <= 1000; i++) {
            productCatalog.add(new Product(i, "Produit_" + i));
        }

        List<String> selectedProductNames = productCatalog.stream()
            .map(Product::getName) // Opération qui sera tracée
            .limit(10)             // Limite le flux aux 10 premiers éléments
            .skip(3)               // Saute les 3 premiers des 10 éléments restants
            .collect(Collectors.toList());

        System.out.println("\nProduits sélectionnés : " + selectedProductNames);
    }
}

L'exécution de cet exemple montrera que la méthode getName() n'est appelée que 10 fois (grâce à limit(10)), et le résultat final ne contient que 7 noms de produits après l'opération skip(3). Les opérations à court-circuit optimisent le traitement en ne calculant que ce qui est strictement nécessaire.

Cependant, cette optimisation peut être compromise si limit ou skip sont placés après une opération de tri (sorted), car l'opération de tri doit souvent évaluer tous les éléments pour déterminer leur ordre.

Exemple 17 : limit et skip sans effet sur le nombre d'appels après sorted

import java.util.ArrayList;
import java.util.Comparator;
import java.util.List;
import java.util.stream.Collectors;

// Product class as defined above
public class LimitSkipSortedExample {
    public static void main(String[] args) {
        List<Product> productCatalog = new ArrayList<>();
        for (int i = 1; i <= 5; i++) {
            productCatalog.add(new Product(i, "Article_" + i));
        }

        List<String> sortedLimitedNames = productCatalog.stream()
            .sorted(Comparator.comparing(Product::getName)) // Tri par nom (exige de voir tous les noms)
            .map(Product::getName)                          // Opération qui sera tracée
            .limit(2)                                       // Limite les résultats triés aux 2 premiers
            .collect(Collectors.toList());

        System.out.println("\nNoms triés et limités : " + sortedLimitedNames);
    }
}

Ici, getName() sera appelé 5 fois, car l'opération sorted doit potentiellement accéder à tous les éléments pour déterminer l'ordre correct avant que limit ne puisse être appliqué. Pour les flux parallèles, limit peut être coûteux si l'ordre des éléments doit être préservé. Dans de tels cas, il peut être préférable d'annuler l'ordre ou d'éviter les flux parallèles.

sorted

L'opération sorted() permet de trier les éléments d'un flux. Son avantage est la possibilité de combiner des opérations comme map, filter, limit ou skip avant le tri, ce qui peut considérablement réduire le nombre d'éléments à trier et améliorer les performances (par exemple, si vous ne voulez que les 10 premiers éléments triés, filtrez d'abord, puis triez les éléments restants).

Exemple 18 : Optimisation : tri après filtration et limitation

import java.util.ArrayList;
import java.util.Comparator;
import java.util.List;
import java.util.stream.Collectors;

// Product class as defined above, with a getPrice() method
class ProductWithPrice {
    public int id;
    private String name;
    private double price;

    public ProductWithPrice(int id, String name, double price) {
        this.id = id;
        this.name = name;
        this.price = price;
    }

    public String getName() {
        System.out.println("Fetching name for: " + name);
        return name;
    }
    public double getPrice() { return price; }

    @Override
    public String toString() {
        return "ProductWithPrice{id=" + id + ", name='" + name + "', price=" + price + "}";
    }
}

public class OptimizedSortExample {
    public static void main(String[] args) {
        List<ProductWithPrice> catalog = new ArrayList<>();
        catalog.add(new ProductWithPrice(1, "Clavier", 75.0));
        catalog.add(new ProductWithPrice(2, "Souris", 25.0));
        catalog.add(new ProductWithPrice(3, "Écran", 300.0));
        catalog.add(new ProductWithPrice(4, "Webcam", 50.0));
        catalog.add(new ProductWithPrice(5, "Microphone", 100.0));

        System.out.println("Noms extraits pour tri :");
        // Tri des 3 produits les moins chers par ordre alphabétique de nom
        List<String> cheapProductsSorted = catalog.stream()
            .sorted(Comparator.comparingDouble(ProductWithPrice::getPrice)) // Tri initial par prix
            .limit(3)                                                       // Ne garde que les 3 moins chers
            .map(ProductWithPrice::getName)                                 // Extrait les noms (appelé 3 fois seulement)
            .sorted()                                                       // Trie les noms restants
            .collect(Collectors.toList());

        System.out.println("\n3 produits les moins chers (noms triés) : " + cheapProductsSorted);
    }
}

L'optimisation n'est pertinente que si la logique métier ne requiert pas un tri complet avant de prendre une sous-sélection.

min, max et distinct

Les opérations min() et max() peuvent être réalisées en triant le flux puis en utilisant findFirst(), mais les méthodes dédiées min() et max() sont plus performantes (complexité O(n) au lieu de O(n log n) pour le tri). Elles sont fréquemment utilisées et donc exposées séparément.

Exemple 19 : Recherche de la chaîne la plus longue dans une liste

import java.util.Arrays;
import java.util.List;
import java.util.OptionalInt;

public class MaxLengthExample {
    public static void main(String[] args) {
        List<String> lines = Arrays.asList(
            "Ceci est une ligne courte.",
            "Voici une ligne beaucoup plus longue, avec plus de mots.",
            "Courte."
        );
        
        OptionalInt longestLength = lines.stream()
            .mapToInt(String::length) // Convertit chaque chaîne en sa longueur
            .max();                   // Trouve la longueur maximale
        
        longestLength.ifPresent(len -> System.out.println("Longueur de la ligne la plus longue : " + len));
    }
}

L'opération distinct() élimine les doublons du flux, ne laissant que les éléments uniques.

Exemple 20 : Extraction de mots uniques triés d'une liste de phrases

import java.util.Arrays;
import java.util.List;
import java.util.stream.Collectors;
import java.util.stream.Stream;

public class UniqueWordsExample {
    public static void main(String[] args) {
        List<String> sentences = Arrays.asList(
            "Le chat est sur le toit.",
            "Le chien est dans le jardin.",
            "Le chat et le chien jouent."
        );

        List<String> uniqueSortedWords = sentences.stream()
            .flatMap(line -> Arrays.stream(line.split("\\s+"))) // Divise chaque phrase en mots
            .map(word -> word.replaceAll("[^a-zA-Z]", "").toLowerCase()) // Nettoie et met en minuscules
            .filter(word -> !word.isEmpty())                          // Élimine les mots vides
            .distinct()                                               // Garde seulement les mots uniques
            .sorted()                                                 // Trie alphabétiquement
            .collect(Collectors.toList());

        System.out.println("Mots uniques et triés : " + uniqueSortedWords);
    }
}

Match (allMatch, anyMatch, noneMatch)

Les méthodes match vérifient si les éléments d'un flux satisfont un prédicat donné :

  • allMatch(Predicate p) : Retourne true si tous les éléments satisfont le prédicat.
  • anyMatch(Predicate p) : Retourne true si au moins un élément satisfait le prédicat.
  • noneMatch(Predicate p) : Retourne true si aucun élément ne satisfait le prédicat.

Ces opérations sont à court-circuit : elles peuvent retourner un résultat sans parcourir l'intégralité du flux. Par exemple, allMatch retourne false dès qu'un élément ne satisfait pas la condition.

Exemple 21 : Vérification des propriétés d'une collection de produits

import java.util.ArrayList;
import java.util.List;

class ProductItem {
    private int id;
    private String label;
    private double price;

    public ProductItem(int id, String label, double price) {
        this.id = id;
        this.label = label;
        this.price = price;
    }

    public double getPrice() { return price; }
}

public class MatchExample {
    public static void main(String[] args) {
        List<ProductItem> items = new ArrayList<>();
        items.add(new ProductItem(1, "Stylo", 2.50));
        items.add(new ProductItem(2, "Cahier", 5.00));
        items.add(new ProductItem(3, "Gomme", 1.20));
        items.add(new ProductItem(4, "Règle", 3.75));

        boolean areAllCheap = items.stream()
            .allMatch(item -> item.getPrice() < 10.0); // Tous les articles sont-ils à moins de 10€ ?
        System.out.println("Tous les articles sont-ils bon marché (<10€) ? " + areAllCheap); // true

        boolean hasExpensiveItem = items.stream()
            .anyMatch(item -> item.getPrice() > 5.0); // Y a-t-il un article cher (>5€) ?
        System.out.println("Y a-t-il un article cher (>5€) ? " + hasExpensiveItem); // false

        boolean noFreeItems = items.stream()
            .noneMatch(item -> item.getPrice() == 0.0); // Aucun article n'est-il gratuit ?
        System.out.println("Aucun article n'est-il gratuit ? " + noFreeItems); // true
    }
}

Génération de Flux Avancée

Stream.generate

La méthode Stream.generate() permet de créer un flux infini à partir d'une interface Supplier. Elle est utile pour générer des nombres aléatoires, des constantes, ou des séquences où chaque élément dépend de l'état précédent. Étant infini, il est crucial d'utiliser des opérations à court-circuit comme limit() pour restreindre sa taille.

Exemple 22 : Génération de 10 entiers aléatoires

import java.util.Random;
import java.util.function.Supplier;
import java.util.stream.IntStream;
import java.util.stream.Stream;

public class StreamGenerateExample {
    public static void main(String[] args) {
        System.out.println("10 nombres aléatoires (méthode 1):");
        Random rnd = new Random();
        Supplier<Integer> randomNumberSupplier = rnd::nextInt;
        Stream.generate(randomNumberSupplier)
              .limit(10)
              .forEach(System.out::println);

        System.out.println("\n10 nombres aléatoires (méthode 2):");
        IntStream.generate(() -> (int) (System.nanoTime() % 100)) // Génère un int basé sur le temps nano
                 .limit(10)
                 .forEach(System.out::println);
    }
}

Stream.generate() accepte également des implémentations personnalisées de Supplier pour des scénarios plus complexes, comme la création de données de test en masse ou le calcul d'éléments en fonction de règles spécifiques.

Exemple 23 : Génération personnalisée de produits

import java.util.Random;
import java.util.function.Supplier;
import java.util.stream.Stream;

class CustomProduct {
    private static int nextId = 1;
    private String description;
    private double unitPrice;

    public CustomProduct(String description, double unitPrice) {
        this.description = description;
        this.unitPrice = unitPrice;
    }

    public int getId() { return nextId++; }
    public String getDescription() { return description; }
    public double getUnitPrice() { return unitPrice; }

    @Override
    public String toString() {
        return "CustomProduct{" +
               "id=" + (nextId-1) + // nextId already incremented
               ", description='" + description + '\'' +
               ", unitPrice=" + String.format("%.2f", unitPrice) +
               '}';
    }
}

class ProductCreator implements Supplier<CustomProduct> {
    private int counter = 0;
    private Random priceRandomizer = new Random();

    @Override
    public CustomProduct get() {
        counter++;
        double price = 10.0 + (priceRandomizer.nextDouble() * 90.0); // Prix entre 10.0 et 100.0
        return new CustomProduct("Article_Test_" + counter, price);
    }
}

public class CustomStreamGeneration {
    public static void main(String[] args) {
        Stream.generate(new ProductCreator())
              .limit(5)
              .forEach(System.out::println);
    }
}

Stream.iterate

L'opération iterate() est similaire à reduce, prenant une valeur "germe" (seed) et un UnaryOperator (par exemple, f). Le germe devient le premier élément du flux, f(germe) le deuxième, f(f(germe)) le troisième, et ainsi de suite. Comme generate(), elle crée un flux infini et nécessite limit().

Exemple 24 : Génération d'une suite arithmétique

import java.util.stream.Stream;

public class StreamIterateExample {
    public static void main(String[] args) {
        System.out.println("Suite arithmétique (départ 5, pas de 7, 8 éléments) :");
        Stream.iterate(5, n -> n + 7) // Commence à 5, ajoute 7 à chaque fois
              .limit(8)
              .forEach(x -> System.out.print(x + " ")); // Affiche : 5 12 19 26 33 40 47 54
        System.out.println();
    }
}

Opérations de Réduction avec Collectors

La classe java.util.stream.Collectors est essentielle pour les opérations de réduction, permettant de transformer les éléments d'un flux en diverses collections ou de les grouper.

groupingBy et partitioningBy

groupingBy regroupe les éléments du flux en fonction d'une clé produite par une fonction de classification. Le résultat est une Map où la clé est le résultat de la fonction et la valeur est une List des éléments correspondants.

Exemple 25 : Regroupement de produits par tranche de prix

import java.util.List;
import java.util.Map;
import java.util.Random;
import java.util.function.Supplier;
import java.util.stream.Collectors;
import java.util.stream.Stream;

// Assumer la classe CustomProduct définie précédemment
class ProductGenerator implements Supplier<CustomProduct> {
    private static int counter = 0;
    private Random random = new Random();

    @Override
    public CustomProduct get() {
        counter++;
        double price = 10.0 + (random.nextDouble() * 190.0); // Prix entre 10.0 et 200.0
        return new CustomProduct("Gadget-" + counter, price);
    }
}

public class GroupingByExample {
    public static void main(String[] args) {
        Map<String, List<CustomProduct>> productsByPriceRange = Stream.generate(new ProductGenerator())
            .limit(50) // Génère 50 produits
            .collect(Collectors.groupingBy(product -> {
                double price = product.getUnitPrice();
                if (price < 50) return "Économique";
                else if (price < 100) return "Standard";
                else if (price < 150) return "Premium";
                else return "Luxe";
            }));

        productsByPriceRange.forEach((range, products) ->
            System.out.println("Catégorie " + range + ": " + products.size() + " produits")
        );
    }
}

partitioningBy est un cas particulier de groupingBy qui divise le flux en deux groupes (true et false) en fonction d'un prédicat. Le résultat est une Map<Boolean, List<T>>.

Exemple 26 : Partitionnement des produits par disponibilité

import java.util.List;
import java.util.Map;
import java.util.Random;
import java.util.function.Supplier;
import java.util.stream.Collectors;
import java.util.stream.Stream;

class InventoryItem {
    private String name;
    private boolean inStock;

    public InventoryItem(String name, boolean inStock) {
        this.name = name;
        this.inStock = inStock;
    }
    public boolean isInStock() { return inStock; }
}

class ItemStockGenerator implements Supplier<InventoryItem> {
    private int counter = 0;
    private Random random = new Random();

    @Override
    public InventoryItem get() {
        counter++;
        return new InventoryItem("Article #" + counter, random.nextBoolean()); // Stock aléatoire
    }
}

public class PartitioningByExample {
    public static void main(String[] args) {
        Map<Boolean, List<InventoryItem>> stockStatus = Stream.generate(new ItemStockGenerator())
            .limit(30)
            .collect(Collectors.partitioningBy(InventoryItem::isInStock));

        System.out.println("Articles en stock : " + stockStatus.get(true).size());
        System.out.println("Articles en rupture de stock : " + stockStatus.get(false).size());
    }
}

Résumé des Caractéristiques des Flux

Les principales caractéristiques des flux sont :

  • Non-structure de données : Un flux ne stocke pas de données en interne. Il opère sur des données provenant d'une source (collection, tableau, fonction génératrice, canal I/O) via un pipeline d'opérations.
  • Immuabilité de la source : Les opérations sur un flux ne modifient jamais les données de la structure sous-jacente. Par exemple, filter produit un nouveau flux sans les éléments filtrés, au lieu de les supprimer de la source.
  • Expressions Lambda : Toutes les opérations de flux prennent des expressions Lambda comme paramètres.
  • Pas d'accès par index : Il n'est pas possible d'accéder directement au deuxième, troisième ou dernier élément par index, bien que des méthodes comme findFirst() existent.
  • Conversion facile : Les flux peuvent être facilement convertis en tableaux ou en listes.
  • Évaluation paresseuse : De nombreuses opérations de flux sont différées jusqu'à ce que l'opération terminale soit appelée, ce qui permet d'optimiser les calculs. Les opérations intermédiaires sont toujours paresseuses.
  • Capacités parallèles : Les flux peuvent être facilement parallélisés, ce qui permet de tirer parti des architectures multi-cœurs sans écrire de code multi-threadé explicite.
  • Flux infinis possibles : Contrairement aux collections, les flux n'ont pas de taille fixe. Des opérations à court-circuit comme limit(n) ou findFirst() permettent de travailler efficacement avec des flux infinis.

Flux Séquentiels vs. Parallèles

Un flux séquentiel traite les éléments avec un seul thread, tandis qu'un flux parallèle utilise plusieurs threads pour le traitement concurrent. Le mode parallèle est généralement plus rapide pour les grands volumes de données.

Exemple : Comparaison de performance pour le tri (Séquentiel vs. Parallèle)

import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Random;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;

public class ParallelStreamPerformance {
    public static void main(String[] args) {
        List<Long> dataSet = new ArrayList<>();
        Random rand = new Random();
        int dataSize = 10_000_000; // Dix millions d'éléments

        for (int i = 0; i < dataSize; i++) {
            dataSet.add(rand.nextLong());
        }

        // Test du flux séquentiel
        long startTimeSequential = System.nanoTime();
        List<Long> sortedSequential = dataSet.stream()
            .sequential()
            .sorted()
            .collect(Collectors.toList());
        long endTimeSequential = System.nanoTime();
        long sequentialTimeMs = TimeUnit.NANOSECONDS.toMillis(endTimeSequential - startTimeSequential);

        // Test du flux parallèle
        long startTimeParallel = System.nanoTime();
        List<Long> sortedParallel = dataSet.stream()
            .parallel()
            .sorted()
            .collect(Collectors.toList());
        long endTimeParallel = System.nanoTime();
        long parallelTimeMs = TimeUnit.NANOSECONDS.toMillis(endTimeParallel - startTimeParallel);

        System.out.println("Temps de tri séquentiel : " + sequentialTimeMs + " ms");
        System.out.println("Temps de tri parallèle : " + parallelTimeMs + " ms");
        
        // Pour éviter les avertissements d'utilisation de sortedSequential et sortedParallel
        System.out.println("Premier élément séquentiel : " + sortedSequential.get(0));
        System.out.println("Premier élément parallèle : " + sortedParallel.get(0));
    }
}

Résultat typique sur une machine multi-cœurs :

Temps de tri séquentiel : 10362 ms
Temps de tri parallèle : 6527 ms

Questions fréquentes sur les flux parallèles :

  • **Combien de threads sont utilisés pour un flux parallèle ?**Le nombre de threads est généralement égal au nombre de cœurs CPU disponibles sur la machine, grâce au principe d'optimalité parallèle du framework Fork/Join.

  • **Les flux parallèles utilisent-ils un pool de threads ?**Oui, ils utilisent un ForkJoinPool par défaut, configuré avec un nombre de threads correspondant au nombre de processeurs disponibles. C'est le pool de threads commun à l'ensemble du système pour les tâches Fork/Join.

  • **Peut-on personnaliser ce pool de threads ?**Oui, pour des cas d'usage spécifiques où le pool par défaut ne serait pas optimal (par exemple, si les tâches sont I/O-bound), il est possible de créer un ForkJoinPool personnalisé et d'exécuter des opérations de flux à l'intérieur de celui-ci :

    import java.util.concurrent.ForkJoinPool;
    // ... (code précédent)
    ForkJoinPool customPool = new ForkJoinPool(4); // Utilise 4 threads
    try {
       List<Long> customParallelSorted = customPool.submit(() -> 
           dataSet.stream()
               .parallel()
               .sorted()
               .collect(Collectors.toList())
       ).get();
       // ...
    } catch (Exception e) {
       e.printStackTrace();
    } finally {
       customPool.shutdown();
    }
    
    

    Cependant, il est généralement recommandé d'utiliser le pool par défaut à moins d'avoir des raisons spécifiques de le modifier, car le ForkJoinPool.commonPool() est partagé et optimisé pour la plupart des charges de travail.

Étiquettes: Java StreamAPI LambdaExpressions FunctionalProgramming Java8

Publié le 24 septembre à 02h21