1. Collectors personnalisés : quand et comment écrire les vôtres
Dans le Java Stream API, l’interface Collector est utilisée pour transformer un flux en collection ou en agrégat. En général, vous utilisez les collectors prêts à l’emploi de la classe Collectors (toList(), toMap(), groupingBy(), etc.), mais parfois vous avez besoin de quelque chose de spécifique — vous pouvez alors écrire votre propre Collector.
Collector — un objet qui décrit comment assembler le résultat final à partir des éléments d’un flux. Il définit quatre (en réalité cinq) composants clés :
- supplier — crée un nouveau conteneur pour collecter les éléments (par exemple, une nouvelle liste ou une map).
- accumulator — ajoute l’élément suivant au conteneur.
- combiner — fusionne deux conteneurs (important pour les streams parallèles!).
- finisher — transforme le conteneur en résultat final (par exemple, le rend immuable ou le convertit en un autre type).
- characteristics — un ensemble de flags décrivant les propriétés du collector (par exemple, prend-il en charge le parallélisme, change-t-il le type de résultat, etc.).
Signature :
Collector<T, A, R>
- T — type des éléments du flux,
- A — type de l’accumulateur intermédiaire,
- R — type du résultat.
2. Exemple : Collector pour MultiMap (Map<K, List<V>>)
Supposons que vous vouliez collecter un flux de paires Pair<K, V> dans une Map<K, List<V>> (multi-map), où chaque clé correspond à une liste de valeurs.
Exemple d’implémentation :
public static <K, V> Collector<Pair<K, V>, ?, Map<K, List<V>>> toMultiMap() {
return Collector.of(
HashMap::new, // supplier
(map, pair) -> map.computeIfAbsent(pair.key(), k -> new ArrayList<>()).add(pair.value()), // accumulator
(map1, map2) -> { // combiner
map2.forEach((k, vList) -> map1.merge(k, vList, (l1, l2) -> { l1.addAll(l2); return l1; }));
return map1;
},
Function.identity(), // finisher
Collector.Characteristics.UNORDERED
);
}
Utilisation :
List<Pair<String, Integer>> pairs = List.of(
new Pair<>("a", 1), new Pair<>("b", 2), new Pair<>("a", 3)
);
Map<String, List<Integer>> multiMap = pairs.stream().collect(toMultiMap());
// multiMap: {a=[1, 3], b=[2]}
3. Exemple : Collector pour le top‑N d’éléments
Supposons que vous vouliez collecter un flux en une liste des N plus grands éléments (par exemple, top‑5 en ordre décroissant).
Implémentation :
public static <T> Collector<T, ?, List<T>> topN(int n, Comparator<? super T> comparator) {
return Collector.of(
() -> new PriorityQueue<>(n, comparator), // supplier
(pq, t) -> {
pq.offer(t);
if (pq.size() > n) pq.poll(); // on supprime le plus petit
},
(pq1, pq2) -> {
pq2.forEach(t -> {
pq1.offer(t);
if (pq1.size() > n) pq1.poll();
});
return pq1;
},
pq -> {
List<T> result = new ArrayList<>(pq);
result.sort(comparator.reversed()); // en ordre décroissant
return result;
},
Collector.Characteristics.UNORDERED
);
}
Utilisation :
List<Integer> top3 = Stream.of(5, 1, 9, 3, 7, 2).collect(topN(3, Comparator.naturalOrder()));
// top3: [9, 7, 5]
4. Quand il NE faut PAS écrire son propre Collector
- Si la tâche peut être exprimée via une combinaison de collectors standard et d’opérations downstream (groupingBy, mapping, flatMapping, collectingAndThen, etc.), il vaut mieux les utiliser.
- Un Collector personnalisé n’est nécessaire que pour des scénarios réellement atypiques (structure de données particulière, agrégation complexe, top‑N, multi‑maps, etc.).
- N’écrivez pas un Collector pour le principe — cela complique la maintenance et les tests.
Exemple :
// Au lieu d'écrire votre propre Collector pour Map<K, Set<V>>:
.collect(Collectors.groupingBy(
Pair::key,
Collectors.mapping(Pair::value, Collectors.toSet())
))
5. Spliterator personnalisé : pourquoi et comment
Spliterator — une interface dédiée pour itérer et découper efficacement des collections (ou d’autres sources de données) en parties, en particulier pour le traitement parallèle. Contrairement à un itérateur classique, un Spliterator peut « scinder » (split) une collection en morceaux indépendants pour un traitement parallèle.
Méthodes clés :
- tryAdvance(Consumer<? super T> action) — traiter l’élément suivant.
- trySplit() — tenter de diviser la collection en deux parties (retourne un nouveau Spliterator pour l’une des parties).
- estimateSize() — estimation du nombre d’éléments restants.
- characteristics() — masque de bits des caractéristiques (ORDERED, SIZED, SUBSIZED, etc.).
trySplit : stratégies de découpage
Découpage équilibré — important pour les streams parallèles : trySplit doit retourner des parties d’une taille approximativement égale, afin que les threads soient chargés uniformément.
S’il n’y a rien à découper (par exemple, trop peu d’éléments) — retournez null.
Exemple : Spliterator pour la lecture de fichier par portions
Supposons que vous ayez un gros fichier et que vous souhaitiez le traiter par paquets de 1000 lignes, afin de ne pas tout garder en mémoire.
public class ChunkedLineSpliterator implements Spliterator<List<String>> {
private final BufferedReader reader;
private final int chunkSize;
public ChunkedLineSpliterator(BufferedReader reader, int chunkSize) {
this.reader = reader;
this.chunkSize = chunkSize;
}
@Override
public boolean tryAdvance(Consumer<? super List<String>> action) {
List<String> chunk = new ArrayList<>(chunkSize);
try {
String line;
for (int i = 0; i < chunkSize && (line = reader.readLine()) != null; i++) {
chunk.add(line);
}
if (chunk.isEmpty()) return false;
action.accept(chunk);
return true;
} catch (IOException e) {
throw new UncheckedIOException(e);
}
}
@Override
public Spliterator<List<String>> trySplit() {
// Pour une lecture en flux depuis un fichier, le découpage n'a pas de sens - retourner null
return null;
}
@Override
public long estimateSize() {
return Long.MAX_VALUE; // inconnu à l'avance
}
@Override
public int characteristics() {
return ORDERED | NONNULL;
}
}
Utilisation :
try (BufferedReader reader = Files.newBufferedReader(Path.of("big.txt"))) {
StreamSupport.stream(new ChunkedLineSpliterator(reader, 1000), false)
.forEach(chunk -> processChunk(chunk));
}
Caractéristiques de Spliterator
- ORDERED — les éléments suivent un ordre défini (par exemple, une liste).
- SIZED — le nombre exact d’éléments est connu.
- SUBSIZED — tous les Spliterator obtenus via trySplit sont aussi SIZED.
- IMMUTABLE — la source ne change pas pendant l’itération.
- CONCURRENT — la source prend en charge une modification parallèle sûre.
- DISTINCT, SORTED, NONNULL — propriétés additionnelles.
Important : indiquer correctement les caractéristiques — cela influe sur l’optimisation des streams.
6. Exemples
- Lecture de fichier par blocs (chunks) — permet de traiter de gros fichiers par morceaux sans tout charger en mémoire.
- Parsing sans allocations inutiles — si vous analysez un flux d’octets/de caractères et souhaitez minimiser la création d’objets temporaires, vous pouvez implémenter un Spliterator qui renvoie des « fenêtres » ou des « tranches » du tableau source.
Exemple : Spliterator pour l’analyse d’un CSV ligne par ligne
public class CsvLineSpliterator implements Spliterator<String[]> {
private final BufferedReader reader;
public CsvLineSpliterator(BufferedReader reader) {
this.reader = reader;
}
@Override
public boolean tryAdvance(Consumer<? super String[]> action) {
try {
String line = reader.readLine();
if (line == null) return false;
action.accept(line.split(","));
return true;
} catch (IOException e) {
throw new UncheckedIOException(e);
}
}
@Override
public Spliterator<String[]> trySplit() {
return null; // parsing séquentiel
}
@Override
public long estimateSize() {
return Long.MAX_VALUE;
}
@Override
public int characteristics() {
return ORDERED | NONNULL;
}
}
7. Intégration avec parallel() — comment le faire en toute sécurité
- Si votre Spliterator prend en charge le découpage parallèle (trySplit ne retourne pas null) et que ses caractéristiques incluent SIZED/SUBSIZED, alors le Stream API pourra paralléliser le traitement efficacement.
- Pour les sources en flux (fichiers, sockets), le découpage n’est généralement pas pris en charge — utilisez des streams séquentiels.
- Pour les collections et les tableaux — implémentez un découpage équilibré (par exemple, coupez le tableau en deux).
Exemple : Spliterator pour un tableau
public class ArraySpliterator<T> implements Spliterator<T> {
private final T[] array;
private int start, end;
public ArraySpliterator(T[] array, int start, int end) {
this.array = array;
this.start = start;
this.end = end;
}
@Override
public boolean tryAdvance(Consumer<? super T> action) {
if (start < end) {
action.accept(array[start++]);
return true;
}
return false;
}
@Override
public Spliterator<T> trySplit() {
int mid = (start + end) >>> 1;
if (mid == start) return null;
ArraySpliterator<T> split = new ArraySpliterator<>(array, start, mid);
start = mid;
return split;
}
@Override
public long estimateSize() {
return end - start;
}
@Override
public int characteristics() {
return ORDERED | SIZED | SUBSIZED | IMMUTABLE;
}
}
Utilisation :
String[] arr = {"a", "b", "c", "d"};
StreamSupport.stream(new ArraySpliterator<>(arr, 0, arr.length), true)
.forEach(System.out::println);
GO TO FULL VERSION