Skip to content

5. Typ Stream<T> w Javie 8

5.1. Przykład-01 – klasa Stream

Operacje na strumieniach Observable mają wiele wspólnych cech ze strumieniami Stream. Różnica polega na tym, że element strumienia Stream nie może zostać przetworzony, dopóki nie zostanie pobrany cały strumień Stream, podczas gdy element strumienia Observable może zostać przetworzony (obserwowany) natychmiast po jego otrzymaniu, bez konieczności oczekiwania na uzyskanie całego strumienia Observable. Kolejną różnicą jest to, że po otrzymaniu strumienia Stream jego wartości są wykorzystywane poprzez pobieranie (pull) ich pojedynczo ze strumienia Stream. W przypadku obserwowalnego obiektu jest inaczej. Gdy tylko wyemituje on wartość, jest ona wysyłana (pushed) do jego subskrybenta.

Kilka klas implementuje pojęcie Stream. Przedstawiamy tutaj klasę Stream<T>:

Image

Klasa Stream zawiera 39 metod. Przedstawimy tutaj kilka z nich. Rozważmy następujący kod:

  

package dvp.java8.streams;

import java.util.List;

import dvp.data.Personne;
import dvp.data.Personnes;

public class Exemple01 {
    public static void main(String[] args) {
        // lista osób
        List<Personne> personnes = Personnes.get();
        // wyświetlanie 1
        personnes.stream().forEach(p -> {
            System.out.println(p);
        });
        System.out.println("----------------");
        // wyświetlanie 2
        personnes.stream().forEach(System.out::println);
    }
}
  • wiersz 11: tworzymy instancję listy osób;
  • wiersz 13: na podstawie tej listy tworzy się obiekt typu Stream. W ten sposób wszystkie kolekcje można przekształcić w strumień Stream. Pozwala to na korzystanie ze wszystkich metod tej klasy, dzięki czemu można przetwarzać elementy kolekcji w sposób bardziej zwięzły niż przy użyciu pętli. Pozwala to również na równoległe przetwarzanie elementów, gdy jest to możliwe;
  • wiersz 13: metoda [Stream.forEach] ma następującą sygnaturę:
 

Widać, że parametrem tej metody jest interfejs funkcjonalny [Consumer<T>] przedstawiony w paragrafie 4.4 – interfejs, którego jedyna metoda wykorzystuje typ T i nie zwraca żadnej wartości.

  • W kodzie:

        personnes.stream().forEach(p -> {
            System.out.println(p);
});
  • [personnes.stream()] generuje strumień elementów typu [Personne], który zasila metodę [forEach]. Parametr p jest typu [Personne], a podana funkcja lambda wyświetla tę osobę;

Powyższy kod można uprościć w następujący sposób (wiersz 18):


personnes.stream().forEach(System.out::println);

Zamiast przekazywać jako parametr wartość funkcji lambda, przekazujemy odwołanie do istniejącej metody, w tym przypadku metody println klasy System.out. Oczywiście metoda ta musi mieć odpowiednią sygnaturę, w tym przypadku sygnaturę metody [Consumer.accept]: void accept(T t). Jak wspomniano wcześniej, parametrem metody [accept] będzie typ [Personne];

Otrzymujemy następujące wyniki:

1
2
3
4
5
6
7
{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"}
{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"}
{"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"}
----------------
{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"}
{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"}
{"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"}

Gdy strumień Stream został już wykorzystany, nie można z niego ponownie korzystać. Aby móc z niego ponownie skorzystać, należy go odtworzyć. Pokazuje to poniższy kod [Exemple01b]:


package dvp.java8.streams;

import java.util.stream.Stream;

import dvp.data.Personne;
import dvp.data.Personnes;

public class Exemple01b {
    public static void main(String[] args) {
        // przepływ osób
        Stream<Personne> personnes = Personnes.get().stream();
        // wyświetlanie 1
        personnes.forEach(p -> {
            System.out.println(p);
        });
        System.out.println("----------------");
        // wyświetlanie 2
        personnes.forEach(System.out::println);
    }
}
  • wiersz 11: w celu optymalizacji kodu postanowiono skonstruować Stream tylko raz. Uzyskane wyniki są następujące:

{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"}
{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"}
{"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"}
----------------
Exception in thread "main" java.lang.IllegalStateException: stream has already been operated upon or closed
    at java.util.stream.AbstractPipeline.sourceStageSpliterator(Unknown Source)
    at java.util.stream.ReferencePipeline$Head.forEach(Unknown Source)
at dvp.java8.streams.Exemple02b.main(Exemple02b.java:18)

Za każdym razem, gdy chcemy wykorzystać Stream, należy go utworzyć, nawet jeśli został już wcześniej utworzony.

5.2. Przykład 02 – równoległe przetwarzanie elementów strumienia

  

Rozważmy następujący kod:


package dvp.java8.streams;

import java.util.List;

import dvp.data.Personne;
import dvp.data.Personnes;

public class Exemple02 {
    public static void main(String[] args) {
        // lista osób
        List<Personne> personnes = Personnes.get();
        // wyświetlanie 1
        personnes.stream().forEach(Exemple02::affiche);
        System.out.println("-----------------");
        // wyświetlanie 2
        personnes.stream().parallel().forEach(Exemple02::affiche);
    }

    public static void affiche(Personne p) {
        System.out.printf("Personne %s sur thread %s%n", p, Thread.currentThread().getName());
    }
}
  • wiersze 19–21: metoda [affiche] wyświetla na konsoli ciąg znaków jSON dotyczący danej osoby oraz nazwę wątku wykonawczego, w którym odbywa się wyświetlanie;
  • wiersz 13: wyświetla listę osób. Należy zauważyć, że parametrem metody [forEach] jest odwołanie do poprzedniej metody statycznej;
  • wiersz 16: wykonujemy to samo, ale za pomocą metody [parallel] żądamy, aby przetwarzanie elementów strumienia odbywało się równolegle w wielu wątkach. Nie każde przetwarzanie może odbywać się równolegle. W tym przypadku należy założyć, że kolejność wyświetlania nie ma znaczenia, ponieważ w przetwarzaniu równoległym nie ma pewności co do kolejności wykonywania wątków. Należy również zwrócić uwagę na składnię, która będzie powszechnie stosowana zarówno w przypadku metod Stream, jak i Observable:
flux.m1(e1->...).m2(e2->..).m3(e3->...)...
  • (ciąg dalszy)
    • flux generuje elementy e1, które są przekazywane do metody m1;
    • flux.m1 jest z kolei strumieniem elementów e2, które zasilają metodę m2;
    • flux.m1.m2 to strumień elementów e3, które zasilają metodę m3;

Typ elementów e1, e2, e3 może ulegać zmianie w trakcie przetwarzania początkowego strumienia.

Wykonanie tego kodu daje następujący wynik:

1
2
3
4
5
6
7
Personne {"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"} sur thread main
Personne {"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"} sur thread main
Personne {"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"} sur thread main
-----------------
Personne {"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"} sur thread main
Personne {"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"} sur thread ForkJoinPool.commonPool-worker-1
Personne {"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"} sur thread ForkJoinPool.commonPool-worker-2

Widać, że przetwarzanie równoległe (wiersze 5–7) odbyło się w trzech różnych wątkach i nie zachowało kolejności elementów zgodnej z wierszami 1–3. W niniejszym dokumencie nie będziemy się zbytnio skupiać na równoległym przetwarzaniu elementów Stream, ponieważ wymagałoby to omówienia warunków umożliwiających takie przetwarzanie. Okazuje się zatem, że niewiele operacji można realizować równolegle. Jedną z tych, które w naturalny sposób nadają się do równoległości, jest suma elementów liczbowych strumienia, którą teraz przedstawimy.

5.3. Przykład-03 – równoległe przetwarzanie elementów strumienia

  

Rozważmy następujący kod (Przykład 03a):


package dvp.java8.streams;

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

public class Exemple03a {
    public static void main(String[] args) {

        final long limite = 10_000_000L;
        // liczba procesorów
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lista liczb
        long début = new Date().getTime();
        List<Long> nombres = new ArrayList<>();
        for (long i = 0; i < limite; i++) {
            nombres.add(i);
        }
        System.out.printf("création de la liste des %s nombres en %s ms%n", limite, new Date().getTime() - début);
        // suma liczb – metoda sekwencyjna
        début = new Date().getTime();
        long somme = nombres.stream().reduce(0L, (s, i) -> s + i);
        System.out.printf("somme séquentielle : somme=%s, durée (ms)=%s%n", somme, new Date().getTime() - début);
    }
}
  • W wierszu 22 wykorzystujemy metodę [reduce], której sygnatura jest następująca:
  • metoda [reduce] obsługuje elementy typu T;
  • metoda [reduce] stosuje to samo przetwarzanie do wszystkich elementów strumienia: początkowa wartość akumulatora jest podawana jako pierwszy parametr. Jako drugi parametr podawana jest metoda instancjonująca interfejs funkcjonalny [BinaryOperator] [2]: na podstawie każdego elementu i akumulatora metoda ta zwraca nową wartość akumulatora. Jego wartość końcowa jest wartością zwracaną przez metodę [reduce]. Kod [3] wyjaśnia ten mechanizm. Metoda [apply] jest metodą interfejsu funkcjonalnego [BinaryOperator] [2];

Wróćmy do przykładowego kodu:

  • wiersz 12: wyświetlana jest liczba rdzeni wykrytych przez metodę JVM;
  • wiersze 15–18: tworzy się listę zawierającą 10 milionów liczb;
  • wiersz 22: suma tych liczb jest obliczana sekwencyjnie przy użyciu jednego wątku;

Otrzymujemy następujące wyniki:

1
2
3
La JVM a détecté [8] coeurs sur votre machine
création de la liste des 10000000 nombres en 4336 ms
somme séquentielle : somme=49999995000000, durée (ms)=225

Teraz zastąpmy wiersz 22 kodu następującym (Przykład03b):


long somme = nombres.stream().parallel().reduce(0L, (s, i) -> s + i);

Zalecamy, aby elementy strumienia były przetwarzane równolegle przy użyciu wielu wątków. Jest to możliwe, ponieważ kolejność sumowania liczb nie ma znaczenia. Można zatem przypisać n1 liczb do wątku T1, n2 liczb do wątku T2, ... i na koniec zsumować wyniki dostarczone przez te różne wątki. Otrzymujemy wówczas następujące wyniki:

1
2
3
La JVM a détecté [8] coeurs sur votre machine
création de la liste des 10000000 nombres en 4332 ms
somme parallèle : somme=49999995000000, durée (ms)=184

Nie ma więc praktycznie żadnego wzrostu wydajności. W kolejnych przykładach często będzie tak właśnie wyglądać. Samo zarządzanie wątkami jest czasochłonne. Operacja wykonywana przez każdy rdzeń musi być wystarczająco złożona, aby dał się zauważyć wzrost wydajności. Pokazuje to poniższy przykład (Przykład03c):


package dvp.java8.streams;

import java.util.ArrayList;
import java.util.Date;
import java.util.List;
import java.util.function.BinaryOperator;

public class Exemple03c {
    public static void main(String[] args) {

        final long limite = 10_000L;
        // liczba procesorów
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lista liczb
        long début = new Date().getTime();
        List<Long> nombres = new ArrayList<>();
        for (long i = 0; i < limite; i++) {
            nombres.add(i);
        }
        System.out.printf("création de la liste des %s nombres en %s ms%n", limite, new Date().getTime() - début);
        // suma liczb – metoda sekwencyjna
        début = new Date().getTime();
        BinaryOperator<Long> bo = (s, i) -> {
            try {
                Thread.sleep(1);
            } catch (InterruptedException e) {
            }
            return s + i;
        };
        long somme = nombres.stream().reduce(0L, bo);
        System.out.printf("somme séquentielle : somme=%s, durée (ms)=%s%n", somme, new Date().getTime() - début);
    }
}
  • wiersz 30: ponownie wykorzystujemy metodę [reduce], której jako parametr przekazujemy odwołanie do metody z wierszy 23–29;
  • wiersz 28: metoda [bo] zwraca sumę swoich dwóch parametrów;
  • wiersze 24–27: sztucznie wstrzymujemy wątek na 1 milisekundę, aby zasymulować intensywną pracę;

Otrzymujemy wówczas następujące wyniki:

1
2
3
La JVM a détecté [8] coeurs sur votre machine
création de la liste des 10000 nombres en 2 ms
somme séquentielle : somme=49995000, durée (ms)=13617

Teraz, jeśli zastąpimy wiersz 30 następującym:


long somme = nombres.stream().parallel().reduce(0L, bo);

otrzymujemy następujące wyniki:

1
2
3
La JVM a détecté [8] coeurs sur votre machine
création de la liste des 10000 nombres en 2 ms
somme séquentielle : somme=49995000, durée (ms)=1598

Wyraźnie widać wzrost wydajności wynikający z równoległego wykonywania obliczeń sumy. W przypadku przetwarzania 8 liczb:

  • wątek sekwencyjny czeka 8 razy po 1 milisekundzie, czyli łącznie 8 ms;
  • 8 wątków równoległych czeka jednocześnie po 1 milisekundzie (dla uproszczenia), czyli łącznie 1 milisekundę dla 8 liczb;

Można zatem oczekiwać, że wykonanie równoległe będzie przebiegało 8 razy szybciej niż sekwencyjne. Tak właśnie jest w tym przypadku.

5.4. Przykład-04 – filtrowanie strumienia

  

Rozważmy następujący kod:


package dvp.java8.streams;

import java.util.List;

import dvp.data.Personne;
import dvp.data.Personnes;

public class Exemple04 {
    public static void main(String[] args) {
        // lista osób
        List<Personne> personnes = Personnes.get();
        // wyświetlenia
        System.out.println("age < 28 ----------------------");
        personnes.stream().filter(p -> p.getAge() < 28).forEach(p -> {
            System.out.println(p);
        });
        System.out.println("poids < 50 ----------------------");
        personnes.stream().filter(p -> p.getPoids() < 50).forEach(p -> {
            System.out.println(p);
        });
        System.out.println("age < 28 ----------------------");
        personnes.stream().filter(p -> p.getAge() < 28).forEach(System.out::println);
        System.out.println("poids < 50 ----------------------");
        personnes.stream().filter(p -> p.getPoids() < 50).forEach(System.out::println);
    }
}
  • wiersz 14: metoda [Stream.filter] ma następującą sygnaturę:
 
  • metoda [filter] oczekuje jako parametr instancję interfejsu funkcjonalnego [Predicate] przedstawionego w paragrafie 4.2, którego jedyną metodą do zaimplementowania jest następująca: boolean test(T t);
  • metoda [filter] zwraca elementy strumienia spełniające warunki Predicate. Służy ona zatem do filtrowania Stream;

Rozważmy następujący kod:


package dvp.java8.streams;

import java.util.List;

import dvp.data.Personne;
import dvp.data.Personnes;

public class Exemple04 {
    public static void main(String[] args) {
        // lista osób
        List<Personne> personnes = Personnes.get();
        // wyświetlenia
        System.out.println("age < 28 ----------------------");
        personnes.stream().filter(p -> p.getAge() < 28).forEach(p -> {
            System.out.println(p);
        });
        System.out.println("poids < 50 ----------------------");
        personnes.stream().filter(p -> p.getPoids() < 50).forEach(p -> {
            System.out.println(p);
        });
        System.out.println("age < 28 ----------------------");
        personnes.stream().filter(p -> p.getAge() < 28).forEach(System.out::println);
        System.out.println("poids < 50 ----------------------");
        personnes.stream().filter(p -> p.getPoids() < 50).forEach(System.out::println);
    }
}
  • wiersze 14–16: wyświetlają osoby w wieku <28;
  • wiersze 18–20: wyświetlają osoby o wadze <50;
  • wiersz 22: wykonuje to samo co wiersze 14–16, ale w bardziej zwięzły sposób;
  • wiersz 24: wykonuje to samo co wiersze 18–20, ale w bardziej zwięzły sposób;

Wyniki wykonania są następujące:

age < 28 ----------------------
{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"}
{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"}
poids < 50 ----------------------
{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"}
age < 28 ----------------------
{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"}
{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"}
poids < 50 ----------------------
{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"}

5.5. Przykład-05 – utworzenie obiektu Stream<T2> na podstawie obiektu Stream<T1>

  

Rozważmy następujący kod:


package dvp.java8.streams;

import java.util.List;

import dvp.data.Personne;
import dvp.data.Personnes;

public class Exemple05 {
  public static void main(String[] args) {
    // lista osób
    List<Personne> personnes = Personnes.get();
    // wyświetlenia
    System.out.println("Personne --> String ----------------------");
    personnes.stream().map(p -> p.getNom()).forEach(System.out::println);
    System.out.println("Personne --> Integer ----------------------");
    personnes.stream().map(p -> p.getAge()).forEach(System.out::println);
  }
}
  • W wierszu 13 metoda [Stream.map] ma następującą sygnaturę:
 

Parametrem metody [Stream.map] jest instancja interfejsu funkcjonalnego [Function] przedstawionego w paragrafie 4.3, którego jedyną metodą do zaimplementowania jest: R apply(T t). Widać, że na podstawie typu T funkcja [apply] generuje typ R. Metoda [Stream.map] wygeneruje zatem strumień Stream typu R na podstawie strumienia typu T (strumień typu T oznacza tutaj – w pewnym nadużyciu językowym, które zachowamy – strumień elementów typu T).

Przeanalizujmy teraz kod z przykładu:

  • wiersz 14: z osoby p zachowujemy jedynie imię. Otrzymujemy zatem strumień o typie String;
  • wiersz 14: z osoby p zachowujemy jedynie imię. Otrzymujemy zatem strumień o postaci Integer;

Otrzymane wyniki są następujące:

1
2
3
4
5
6
7
8
Personne --> String ----------------------
jean
marie
camille
Personne --> Integer ----------------------
20
10
30

5.6. Przykład-06 – inne metody klasy Stream<T>

  

Niektóre z 39 metod klasy Stream ilustrujemy za pomocą poniższego kodu:


package dvp.java8.streams;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import dvp.data.Personne;
import dvp.data.Personnes;

import java.util.Comparator;
import java.util.List;
import java.util.stream.Collectors;
import java.util.stream.DoubleStream;
import java.util.stream.IntStream;
import java.util.stream.Stream;

public class Exemple06 {

    // mapper jSON
    static private ObjectMapper jsonMapper = new ObjectMapper();

    public static void main(String[] args) throws JsonProcessingException {
        // lista osób
        List<Personne> personnes = Personnes.get();
        // wszystkie osoby
        affiche("all", personnes);
        // pierwsza osoba
        affiche("findFirst", personnes.stream().findFirst().get());
        // dowolna osoba
        affiche("findAny", personnes.stream().findAny().get());
        // osoby z wyjątkiem pierwszej
        affiche("skip 1", personnes.stream().skip(1L).collect(Collectors.toList()));
        // dwie pierwsze osoby
        affiche("limit 2", personnes.stream().limit(2L).collect(Collectors.toList()));
        // liczba osób
        affiche("count", personnes.stream().count());
        // najstarsza osoba
        affiche("age max", personnes.stream().max(Comparator.comparingInt(Personne::getAge)).get());
        // osoba o najmniejszej masie ciała
        affiche("poids min", personnes.stream().min(Comparator.comparingDouble(Personne::getPoids)).get());
        // ostatnia osoba w porządku alfabetycznym według nazwisk
        affiche("nom max", personnes.stream().max((p1, p2) -> p1.getNom().compareToIgnoreCase(p2.getNom())).get());
        // łączny wiek wszystkich osób
        affiche("âge total (reduce)", personnes.stream().map(p -> p.getAge()).reduce(0, (a1, a2) -> a1 + a2));
        // osoby w porządku rosnącym według wieku
        affiche("personnes par âge croissant",
                personnes.stream().sorted(Comparator.comparingInt(Personne::getAge)).collect(Collectors.toList()));
        // czy są osoby powyżej 100 lat?
        affiche("des personnes de + de 100 ans (anyMatch)", personnes.stream().anyMatch(p -> p.getAge() > 100));
        // czy wszystkie osoby mają co najwyżej 100 lat?
        affiche("des personnes de + de 100 ans (noneMatch)", personnes.stream().noneMatch(p -> p.getAge() > 100));
        // czy wszystkie osoby mają więcej niż 8 lat
        affiche("des personnes de + de 8 ans (allMatch)", personnes.stream().allMatch(p -> p.getAge() > 8));
        // osoby są pogrupowane według płci
        affiche("personnes regroupées par sexe", personnes.stream().collect(Collectors.groupingBy(p -> p.getSexe())));
        // usuwanie zduplikowanych elementów z listy
        affiche("distinct", Stream.of(1, 2, 1).distinct().collect(Collectors.toList()));
        // z Stream<Stream<T>> tworzymy Stream<T>
        affiche("flatMap", Stream.of(1, 2, 3).flatMap(i -> Stream.of(i, i + 10)).collect(Collectors.toList()));
        // z Stream<Stream<Integer>> tworzymy IntStream, którego sumę obliczamy
        affiche("flatMapToInt", Stream.of(1, 2, 3).flatMapToInt(i -> IntStream.of(i, i + 10)).sum());
        // z Stream<Stream<Integer>> tworzy się DoubleStream, a następnie tablicę
        affiche("flatMapToDouble", Stream.of(1, 2, 3).flatMapToDouble(i -> DoubleStream.of(i, i * 1.2)).toArray());
        // maksymalną wartość ze strumienia liczb całkowitych
        affiche("reduce Integer::max", Stream.of(1, 10, 8).reduce(Integer::max).get());
        // min strumienia liczb typu Double
        affiche("reduce Integer::min", Stream.of(1.5, 10.4, 8.9).reduce(Double::min).get());
        // średnia ze strumienia liczb całkowitych
        affiche("IntStream average", IntStream.of(1, 10, 8).average().getAsDouble());
        // statystyki strumienia liczb całkowitych
        affiche("IntStream summaryStatistics", IntStream.of(1, 10, 8).summaryStatistics());
    }

    public static <T> void affiche(String message, T value) throws JsonProcessingException {
        System.out.println(String.format("%s ----", message));
        System.out.println(jsonMapper.writeValueAsString(value));
    }
}
  • wiersze 72, 75: wyświetlają ciąg znaków jSON z drugiego parametru metody;
  • wiersz 24: wyświetla ciąg znaków jSON dla wszystkich osób. Otrzymujemy następujący wynik:
all ----
[{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"},{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"},{"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"}]

5.6.1. [findFirst]


// pierwsza osoba
affiche("findFirst", personnes.stream().findFirst().get());

Metoda [findFirst] zwraca pierwszy element strumienia, jeśli taki istnieje. Jej sygnatura jest następująca:

Wynik jest typu Optional<T>, typu wprowadzonego w Javie 8:

Klasa Optional<T> pozwala na odmienne traktowanie wskaźników null. Metoda, która powinna zwracać typ T mogący przyjmować wartość null, może zdecydować się na zwracanie typu Optional<T>. Metoda [Optional<T>.isPresent()] pozwala sprawdzić, czy metoda zwróciła wartość, czy nie. Poniższy kod [Exemple06b] ilustruje część działania klasy Optional<T>:


package dvp.java8.streams;

import java.util.Optional;

import com.fasterxml.jackson.core.JsonProcessingException;

public class Exemple06b {

    public static void main(String[] args) throws JsonProcessingException {
        // opcjonalne bez wartości
        Optional<Integer> o1 = m1();
        System.out.println(o1.isPresent());
        affiche(o1);
        // opcjonalny z wartością
        Optional<Integer> o2 = m2();
        System.out.println(o2.isPresent());
        affiche(o2);
    }

    private static void affiche(Optional<Integer> o1) {
        try {
            // pobieramy wartość parametru opcjonalnego
            // rzuca 1 wyjątek, jeśli nie ma wartości
            System.out.println(o1.get());
        } catch (Throwable th) {
            System.out.printf("%s : %s%n", th.getClass().getName(), th.getMessage());
        }

    }

    public static Optional<Integer> m1() {
        // brak wartości
        return Optional.empty();
    }

    public static Optional<Integer> m2() {
        // wartość
        return Optional.of(10);
    }
}

Uzyskane wyniki są następujące:


false
java.util.NoSuchElementException : No value present
true
10

Wróćmy do kodu ilustrującego działanie metody [findFirst]:


// pierwsza osoba
affiche("findFirst", personnes.stream().findFirst().get());
  • wiersz 2: aby uprościć kod, stosujemy metodę [get] na wyniku Optional<Personne> wygenerowanym przez metodę [findFirst]. Zgodnie z zasadami dobrej praktyki programistycznej należałoby wywołać metodę [Optional<Personne>.isPresent()] przed wywołaniem metody [get];

Otrzymany wynik jest następujący:

findFirst ----
{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"}

5.6.2. [findAny]


// dowolna osoba
affiche("findAny", personnes.stream().findAny().get());

Metoda [findAny] ma następującą sygnaturę:

 

Metoda [findAny] może zwrócić dowolny element strumienia. Podczas testów zauważono, że wykonanie sekwencyjne zwraca pierwszy element strumienia, podczas gdy wykonanie równoległe może faktycznie zwrócić dowolny element. Pokazuje to poniższy kod [Exemple06c]:


package dvp.java8.streams;

import java.util.List;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;

import dvp.data.Personne;
import dvp.data.Personnes;

public class Exemple06c {

    // mapper jSON
    static private ObjectMapper jsonMapper = new ObjectMapper();

    public static void main(String[] args) throws JsonProcessingException {
        // lista osób
        List<Personne> personnes = Personnes.get();
        // wszystkie osoby
        affiche("all", personnes);
        // dowolna osoba
        affiche("findAny parallèle", personnes.stream().parallel().findAny().get());
        // dowolna osoba
        affiche("findAny séquentiel", personnes.stream().findAny().get());
    }

    public static <T> void affiche(String message, T value) throws JsonProcessingException {
        System.out.println(String.format("%s ----", message));
        System.out.println(jsonMapper.writeValueAsString(value));
    }
}
  • wiersz 22: findAny wykonywany równolegle;
  • wiersz 24: findAny wykonywany sekwencyjnie;

Uzyskane wyniki są następujące:

1
2
3
4
5
6
all ----
[{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"},{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"},{"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"}]
findAny parallèle ----
{"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"}
findAny séquentiel ----
{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"}
  • wiersz 4: wykonanie równoległe zwróciło element nr 2 z listy osób. Mógł to być również inny element;
  • wiersz 6: wykonanie sekwencyjne zwróciło pierwszy element listy osób;

Zastosowanie metody [findAny] wydaje się mieć sens jedynie w przypadku równoległego przetwarzania strumienia.

5.6.3. [skip]


// osoby z wyjątkiem pierwszej
affiche("skip 1", personnes.stream().skip(1L).collect(Collectors.toList()));

Metoda [skip] ma następującą sygnaturę:

 

Metoda [skip] pomija pierwsze n elementów strumienia. Jak wskazuje powyższa dokumentacja, równoległe wykonywanie tej metody przynosi niewielki wzrost wydajności, a nawet może ją obniżyć. W rzeczywistości, aby pominąć pierwsze n elementów, wątki muszą się koordynować, co niweluje korzyści wydajnościowe wynikające z równoległości.

Metoda [skip] zwraca strumień Stream<Personne>, który jest przekształcany na typ List<Personne> przez metodę [collect], której sygnatura jest następująca:

 

Metoda [collect] przyjmuje jako parametr instancję typu [Collector], której sygnatura jest złożona. Istnieją predefiniowane implementacje typu [Collector], które w większości przypadków pozwalają uniknąć samodzielnej implementacji. W tym przypadku wykorzystano implementację [Collectors.toList()]. [Collectors] to klasa posiadająca wiele metod statycznych implementujących typ [Collector<T,A,R>]. To właśnie tam należy szukać, gdy chce się przekształcić Stream w standardową kolekcję Javy:

 

Niektóre z tych metod wykorzystamy w dalszej części.

Wynik wykonania jest następujący:

skip 1 ----
[{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"},{"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"}]

Pierwszy element listy (jean) został pominięty.

5.6.4. [limit]


// dwie pierwsze osoby
affiche("limit 2", personnes.stream().limit(2L).collect(Collectors.toList()));

Metoda [limit] ma następującą sygnaturę:

 

Metoda [limit] pozwala zachować tylko n pierwszych elementów strumienia. Nie nadaje się ona do przetwarzania równoległego.

Wynik wykonania jest następujący:

limit 2 ----
[{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"},{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"}]

5.6.5. [count]


// liczba osób
affiche("count", personnes.stream().count());

Metoda [count] ma następującą sygnaturę:

 

Metoda [count] zwraca liczbę elementów z Stream. Równoległe wykonywanie tej metody nie przynosi wzrostu wydajności, co pokazuje poniższy kod (Przykład06d1):


package dvp.java8.streams;

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

public class Exemple06d1 {
    public static void main(String[] args) {

        final long limite = 10_000_000L;
        // liczba procesorów
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lista liczb
        long début = new Date().getTime();
        List<Long> nombres = new ArrayList<>();
        for (long i = 0; i < limite; i++) {
            nombres.add(i);
        }
        System.out.printf("création de la liste des %s nombres en %s ms%n", limite, new Date().getTime() - début);
        // liczenie liczb – metoda sekwencyjna
        Stream<Long> sNombres = nombres.stream();
        début = new Date().getTime();
        long count = sNombres.count();
        System.out.printf("comptage séquentiel : compteur=%s, durée (ms)=%s%n", count, new Date().getTime() - début);
    }
}
  • wiersze 11–22: tworzy się plik Stream zawierający 10 milionów liczb;
  • wiersze 22–24: zliczanie Stream;

Wynik wykonania jest następujący:

1
2
3
La JVM a détecté [8] coeurs sur votre machine
création de la liste des 10000000 nombres en 4407 ms
comptage séquentiel : compteur=10000000, durée (ms)=67

Jeśli zastąpimy wiersz 22 kodu następującym (Przykład06d2):


Stream<Long> sNombres = nombres.stream().parallel();

otrzymujemy następujące wyniki:

1
2
3
La JVM a détecté [8] coeurs sur votre machine
création de la liste des 10000000 nombres en 4341 ms
comptage parallèle : compteur=10000000, durée (ms)=100

5.6.6. [max, min]


// najstarsza osoba
affiche("age max", personnes.stream().max(Comparator.comparingInt(Personne::getAge)).get());

Metoda [max] ma następujący podpis:

 

Metoda [max] zwraca maksymalną wartość strumienia przy użyciu komparatora przekazanego jej jako parametr. Comparator jest interfejsem funkcjonalnym, którego jedyna metoda do zaimplementowania ma następującą sygnaturę: int compare (T o1, T o2). Metoda ta musi zwracać -1, jeśli o1 < o2, 0, jeśli o1 ≥ o2, oraz +1, jeśli o1 > o2. Interfejs funkcjonalny Comparator posiada domyślnie wiele metod statycznych, które implementują interfejs Comparator dla najczęstszych przypadków. Tak więc w instrukcji:


affiche("age max", personnes.stream().max(Comparator.comparingInt(Personne::getAge)).get());

używamy metody statycznej [Comparator.comparingInt], której sygnatura jest następująca:

 

Typ ToIntFunction jest interfejsem funkcjonalnym:

 

Metoda [applyAsInt] interfejsu funkcjonalnego ToIntFunction generuje typ int na podstawie typu T. Wróćmy do naszego kodu:


affiche("age max", personnes.stream().max(Comparator.comparingInt(Personne::getAge)).get());

Rzeczywistym parametrem metody [Comparator.comparingInt] musi być tutaj lambda typu Personne --> int. Przekazujemy odwołanie do metody [Personne.getAge], która ma właśnie taką sygnaturę. W rezultacie otrzymamy osobę o największym wieku. Otrzymujemy typ Optional<Personne>, z którego wyodrębniamy wartość za pomocą metody [Optional.get]. Otrzymujemy następujący wynik:

age max ----
{"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"}

Równoległe obliczanie max nie przynosi korzyści w zakresie wydajności, jak pokazuje poniższy przykład: (Przykład06e1):


package dvp.java8.streams;

import java.util.ArrayList;
import java.util.Comparator;
import java.util.Date;
import java.util.List;
import java.util.Random;
import java.util.stream.Stream;

public class Exemple06e1 {
    public static void main(String[] args) {

        // data
        // final long limit = 100L;
        // final boolean verbose = true;
        final long limite = 10_000_000L;
        final boolean verbose = false;

        // liczba procesorów
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lista liczb
        long début = new Date().getTime();
        List<Long> nombres = new ArrayList<>();
        for (long i = 0; i < limite; i++) {
            nombres.add(new Random().nextLong());
        }
        System.out.printf("création de la liste des %s nombres en %s ms%n", limite, new Date().getTime() - début);
        // maksymalna wartość liczb – metoda sekwencyjna
        Stream<Long> sNombres = nombres.stream();
        Comparator<Long> compLong = (l1, l2) -> {
            if (verbose) {
                // wątek
                System.out.printf("[%s]", Thread.currentThread().getName());
            }
            // porównanie
            long v1 = l1.longValue();
            long v2 = l2.longValue();
            if (v1 < v2) {
                return -1;
            } else {
                if (v1 == v2) {
                    return 0;
                } else {
                    return +1;
                }
            }
        };
        début = new Date().getTime();
        // maksymalna długość = sNombres.max(Comparator.naturalOrder()).get();
        long max = sNombres.max(compLong).get();
        System.out.printf("%nmax séquentiel : max=%s, durée (ms)=%s%n", max, new Date().getTime() - début);
    }
}
  • wiersz 29: mamy strumień liczb losowych typu Long o nazwie limite;
  • wiersze 30–47: zmienna lambda compLong implementuje interfejs Comparator<Long>. Interfejs ten jest zazwyczaj implementowany przez metodę [Comparator.naturalOrder()] w wierszu 49. Jednak w tym przypadku chcemy wyświetlić wątek wykonania (wiersze 31–33). Dlatego sami implementujemy ten interfejs;
  • wiersz 50: wyszukiwanie metody max;

Otrzymujemy następujące wyniki:

 

Jeśli teraz zastąpimy wiersz 27 następującym (Przykład06e2):


Stream<Long> sNombres = nombres.stream().parallel();

otrzymujemy następujące wyniki:

 

Wykonywanie równoległe przebiegło zatem wolniej. Jeśli zwiększymy liczbę liczb do 10 milionów przy użyciu verbose=false, otrzymamy następujące wyniki:

1
2
3
4
La JVM a détecté [8] coeurs sur votre machine
création de la liste des 10000000 nombres en 3764 ms

max séquentiel : max=9223370471463514417, durée (ms)=53

dla wykonania sekwencyjnego:

1
2
3
4
La JVM a détecté [8] coeurs sur votre machine
création de la liste des 10000000 nombres en 3760 ms

max parallèle : max=9223365260999360873, durée (ms)=77

dla wykonania równoległego, które pozostaje zatem wolniejsze.

Metodę [Stream.min] stosuje się w podobny sposób:


// najlżejsza osoba
affiche("poids min", personnes.stream().min(Comparator.comparingDouble(Personne::getPoids)).get());

5.6.7. [reduce]


// łączny wiek wszystkich osób
affiche("âge total (reduce)", personnes.stream().map(p -> p.getAge()).reduce(0, (a1, a2) -> a1 + a2));

Metoda [reduce] została przedstawiona w punkcie 5.3. Wiersz 2 powyżej sumuje wiek wszystkich osób. Wynik jest następujący:

âge total (reduce) ----
60

5.6.8. [sorted]


// osoby w porządku rosnącym według wieku
affiche("personnes par âge croissant",
                personnes.stream().sorted(Comparator.comparingInt(Personne::getAge)).collect(Collectors.toList()));
// osoby w porządku alfabetycznym według nazwisk
List<Personne> lPersonnes=personnes.stream().sorted((p1, p2) -> p1.getNom().compareTo(p2.getNom())).collect(Collectors.toList());
affiche("personnes par ordre alphabétique des noms", lPersonnes);

Metoda [sorted] (wiersze 3 i 5) ma następującą sygnaturę:

 

Metoda [sorted] przyjmuje jako parametr typ [Comparator] opisany w paragrafie 5.6.6 dla metod min i max. Pozwala ona na sortowanie obiektu typu Stream zgodnie z kolejnością określona przez komparator przekazany jako parametr. Widzieliśmy już, że interfejs [Comparator] oferuje domyślnie kilka metod statycznych implementujących popularne komparatory, w szczególności dla liczb i ciągów znaków. W tym przypadku używamy metody [Comparator.comparingInt], która przyjmuje jako parametr typ ToIntFunction, będący funkcjonalnym interfejsem metody [applyAsInt] o następującej sygnaturze: int applyAsInt(T t). W tym przypadku rzeczywistym parametrem przekazanym do metody [Comparator.comparingInt] w wierszu 3 jest odwołanie do metody [Personne.age], która zwraca wiek osoby.

Interfejs [Comparator] nie udostępnia metod statycznych do porównywania ciągów znaków. W wierszu 5 samodzielnie tworzymy lambda implementującą jedyną metodę tego interfejsu: int compare(T t1, T t2)


(p1, p2) -> p1.getNom().compareTo(p2.getNom())

Ta funkcja lambda porównuje imiona osób. Uzyskane wyniki są następujące:

1
2
3
4
personnes par âge croissant ----
[{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"},{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"},{"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"}]
personnes par ordre alphabétique des noms ----
[{"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"},{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"},{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"}]

Równoległe wykonywanie sortowania nie wydaje się możliwe, jak pokazuje poniższy kod (Przykład06f1):


package dvp.java8.streams;

import java.util.ArrayList;
import java.util.Comparator;
import java.util.Date;
import java.util.List;
import java.util.Random;
import java.util.stream.Collectors;
import java.util.stream.Stream;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;

public class Exemple06f1 {
    // mapper jSON
    static ObjectMapper jsonMapper = new ObjectMapper();

    public static void main(String[] args) throws JsonProcessingException {

        // dane
        final long limite = 100L;
        final boolean verbose = true;
//         długość końcowa = 10_000_000L;
//         final boolean verbose = false;

        // liczba procesorów
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lista liczb
        long début = new Date().getTime();
        List<Integer> nombres = new ArrayList<>();
        for (long i = 0; i < limite; i++) {
            nombres.add(new Random().nextInt(1000));
        }
        System.out.printf("création de la liste des %s nombres en %s ms%n", limite, new Date().getTime() - début);
        // sortowanie liczb – metoda sekwencyjna
        Stream<Integer> sNombres = nombres.stream();
        début = new Date().getTime();
        Comparator<Integer> compInt = (i1, i2) -> {
            if (verbose) {
                // wątek
                System.out.printf("[%s]", Thread.currentThread().getName());
            }
            // porównanie
            int v1 = i1.intValue();
            int v2 = i2.intValue();
            if (v1 < v2) {
                return +1;
            } else {
                if (v1 == v2) {
                    return 0;
                } else {
                    return -1;
                }
            }
        };
        if (verbose) {
            affiche("nombres", sNombres.sorted(compInt).collect(Collectors.toList()));
        }
        System.out.printf("tri séquentiel : durée (ms)=%s%n", new Date().getTime() - début);
    }

    public static <T> void affiche(String message, T value) throws JsonProcessingException {
        System.out.println(String.format("%s ----", message));
        System.out.println(jsonMapper.writeValueAsString(value));
    }

}
  • wiersze 30–36: tworzy się strumień liczb losowych limite;
  • w wierszu 32 przekazujemy lambda compInt (wiersze 38–55) do metody [sorted]. Ta lambda sortuje liczby w porządku malejącym i wyświetla wątek, który ją wykonuje.

Otrzymane wyniki są następujące:

 

Jeśli w powyższym kodzie zastąpimy wiersz 36 następującym (Przykład06f2):


        Stream<Integer> sNombres = nombres.stream().parallel();        

otrzymujemy następujące wyniki:

 

Zaskakujące jest to, że sortowanie strumienia liczb odbyło się przy użyciu tylko jednego wątku. Nie wystąpiła żadna równoległość. A może coś mi umyka?

5.6.9. [anyMatch, noneMatch, allMatch]


// czy są osoby w wieku powyżej 100 lat?
affiche("des personnes de + de 100 ans (anyMatch)", personnes.stream().anyMatch(p -> p.getAge() > 100));
// czy wszystkie osoby mają co najwyżej 100 lat?
affiche("des personnes de + de 100 ans (noneMatch)", personnes.stream().noneMatch(p -> p.getAge() > 100));
// czy wszystkie osoby mają więcej niż 8 lat?
affiche("des personnes de + de 8 ans (allMatch)", personnes.stream().allMatch(p -> p.getAge() > 8));

W wierszach 2, 4 i 6 metody [anyMatch, noneMatch, allMatch] przyjmują jako parametr typ Predicate opisany w paragrafie 4.2. Wykonują one zatem filtrowanie. Wszystkie trzy zwracają wartość logiczną:

  • anyMatch zwraca true, jeśli istnieje co najmniej jeden element typu Stream spełniający kryteria filtru;
  • noneMatch zwraca true, jeśli nie istnieje żaden element z Stream spełniający kryteria filtru;
  • allMatch zwraca true, jeśli wszystkie elementy z Stream spełniają kryteria filtru;

Otrzymane wyniki są następujące:

1
2
3
4
5
6
des personnes de + de 100 ans (anyMatch) ----
false
des personnes de + de 100 ans (noneMatch) ----
true
des personnes de + de 8 ans (allMatch) ----
true

5.6.10. [collect(Collectors.groupingBy)]


// osoby są pogrupowane według płci
affiche("personnes regroupées par sexe", personnes.stream().collect(Collectors.groupingBy(p -> p.getSexe())));

Metoda [collect] została przedstawiona w paragrafie 5.6.3. Jej parametrem jest implementacja interfejsu [Collector]. Klasa [Collectors] udostępnia szereg metod statycznych implementujących interfejs [Collector]. Do tej pory korzystaliśmy z metody [Collectors.toList()]. W tym miejscu wykorzystujemy metodę statyczną [Collectors.groupingBy], która tworzy słownik na podstawie Stream. Jej sygnatura jest następująca:

 

Metoda [groupingBy] tworzy na podstawie typu Stream<T> typ Map<K,List<T>>. Klucz K jest dostarczany przez parametr metody [groupingBy] typu Function<T,K>, której jedyna metoda ma sygnaturę: K apply(T t). Jeśli chcemy utworzyć słownik indeksowany według płci osób, musimy podać funkcję generującą płeć na podstawie osoby. W tym przypadku jako rzeczywisty parametr metody [groupingBy] przekazujemy odwołanie do metody [Personne.getSexe]. Uzyskane wyniki są następujące:

personnes regroupées par sexe ----
{"HOMME":[{"nom":"jean","age":20,"poids":70.0,"sexe":"HOMME"}],"FEMME":[{"nom":"marie","age":10,"poids":30.0,"sexe":"FEMME"},{"nom":"camille","age":30,"poids":55.0,"sexe":"FEMME"}]}

W wierszu 2 znajduje się ciąg znaków jSON pochodzący ze słownika indeksowanego dwoma kluczami: HOMME i FEMME.

Obliczenia równoległe nie przynoszą wzrostu wydajności, co pokazuje poniższy przykład (Przykład06g1):


package dvp.java8.streams;

import java.util.ArrayList;
import java.util.Date;
import java.util.List;
import java.util.Map;
import java.util.Random;
import java.util.function.Function;
import java.util.stream.Collectors;
import java.util.stream.Stream;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;

public class Exemple06g1 {

    // mppeur jSON
    static ObjectMapper jsonMapper = new ObjectMapper();

    public static void main(String[] args) throws JsonProcessingException {

        // dane
        final long limite = 100L;
        final boolean verbose = true;
//         długie ograniczenie końcowe = 10_000_000L;
//         final boolean verbose = false;

        // liczba procesorów
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lista liczb
        long début = new Date().getTime();
        List<Integer> nombres = new ArrayList<>();
        for (long i = 0; i < limite; i++) {
            nombres.add(new Random().nextInt(1000));
        }
        System.out.printf("création de la liste des %s nombres en %s ms%n", limite, new Date().getTime() - début);
        // grupowanie liczb po sto – metoda sekwencyjna
        Stream<Integer> sNombres = nombres.stream();
        Function<Integer, Integer> groupByCent = n -> {
            if (verbose) {
                System.out.printf("[%s]", Thread.currentThread().getName());
            }
            return n / 100;
        };
        début = new Date().getTime();
        // Map<Integer, List<Integer>> lNombres = sNombres.collect(Collectors.groupingBy(liczba -> liczba / 100));
        Map<Integer, List<Integer>> lNombres = sNombres.collect(Collectors.groupingBy(groupByCent));
        System.out.printf("%nregroupement séquentiel : durée (ms)=%s%n", new Date().getTime() - début);
        // wyniki
        if (verbose) {
            affiche("nombres regroupés", lNombres);
        }
    }

    public static <T> void affiche(String message, T value) throws JsonProcessingException {
        System.out.println(String.format("%s ----", message));
        System.out.println(jsonMapper.writeValueAsString(value));
    }

}
  • wiersze 23–38: tworzenie strumienia liczb limite;
  • w wierszu 47 liczby są grupowane po stu. Wykorzystuje się funkcję lambda z wierszy 39–44, aby móc wyświetlić wątek wykonania;

Wyniki wykonania są następujące:

 

Jeśli w kodzie zastąpimy wiersz 38 następującym wierszem (Przykład06g2):


Stream<Integer> sNombres = nombres.stream().parallel();            

otrzymujemy następujące wyniki:

 

Widać, że równoległe wykonywanie grupowania spowodowało spadek wydajności.

5.6.11. [distinct]


// usuwanie zduplikowanych elementów z listy
affiche("distinct", Stream.of(1, 2, 1).distinct().collect(Collectors.toList()));

Metoda [distinct] ma następującą sygnaturę:

 

Pozwala ona na usunięcie duplikatów ze strumienia. Metoda [Stream.of] (wiersz 2) ma następującą sygnaturę:

 

Pozwala ona na utworzenie obiektu Stream na podstawie wyraźnie podanych wartości. Wyniki wykonania są następujące:

distinct ----
[1,2]

5.6.12. [flatMap]


// ze Stream<Stream<T>> tworzy się Stream<T>
affiche("flatMap", Stream.of(1, 2, 3).flatMap(i -> Stream.of(i, i + 10)).collect(Collectors.toList()));

Metoda [flatMap] ma następującą sygnaturę:

 

Metoda [flatMap] przyjmuje jako parametr funkcję, która:

  • przyjmuje jako parametr element typu T z Stream;
  • zwraca strumień typu Stream<R>;

Gdyby zamiast metody [flatMap] zastosowano by metodę [map] opisaną w paragrafie 5.5, wynikiem byłby typ Stream<Stream<R>>, w którym każdy element typu T z początkowego strumienia dałby początek elementowi Stream<R>. Metoda [flatMap] zwraca natomiast typ Stream<R>. Spłaszcza (flatten) ona różne strumienie Stream<R> do jednego strumienia. Pokazują to wyniki wykonania poprzedniego kodu:

flatMap ----
[1,11,2,12,3,13]

Istnieją wyspecjalizowane warianty funkcji [flatMap]:


// z Stream<IntStream> tworzymy IntStream, którego sumę obliczamy
affiche("flatMapToInt", Stream.of(1, 2, 3).flatMapToInt(i -> IntStream.of(i, i + 10)).sum());

Metoda [flatMapToInt] ma następującą sygnaturę:

 

Metoda [flatMapToInt] przyjmuje jako parametr funkcję, której wynikiem jest następujący typ IntStream:

 

IntStream jest strumieniem pochodnym od int. Ten typ jest lepszy od typu Stream<Integer>, ponieważ jego przetwarzanie pozwala uniknąć operacji boxing/unboxing między typami Integer i int. Interfejs ten przejmuje wiele metod typu Stream<T> i dodaje do nich inne, w tym wspomnianą powyżej metodę [sum], która sumuje elementy typu IntStream.

Poniższy kod ilustruje użycie analogicznej metody [flatMapToDouble]:


// z Stream<DoubleStream> tworzy się DoubleStream, a następnie tablicę
affiche("flatMapToDouble", Stream.of(1, 2, 3).flatMapToDouble(i -> DoubleStream.of(i, i * 1.2)).toArray());

Metoda [DoubleStream.toArray] umożliwia przejście z typu DoubleStream do typu double[].

Wyniki dla tych dwóch przykładów są następujące:

1
2
3
4
flatMapToInt ----
42
flatMapToDouble ----
[1.0,1.2,2.0,2.4,3.0,3.5999999999999996]

Poniższy przykład pokazuje wzrost wydajności uzyskany dzięki przejściu z typu Stream<Long> na typ LongStream (Przykład06i1):


package dvp.java8.streams;

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

public class Exemple06i1 {
    public static void main(String[] args) {

        final long limite = 10_000_000L;
        // liczba procesorów
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lista liczb
        long début = new Date().getTime();
        List<Long> nombres = new ArrayList<>();
        for (long i = 0; i < limite; i++) {
            nombres.add(i);
        }
        System.out.printf("création de la liste des %s nombres en %s ms%n", limite, new Date().getTime() - début);
        // suma liczb – metoda sekwencyjna
        début = new Date().getTime();
        long somme = nombres.stream().reduce(0L, (s, i) -> s + i);
        System.out.printf("somme séquentielle du Stream<Integer> : somme=%s, durée (ms)=%s%n", somme, new Date().getTime() - début);
    }
}
  • wiersz 22: obliczenie sumy strumienia liczb typu Long;

Otrzymujemy następujące wyniki:

1
2
3
La JVM a détecté [8] coeurs sur votre machine
création de la liste des 10000000 nombres en 4537 ms
somme séquentielle du Stream<Integer> : somme=49999995000000, durée (ms)=226

Teraz zastąpmy wiersz 22 następującym (Przykład06i2):


long somme = nombres.stream().mapToLong(n -> n.longValue()).sum();

Metoda Stream<Integer>.mapToLong pozwala nam uzyskać strumień typu LongStream zawierający elementy typu pierwotnego long, który następnie sumujemy za pomocą funkcji sum. Otrzymujemy wówczas następujące wyniki:

1
2
3
La JVM a détecté [8] coeurs sur votre machine
création de la liste des 10000000 nombres en 4511 ms
somme séquentielle du LongStream : somme=49999995000000, durée (ms)=99

Wzrost wydajności jest wyraźny.

5.6.13. Metody przepływu liczb pierwotnych


// maksymalna wartość z ciągu liczb całkowitych
affiche("IntStream max", IntStream.of(1, 10, 8).max());
// minimum z ciągu liczb typu double
affiche("DoubleStream min", DoubleStream.of(1.5, 10.4, 8.9).min());
// średnia z ciągu liczb typu int
affiche("IntStream average", IntStream.of(1, 10, 8).average().getAsDouble());
// statystyki z ciągu liczb typu int
affiche("IntStream summaryStatistics", IntStream.of(1, 10, 8).summaryStatistics());

Strumienie wartości pierwotnych (int, long, double) oferują metody dostosowane do tych typów. Wynik wykonania powyższego kodu jest następujący:

1
2
3
4
5
6
7
8
IntStream max ----
{"asInt":10,"present":true}
DoubleStream min ----
{"asDouble":1.5,"present":true}
IntStream average ----
6.333333333333333
IntStream summaryStatistics ----
{"count":3,"sum":19,"min":1,"max":10,"average":6.333333333333333}
  • Wynikiem linii 2 kodu jest typ OptionalInt, analogiczny do typu Optional<Integer>. Wartość przechowywaną w tym obiekcie można uzyskać za pomocą metody [getAsInt()]. Obecność wartości można sprawdzić za pomocą metody [isPresent()]. Wiersz 2 wyników nie oznacza, że klasa [OptionalInt] posiada pola o nazwach [asInt, present]. Domyślnie biblioteka jSON wykorzystuje wszystkie publiczne metody getX i isY obiektu, który ma zostać zserializowany do jSON. W tym przypadku istnieje metoda [getAsInt] oraz inna metoda [isPresent], podczas gdy pola [asInt, present] w ogóle nie istnieją;
  • wynikiem wiersza 4 kodu jest typ OptionalDouble analogiczny do typu Optional<Double>;
  • wynikiem linii 6 kodu jest typ OptionalDouble, którego wartość można uzyskać za pomocą metody [getAsDouble()]. Metoda [average] oblicza średnią z ciągu liczb;
  • wynikiem linii 8 kodu jest typ IntSummaryStatistics zdefiniowany w następujący sposób:
 

Widać, że uzyskany obiekt IntSummaryStatistics dostarcza różnych informacji o strumieniu liczb, takich jak liczba wartości, suma, maksymalna wartość, minimalna wartość oraz średnia.