1. Własne Collector’y: kiedy i jak pisać własne
W Java Stream API do przekształcania strumienia w kolekcję lub agregat używa się interfejsu Collector. Zazwyczaj korzystasz z gotowych collectorów z klasy Collectors (toList(), toMap(), groupingBy() i inne), ale czasem potrzebujesz czegoś specjalnego — wtedy możesz napisać własny Collector.
Collector — to obiekt, który opisuje, jak z elementów strumienia złożyć wynik końcowy. Określa cztery (a właściwie pięć) kluczowe komponenty:
- supplier — tworzy nowy kontener do zbierania elementów (np. nową listę lub mapę).
- accumulator — dodaje kolejny element do kontenera.
- combiner — łączy dwa kontenery (ważne dla strumieni równoległych!).
- finisher — przekształca kontener w wynik końcowy (np. czyni go niemodyfikowalnym lub konwertuje na inny typ).
- characteristics — zestaw flag opisujących właściwości collectora (np. czy wspiera równoległość, czy zmienia typ wyniku itd.).
Sygnatura:
Collector<T, A, R>
- T — typ elementów strumienia,
- A — typ pośredniego akumulatora,
- R — typ wyniku.
2. Przykład: Collector dla MultiMap (Map<K, List<V>>)
Załóżmy, że chcesz zebrać strumień par Pair<K, V> do Map<K, List<V>> (multi-mapy), w której każdemu kluczowi odpowiada lista wartości.
Przykład implementacji:
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
);
}
Użycie:
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. Przykład: Collector dla top-N elementów
Załóżmy, że chcesz zebrać strumień do listy z N największych elementów (np. top-5 w kolejności malejącej).
Implementacja:
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(); // usuwamy najmniejszy
},
(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()); // malejąco
return result;
},
Collector.Characteristics.UNORDERED
);
}
Użycie:
List<Integer> top3 = Stream.of(5, 1, 9, 3, 7, 2).collect(topN(3, Comparator.naturalOrder()));
// top3: [9, 7, 5]
4. Kiedy NIE warto pisać własnego Collectora
- Jeśli da się wyrazić zadanie przez kombinację standardowych collectorów i operacji downstream (groupingBy, mapping, flatMapping, collectingAndThen itd.), lepiej użyć ich.
- Własny Collector jest potrzebny tylko dla naprawdę niestandardowych scenariuszy (specjalna struktura danych, złożona agregacja, top-N, multi-mapy itd.).
- Nie pisz Collectora dla samego Collectora — to utrudnia utrzymanie i testowanie.
Przykład:
// Zamiast własnego Collectora dla Map<K, Set<V>>:
.collect(Collectors.groupingBy(
Pair::key,
Collectors.mapping(Pair::value, Collectors.toSet())
))
5. Własny Spliterator: po co i jak
Spliterator — to specjalny interfejs do efektywnego iterowania i dzielenia kolekcji (lub innych źródeł danych) na części, zwłaszcza dla przetwarzania równoległego. W odróżnieniu od zwykłego iteratora, Spliterator może „dzielić” (split) kolekcję na niezależne kawałki do równoległego przetwarzania.
Kluczowe metody:
- tryAdvance(Consumer<? super T> action) — przetworzyć następny element.
- trySplit() — spróbować podzielić kolekcję na dwie części (zwraca nowy Spliterator dla jednej z części).
- estimateSize() — oszacowanie pozostałej liczby elementów.
- characteristics() — maska bitowa cech (ORDERED, SIZED, SUBSIZED i inne).
trySplit: strategie dzielenia
Zrównoważony podział — istotny dla strumieni równoległych: trySplit powinien zwracać części o zbliżonym rozmiarze, aby obciążenie było równomierne.
Jeśli nie ma czego dzielić (np. mało elementów) — zwracamy null.
Przykład: Spliterator do porcjowego czytania pliku
Załóżmy, że masz duży plik i chcesz przetwarzać go po 1000 wierszy naraz (partiami), aby nie trzymać wszystkiego w pamięci.
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() {
// Dla strumieniowego czytania z pliku podział nie ma sensu — zwracamy null
return null;
}
@Override
public long estimateSize() {
return Long.MAX_VALUE; // nie wiadomo z góry
}
@Override
public int characteristics() {
return ORDERED | NONNULL;
}
}
Użycie:
try (BufferedReader reader = Files.newBufferedReader(Path.of("big.txt"))) {
StreamSupport.stream(new ChunkedLineSpliterator(reader, 1000), false)
.forEach(chunk -> processChunk(chunk));
}
Charakterystyki Spliteratora
- ORDERED — elementy występują w określonym porządku (np. lista).
- SIZED — znana jest dokładna liczba elementów.
- SUBSIZED — wszystkie Spliteratory otrzymane przez trySplit też są SIZED.
- IMMUTABLE — źródło nie zmienia się podczas iteracji.
- CONCURRENT — źródło wspiera bezpieczną modyfikację równoległą.
- DISTINCT, SORTED, NONNULL — dodatkowe właściwości.
Ważne: poprawne wskazanie charakterystyk wpływa na optymalizację strumieni.
6. Przykłady
- Czytanie pliku porcjami (chunkami) — pozwala przetwarzać duże pliki partiami, bez ładowania wszystkiego do pamięci.
- Parsowanie bez zbędnych alokacji — jeśli parsujesz strumień bajtów/znaków i chcesz zminimalizować tworzenie obiektów tymczasowych, możesz zaimplementować Spliterator, który zwraca „okna” lub „wycinki” źródłowej tablicy.
Przykład: Spliterator do parsowania CSV po wierszach
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; // parsowanie sekwencyjne
}
@Override
public long estimateSize() {
return Long.MAX_VALUE;
}
@Override
public int characteristics() {
return ORDERED | NONNULL;
}
}
7. Integracja z parallel() — jak zrobić to bezpiecznie
- Jeśli Twój Spliterator wspiera równoległy podział (trySplit nie zwraca null), a charakterystyki zawierają SIZED/SUBSIZED, Stream API będzie w stanie efektywnie zrównoleglić przetwarzanie.
- Dla źródeł strumieniowych (pliki, gniazda) zazwyczaj podział nie jest wspierany — używaj strumieni sekwencyjnych.
- Dla kolekcji i tablic — zaimplementuj zrównoważony podział (np. dziel tablicę na połowy).
Przykład: Spliterator dla tablicy
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;
}
}
Użycie:
String[] arr = {"a", "b", "c", "d"};
StreamSupport.stream(new ArraySpliterator<>(arr, 0, arr.length), true)
.forEach(System.out::println);
GO TO FULL VERSION