Skip to content

5. Het type Stream<T> uit Java 8

5.1. Voorbeeld-01 - de klasse Stream

De bewerkingen op de streams Observable hebben veel overeenkomsten met de streams Stream. Een verschil is dat een element van een stream Stream pas kan worden verwerkt nadat de volledige stream Stream is opgehaald, terwijl een element van een stream Observable (waargenomen) zodra het is ontvangen, zonder te wachten tot de volledige stream Observable is ontvangen. Een ander verschil is dat, zodra de Stream is ontvangen, de waarden ervan worden benut door ze één voor één uit de Stream te halen (pull). Voor de observable is dit anders. Zodra deze een waarde uitzendt, wordt deze naar de abonnee gepusht (push).

Verschillende klassen implementeren het concept van Stream. Hier presenteren we de klasse Stream<T>:

Image

De klasse Stream bevat maar liefst 39 methoden. We zullen er een aantal bespreken. Laten we de volgende code eens bekijken:

  

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) {
        // lijst met personen
        List<Personne> personnes = Personnes.get();
        // weergave 1
        personnes.stream().forEach(p -> {
            System.out.println(p);
        });
        System.out.println("----------------");
        // weergave 2
        personnes.stream().forEach(System.out::println);
    }
}
  • regel 11: er wordt een lijst met personen aangemaakt;
  • regel 13: op basis van deze lijst wordt een Stream aangemaakt. Alle collecties kunnen op deze manier worden omgezet in Stream-streams. Hierdoor kunt u profiteren van alle methoden van deze klasse, waarmee de elementen van de collectie beknopter kunnen worden verwerkt dan met lussen. Ook kunt u zo profiteren van de parallelle verwerking van de elementen wanneer dit mogelijk is;
  • regel 13: de methode [Stream.forEach] heeft de volgende signatuur:
 

We zien dat de parameter van de methode de functionele interface [Consumer<T>] is, die in paragraaf 4.4 is beschreven; een interface waarvan de enige methode gebruikmaakt van het type T en niets retourneert.

  • in de code:

        personnes.stream().forEach(p -> {
            System.out.println(p);
});
  • [personnes.stream()] genereert een stroom van elementen van het type [Personne] die de methode [forEach] voeden. De parameter p is van het type [Personne] en de opgegeven lambda-functie geeft deze persoon weer;

De bovenstaande code kan als volgt worden vereenvoudigd (regel 18):


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

In plaats van de waarde van een lambda-functie als parameter door te geven, geven we de referentie van een bestaande methode door, in dit geval de methode println van de klasse System.out. Deze methode moet uiteraard de juiste signatuur hebben, in dit geval de signatuur van de methode [Consumer.accept]: void accept(T t). Zoals eerder vermeld, zal de parameter van de methode [accept] van het type [Personne] zijn;

We krijgen de volgende resultaten:

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"}

Zodra een Stream is verwerkt, kan deze niet meer worden gebruikt. Deze moet opnieuw worden aangemaakt als men deze opnieuw wil gebruiken. Dit wordt geïllustreerd door de volgende code [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) {
        // personenstroom
        Stream<Personne> personnes = Personnes.get().stream();
        // weergave 1
        personnes.forEach(p -> {
            System.out.println(p);
        });
        System.out.println("----------------");
        // weergave 2
        personnes.forEach(System.out::println);
    }
}
  • regel 11: om de code te optimaliseren, wordt besloten om de Stream slechts één keer op te bouwen. De verkregen resultaten zijn dan als volgt:

{"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)

Telkens wanneer men een Stream wil gebruiken, moet deze opnieuw worden gegenereerd, zelfs als deze al eerder is gegenereerd.

5.2. Voorbeeld-02 - parallelle verwerking van elementen van een stream

  

Laten we de volgende code eens bekijken:


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) {
        // lijst met personen
        List<Personne> personnes = Personnes.get();
        // weergave 1
        personnes.stream().forEach(Exemple02::affiche);
        System.out.println("-----------------");
        // weergave 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());
    }
}
  • regels 19-21: de methode [affiche] schrijft de tekenreeks jSON van een persoon naar de console, evenals de naam van de uitvoeringsthread waarin de weergave plaatsvindt;
  • regel 13: geeft een lijst met personen weer. Merk op dat de parameter van de methode [forEach] de referentie is van de voorgaande statische methode;
  • regel 16: we doen hetzelfde, maar met de methode [parallel] vragen we dat de verwerking van de elementen van de stream parallel plaatsvindt in meerdere threads. Niet elke verwerking kan parallel plaatsvinden. Hier moeten we ervan uitgaan dat de weergavevolgorde niet van belang is, omdat bij een parallelle verwerking de uitvoervolgorde van de threads niet gegarandeerd is. We zien overigens een syntaxis die zowel voor de Stream als voor de Observable alomtegenwoordig zal worden:
flux.m1(e1->...).m2(e2->..).m3(e3->...)...
  • (vervolg)
    • flux genereert e1-elementen die als invoer dienen voor de methode m1;
    • flux.m1 is op zijn beurt een stroom van e2-elementen die de methode m2 voeden;
    • flux.m1.m2 is een stroom van e3-elementen die de methode m3 voeden;

Het type van de elementen e1, e2 en e3 kan veranderen naarmate de oorspronkelijke stroom wordt verwerkt.

De uitvoering van deze code levert het volgende resultaat op:

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

We zien dat de parallelle uitvoering (regels 5-7) plaatsvond op drie verschillende threads en niet de volgorde van de elementen heeft gerespecteerd die overeenkomt met die van de regels 1-3. In dit document zullen we weinig aandacht besteden aan de parallelle verwerking van de elementen van een Stream, omdat we dan moeten ingaan op de voorwaarden die deze verwerking mogelijk maken. We ontdekken dan dat maar weinig bewerkingen parallel kunnen worden uitgevoerd. Een van de bewerkingen die zich van nature leent voor parallellisme is de som van de numerieke elementen van een stream, die we nu zullen presenteren.

5.3. Voorbeeld-03 – parallelle verwerking van de elementen van een stream

  

Laten we de volgende code eens bekijken (Voorbeeld 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;
        // aantal processors
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lijst met getallen
        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);
        // som van getallen – sequentiële methode
        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);
    }
}
  • In regel 22 gebruiken we de methode [reduce], waarvan de signatuur als volgt is:
  • de methode [reduce] werkt met elementen van het type T;
  • de methode [reduce] past dezelfde verwerking toe op alle elementen van een stroom: de beginwaarde van een accumulator wordt als eerste parameter opgegeven. Een methode die de functionele interface [BinaryOperator] [2] implementeert, wordt als tweede parameter doorgegeven: op basis van elk element en de accumulator levert deze methode een nieuwe waarde voor de accumulator op. De eindwaarde hiervan is de waarde die door de methode [reduce] wordt geretourneerd. De code [3] illustreert dit mechanisme. De methode [apply] is de methode van de functionele interface [BinaryOperator] [2];

Laten we teruggaan naar de voorbeeldcode:

  • regel 12: we geven het aantal cores weer zoals waargenomen door de JVM;
  • regels 15-18: er wordt een lijst met 10 miljoen getallen aangemaakt;
  • regel 22: de som van deze getallen wordt sequentieel berekend met één enkele thread;

We krijgen de volgende resultaten:

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

Laten we nu regel 22 van de code vervangen door het volgende (Voorbeeld03b):


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

We vragen dat de elementen van de Stream parallel worden verwerkt met behulp van meerdere threads. Dit is mogelijk omdat de volgorde waarin de getallen worden opgeteld niet van belang is. We kunnen dus n1 getallen toewijzen aan een thread T1, n2 getallen aan een thread T2, ... en uiteindelijk de sommen van deze verschillende threads bij elkaar optellen. We krijgen dan de volgende resultaten:

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

Er is dus vrijwel geen prestatiewinst. In de volgende voorbeelden zal dit vaak het geval zijn. Het beheer van threads kost zelf al veel tijd. De bewerking die door elke kern wordt uitgevoerd, moet voldoende complex zijn om prestatiewinst te realiseren. Dit blijkt uit het volgende voorbeeld (Voorbeeld03c):


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;
        // aantal processors
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lijst met getallen
        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);
        // som van de getallen - sequentiële methode
        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);
    }
}
  • regel 30: we gebruiken opnieuw de methode [reduce], waaraan we als parameter de referentie van de methode uit de regels 23-29 doorgeven;
  • regel 28: de methode [bo] levert de som van haar twee parameters op;
  • regels 24-27: kunstmatig laten we de thread 1 milliseconde wachten om intensief werk te simuleren;

Dit levert dan de volgende resultaten op:

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

Als we nu regel 30 vervangen door de volgende:


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

krijgen we de volgende resultaten:

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

We zien duidelijk de prestatiewinst die wordt behaald door de somberekening parallel uit te voeren. Bij de verwerking van 8 getallen:

  • wacht de sequentiële thread 8 keer 1 milliseconde, dus 8 ms;
  • de 8 parallelle threads wachten elk tegelijkertijd 1 milliseconde (voor het gemak in gedachten), dus in totaal 1 milliseconde voor de 8 getallen;

We kunnen dus verwachten dat de parallelle uitvoering 8 keer sneller verloopt dan de sequentiële uitvoering. Dat is hier ongeveer het geval.

5.4. Voorbeeld-04 - een stream filteren

  

Laten we de volgende code eens bekijken:


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) {
        // lijst met personen
        List<Personne> personnes = Personnes.get();
        // weergaven
        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);
    }
}
  • regel 14: de methode [Stream.filter] heeft de volgende signatuur:
 
  • de methode [filter] verwacht als parameter een instantie van de functionele interface [Predicate], die in paragraaf 4.2 is beschreven en waarvan de enige te implementeren methode de volgende is: boolean test(T t);
  • de methode [filter] retourneert de elementen van de Stream die voldoen aan Predicate. Ze dient dus om Stream te filteren;

Laten we de volgende code eens bekijken:


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) {
        // lijst met personen
        List<Personne> personnes = Personnes.get();
        // weergaven
        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);
    }
}
  • regels 14-16: geven de personen weer met een leeftijd <28;
  • regels 18-20: geven de personen weer met een gewicht <50;
  • regel 22: doet hetzelfde als de regels 14-16, maar op een beknoptere manier;
  • regel 24: doet hetzelfde als de regels 18-20, maar op een beknoptere manier;

De uitvoer is als volgt:

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. Voorbeeld-05 - een Stream<T2> aanmaken op basis van een Stream<T1>

  

Laten we de volgende code eens bekijken:


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) {
    // lijst met personen
    List<Personne> personnes = Personnes.get();
    // weergaven
    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);
  }
}
  • regel 13: de methode [Stream.map] heeft de volgende signatuur:
 

De parameter van de methode [Stream.map] is een instantie van de functionele interface [Function] die in paragraaf 4.3 wordt beschreven en waarvan de enige te implementeren methode is: R apply(T t). We zien dat de functie [apply] op basis van een type T een type R produceert. De methode [Stream.map] zal dus een stroom Stream van het type R genereren op basis van een stroom van het type T (een stroom van het type T betekent hier, in een figuurlijke betekenis die we zullen aanhouden, een stroom van elementen van het type T).

Laten we nu de code van het voorbeeld bekijken:

  • regel 14: van een persoon p wordt alleen de naam behouden. We krijgen dus een stroom van het type String;
  • regel 14: van een persoon p wordt alleen de naam behouden. We krijgen dus een stroom van Integer;

De verkregen resultaten zijn als volgt:

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

5.6. Voorbeeld-06 - andere methoden van de klasse Stream<T>

  

We illustreren enkele van de 39 methoden van de klasse Stream met de volgende code:


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 {
        // lijst met personen
        List<Personne> personnes = Personnes.get();
        // alle personen
        affiche("all", personnes);
        // de eerste persoon
        affiche("findFirst", personnes.stream().findFirst().get());
        // een willekeurige persoon
        affiche("findAny", personnes.stream().findAny().get());
        // de personen zonder de eerste
        affiche("skip 1", personnes.stream().skip(1L).collect(Collectors.toList()));
        // de eerste twee personen
        affiche("limit 2", personnes.stream().limit(2L).collect(Collectors.toList()));
        // het aantal personen
        affiche("count", personnes.stream().count());
        // de oudste persoon
        affiche("age max", personnes.stream().max(Comparator.comparingInt(Personne::getAge)).get());
        // de lichtste persoon
        affiche("poids min", personnes.stream().min(Comparator.comparingDouble(Personne::getPoids)).get());
        // de laatste persoon in alfabetische volgorde van de namen
        affiche("nom max", personnes.stream().max((p1, p2) -> p1.getNom().compareToIgnoreCase(p2.getNom())).get());
        // de totale leeftijd van alle personen
        affiche("âge total (reduce)", personnes.stream().map(p -> p.getAge()).reduce(0, (a1, a2) -> a1 + a2));
        // de personen in oplopende volgorde van leeftijd
        affiche("personnes par âge croissant",
                personnes.stream().sorted(Comparator.comparingInt(Personne::getAge)).collect(Collectors.toList()));
        // zijn er personen ouder dan 100 jaar?
        affiche("des personnes de + de 100 ans (anyMatch)", personnes.stream().anyMatch(p -> p.getAge() > 100));
        // zijn alle personen maximaal 100 jaar oud?
        affiche("des personnes de + de 100 ans (noneMatch)", personnes.stream().noneMatch(p -> p.getAge() > 100));
        // zijn alle personen ouder dan 8 jaar?
        affiche("des personnes de + de 8 ans (allMatch)", personnes.stream().allMatch(p -> p.getAge() > 8));
        // De mensen worden ingedeeld naar geslacht
        affiche("personnes regroupées par sexe", personnes.stream().collect(Collectors.groupingBy(p -> p.getSexe())));
        // dubbele elementen uit een lijst verwijderen
        affiche("distinct", Stream.of(1, 2, 1).distinct().collect(Collectors.toList()));
        // van een Stream<Stream<T>> maken we een Stream<T>
        affiche("flatMap", Stream.of(1, 2, 3).flatMap(i -> Stream.of(i, i + 10)).collect(Collectors.toList()));
        // van een Stream<Stream<Integer>> maken we een IntStream waarvan we de som berekenen
        affiche("flatMapToInt", Stream.of(1, 2, 3).flatMapToInt(i -> IntStream.of(i, i + 10)).sum());
        // van een Stream<Stream<Integer>>, maken we een DoubleStream en vervolgens een array
        affiche("flatMapToDouble", Stream.of(1, 2, 3).flatMapToDouble(i -> DoubleStream.of(i, i * 1.2)).toArray());
        // max van een stroom van gehele getallen
        affiche("reduce Integer::max", Stream.of(1, 10, 8).reduce(Integer::max).get());
        // min van een stroom van Double
        affiche("reduce Integer::min", Stream.of(1.5, 10.4, 8.9).reduce(Double::min).get());
        // gemiddelde van een stroom van gehele getallen
        affiche("IntStream average", IntStream.of(1, 10, 8).average().getAsDouble());
        // statistieken van een stroom van gehele getallen
        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));
    }
}
  • regels 72, 75: geven de tekenreeks jSON van de tweede parameter van de methode weer;
  • regel 24: geeft de tekenreeks jSON weer voor alle personen. Dit levert het volgende resultaat op:
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]


// de eerste persoon
affiche("findFirst", personnes.stream().findFirst().get());

De methode [findFirst] retourneert het eerste element van een stream, indien aanwezig. De signatuur is als volgt:

Het resultaat is van het type Optional<T>, een type dat in Java 8 is geïntroduceerd:

Met de klasse Optional<T> kunnen null-verwijzingen op een andere manier worden beheerd. Een methode die een type T moet retourneren dat de waarde null kan hebben, kan ervoor kiezen om een type Optional<T> te retourneren. Met de methode [Optional<T>.isPresent()] kun je nagaan of de methode al dan niet een waarde heeft geretourneerd. De volgende code [Exemple06b] illustreert een deel van de werking van 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 {
        // optioneel zonder waarde
        Optional<Integer> o1 = m1();
        System.out.println(o1.isPresent());
        affiche(o1);
        // optioneel met waarde
        Optional<Integer> o2 = m2();
        System.out.println(o2.isPresent());
        affiche(o2);
    }

    private static void affiche(Optional<Integer> o1) {
        try {
            // de waarde van de optionele parameter wordt opgehaald
            // genereert 1 uitzondering als er geen waarde is
            System.out.println(o1.get());
        } catch (Throwable th) {
            System.out.printf("%s : %s%n", th.getClass().getName(), th.getMessage());
        }

    }

    public static Optional<Integer> m1() {
        // geen waarde
        return Optional.empty();
    }

    public static Optional<Integer> m2() {
        // een waarde
        return Optional.of(10);
    }
}

De verkregen resultaten zijn als volgt:


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

Laten we teruggaan naar de voorbeeldcode van de methode [findFirst]:


// de eerste persoon
affiche("findFirst", personnes.stream().findFirst().get());
  • regel 2: om de code te vereenvoudigen, gebruiken we de methode [get] op de Optional<Personne> die door de methode [findFirst] is gegenereerd. In een nette code zouden we de methode [Optional<Personne>.isPresent()] moeten aanroepen voordat we de methode [get] aanroepen;

Het verkregen resultaat is als volgt:

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

5.6.2. [findAny]


// een willekeurig persoon
affiche("findAny", personnes.stream().findAny().get());

De methode [findAny] heeft de volgende handtekening:

 

De methode [findAny] kan elk element uit de stream weergeven. Tijdens het testen valt op dat bij sequentiële uitvoering het eerste element uit de stream wordt weergegeven, terwijl bij parallelle uitvoering inderdaad elk element kan worden weergegeven. Dit wordt geïllustreerd door de volgende code [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 {
        // lijst met personen
        List<Personne> personnes = Personnes.get();
        // alle personen
        affiche("all", personnes);
        // een willekeurige persoon
        affiche("findAny parallèle", personnes.stream().parallel().findAny().get());
        // een willekeurige persoon
        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));
    }
}
  • regel 22: findAny parallel uitgevoerd;
  • regel 24: findAny sequentieel uitgevoerd;

De verkregen resultaten zijn als volgt:

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"}
  • regel 4: de parallelle uitvoering heeft element 2 van de personenlijst opgeleverd. Dit had ook een ander element kunnen zijn;
  • regel 6: de sequentiële uitvoering heeft het eerste element van de lijst met personen geretourneerd;

Het gebruik van de methode [findAny] lijkt alleen zin te hebben bij de parallelle verwerking van een stroom.

5.6.3. [skip]


// de personen zonder de eerste
affiche("skip 1", personnes.stream().skip(1L).collect(Collectors.toList()));

De methode [skip] heeft de volgende signatuur:

 

De methode [skip] negeert de eerste n elementen van een stream. Zoals in de bovenstaande documentatie wordt aangegeven, levert de parallelle uitvoering van deze methode weinig prestatiewinst op en kan deze zelfs tot prestatieverlies leiden. Om de eerste n elementen te negeren, moeten de threads namelijk met elkaar afstemmen, waardoor de prestatiewinst die door de parallelliteit wordt behaald, teniet wordt gedaan.

De methode [skip] retourneert een stream Stream<Personne> die door de methode [collect] wordt omgezet naar het type List<Personne>. De signatuur van deze methode is als volgt:

 

De methode [collect] accepteert als parameter een instantie van het type [Collector], waarvan de signatuur complex is. Er bestaan vooraf gedefinieerde implementaties van het type [Collector], waardoor het meestal niet nodig is om deze zelf te implementeren. Hier wordt de implementatie [Collectors.toList()] gebruikt. [Collectors] is een klasse met talrijke statische methoden die het type [Collector<T,A,R>] implementeren. Dat is de eerste plaats waar je moet zoeken als je een Stream wilt omzetten in een standaard Java-collectie:

 

We zullen later enkele van deze methoden gebruiken.

De uitvoering levert het volgende resultaat op:

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

Het eerste element van de lijst (jean) is weggelaten.

5.6.4. [limit]


// de eerste twee personen
affiche("limit 2", personnes.stream().limit(2L).collect(Collectors.toList()));

De methode [limit] heeft de volgende handtekening:

 

Met de methode [limit] kunnen alleen de eerste n elementen van een stream worden behouden. Deze methode is niet geschikt voor parallelle verwerking.

De uitvoering levert het volgende resultaat op:

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

5.6.5. [count]


// het aantal personen
affiche("count", personnes.stream().count());

De methode [count] heeft de volgende signatuur:

 

De methode [count] geeft het aantal elementen van een Stream weer. Het parallel uitvoeren van de methode levert geen prestatiewinst op, zoals blijkt uit de volgende code (Voorbeeld06d1):


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;
        // aantal processors
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lijst met getallen
        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);
        // telling van de getallen – sequentiële methode
        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);
    }
}
  • regels 11-22: er wordt een Stream met 10 miljoen getallen aangemaakt;
  • regels 22-24: het tellen van de Stream;

De uitvoering levert het volgende resultaat op:

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

Als we regel 22 van de code vervangen door de volgende (Voorbeeld06d2):


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

krijgt men de volgende resultaten:

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]


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

De methode [max] heeft de volgende handtekening:

 

De methode [max] retourneert de maximale waarde van een stream met behulp van de vergelijker die als parameter wordt doorgegeven. Comparator is een functionele interface waarvan de enige te implementeren methode de volgende signatuur heeft: int compare (T o1, T o2). Deze methode moet -1 retourneren als o1 < o2, 0 als o1.equals(o2), en +1 als o1 > o2. De functionele interface Comparator heeft standaard talrijke statische methoden die de interface Comparator implementeren voor de meest voorkomende gevallen. Zo geldt in de instructie:


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

gebruiken we de statische methode [Comparator.comparingInt], waarvan de signatuur als volgt is:

 

Het type ToIntFunction is een functionele interface:

 

De methode [applyAsInt] van de functionele interface ToIntFunction genereert een type int op basis van een type T. Laten we teruggaan naar onze code:


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

De effectieve parameter van de methode [Comparator.comparingInt] moet hier een lambda van het type Personne --> int zijn. We geven de referentie door van de methode [Personne.getAge], die inderdaad deze signatuur heeft. Uiteindelijk krijgen we de persoon met de hoogste leeftijd. We krijgen een type Optional<Personne>, waaruit we de waarde extraheren met de methode [Optional.get]. We krijgen het volgende resultaat:

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

Het parallel berekenen van max levert geen prestatiewinst op, zoals het volgende voorbeeld laat zien: (Voorbeeld06e1):


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) {

        // gegevens
        // final long limiet = 100L;
        // final boolean verbose = true;
        final long limite = 10_000_000L;
        final boolean verbose = false;

        // aantal processors
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lijst met getallen
        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);
        // maximum van de getallen - sequentiële methode
        Stream<Long> sNombres = nombres.stream();
        Comparator<Long> compLong = (l1, l2) -> {
            if (verbose) {
                // thread
                System.out.printf("[%s]", Thread.currentThread().getName());
            }
            // vergelijking
            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();
        // maximale lengte = 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);
    }
}
  • regel 29: er is een stroom van limite willekeurige getallen van het type Long;
  • regels 30-47: de lambda-variabele compLong implementeert de interface Comparator<Long>. Deze interface wordt normaal gesproken geïmplementeerd door de methode [Comparator.naturalOrder()] in regel 49. Maar hier willen we de uitvoeringsthread weergeven (regels 31-33). Daarom implementeren we de interface zelf;
  • regel 50: zoeken naar max;

We krijgen de volgende resultaten:

 

Als we nu regel 27 vervangen door het volgende (Voorbeeld06e2):


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

krijgen we de volgende resultaten:

 

De parallelle uitvoering verliep dus trager. Als we met verbose=false naar 10 miljoen getallen gaan, krijgen we de volgende resultaten:

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

voor de sequentiële uitvoering:

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

voor de parallelle uitvoering, die dus nog steeds trager is.

De methode [Stream.min] wordt op dezelfde manier gebruikt:


// de lichtste persoon
affiche("poids min", personnes.stream().min(Comparator.comparingDouble(Personne::getPoids)).get());

5.6.7. [reduce]


// de totale leeftijd van alle personen
affiche("âge total (reduce)", personnes.stream().map(p -> p.getAge()).reduce(0, (a1, a2) -> a1 + a2));

De methode [reduce] is in paragraaf 5.3 beschreven. In regel 2 hierboven worden de leeftijden van alle personen bij elkaar opgeteld. Het resultaat is als volgt:

âge total (reduce) ----
60

5.6.8. [sorted]


// de personen in oplopende volgorde van leeftijd
affiche("personnes par âge croissant",
                personnes.stream().sorted(Comparator.comparingInt(Personne::getAge)).collect(Collectors.toList()));
// personen in alfabetische volgorde van de namen
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);

De methode [sorted] (regels 3 en 5) heeft de volgende handtekening:

 

De methode [sorted] accepteert als parameter het type [Comparator], zoals beschreven in paragraaf 5.6.6 voor de methoden min en max. Hiermee kan een Stream worden gesorteerd volgens de volgorde van de vergelijker die als parameter wordt doorgegeven. We hebben gezien dat de interface [Comparator] standaard verschillende statische methoden biedt die de gangbare vergelijkers implementeren, met name voor getallen en tekenreeksen. Hier gebruiken we de methode [Comparator.comparingInt], die als parameter een type ToIntFunction accepteert. Dit is een functionele interface van de methode [applyAsInt] met de volgende signatuur: int applyAsInt(T t). Hier is de daadwerkelijke parameter die in regel 3 aan de methode [Comparator.comparingInt] wordt doorgegeven, de referentie van de methode [Personne.age], die de leeftijd van de persoon opgeeft.

De interface [Comparator] biedt geen statische methoden om tekenreeksen te vergelijken. In regel 5 construeren we zelf een lambda die de enige methode van deze interface implementeert: int compare(T t1, T t2)


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

Deze lambda vergelijkt de namen van de personen. De verkregen resultaten zijn als volgt:

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"}]

Het lijkt niet mogelijk om het sorteren parallel uit te voeren, zoals blijkt uit de volgende code (Voorbeeld06f1):


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 {

        // gegevens
        final long limite = 100L;
        final boolean verbose = true;
//         eindlimiet = 10_000_000L;
//         final boolean verbose = false;

        // aantal processors
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lijst met getallen
        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);
        // sorteren van getallen - sequentiële methode
        Stream<Integer> sNombres = nombres.stream();
        début = new Date().getTime();
        Comparator<Integer> compInt = (i1, i2) -> {
            if (verbose) {
                // thread
                System.out.printf("[%s]", Thread.currentThread().getName());
            }
            // vergelijking
            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));
    }

}
  • regels 30-36: er wordt een stroom van limite willekeurige getallen gegenereerd;
  • regel 32: de lambda compInt (regels 38-55) wordt doorgegeven aan de methode [sorted]. Deze lambda sorteert de getallen in aflopende volgorde en geeft de thread weer die deze uitvoert.

De verkregen resultaten zijn als volgt:

 

Als we in de bovenstaande code regel 36 vervangen door de volgende (Voorbeeld06f2):


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

krijgen we de volgende resultaten:

 

We zien dat, verrassend genoeg, het sorteren van de reeks getallen met slechts één thread is uitgevoerd. Er was geen sprake van parallellisme. Of zie ik iets over het hoofd?

5.6.9. [anyMatch, noneMatch, allMatch]


// zijn er mensen ouder dan 100 jaar?
affiche("des personnes de + de 100 ans (anyMatch)", personnes.stream().anyMatch(p -> p.getAge() > 100));
// zijn alle mensen maximaal 100 jaar oud?
affiche("des personnes de + de 100 ans (noneMatch)", personnes.stream().noneMatch(p -> p.getAge() > 100));
// zijn alle mensen ouder dan 8 jaar?
affiche("des personnes de + de 8 ans (allMatch)", personnes.stream().allMatch(p -> p.getAge() > 8));

In de regels 2, 4 en 6 hebben de methoden [anyMatch, noneMatch, allMatch] als parameter een type Predicate, zoals beschreven in paragraaf 4.2. Ze voeren dus een filtering uit. Ze retourneren alle drie een booleaanse waarde:

  • anyMatch retourneert true als er ten minste één element van het type Stream is dat aan het filter voldoet;
  • noneMatch levert true op als er geen element in Stream is dat aan het filter voldoet;
  • allMatch levert true op als alle elementen van Stream aan het filter voldoen;

De verkregen resultaten zijn als volgt:

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)]


// De mensen worden ingedeeld naar geslacht
affiche("personnes regroupées par sexe", personnes.stream().collect(Collectors.groupingBy(p -> p.getSexe())));

De methode [collect] is beschreven in paragraaf 5.6.3. De parameter ervan is een implementatie van de interface [Collector]. De klasse [Collectors] biedt een aantal statische methoden die de interface [Collector] implementeren. Tot nu toe hebben we de methode [Collectors.toList()] gebruikt. Hier gebruiken we de statische methode [Collectors.groupingBy], die een woordenboek aanmaakt op basis van Stream. De signatuur ervan is als volgt:

 

De methode [groupingBy] creëert op basis van een type Stream<T> een type Map<K,List<T>>. De sleutel K wordt geleverd door de parameter van de methode [groupingBy] van het type Function<T,K>, waarvan de enige methode de volgende handtekening heeft: K apply(T t). Als we een woordenboek willen aanmaken dat geïndexeerd is op basis van het geslacht van personen, moeten we een functie opgeven die het geslacht genereert op basis van een persoon. Hier geven we als effectieve parameter van de methode [groupingBy] de referentie van de methode [Personne.getSexe] door. De verkregen resultaten zijn als volgt:

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"}]}

Op regel 2 staat de tekenreeks jSON uit een woordenboek dat is geïndexeerd op basis van twee sleutels: HOMME en FEMME.

De parallelle berekening levert geen prestatiewinst op, zoals blijkt uit het volgende voorbeeld (Voorbeeld06g1):


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 {

        // gegevens
        final long limite = 100L;
        final boolean verbose = true;
//         eindlimiet = 10_000_000L;
//         final boolean verbose = false;

        // aantal processors
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lijst met getallen
        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);
        // groepering van getallen per honderd - sequentiële methode
        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(getal -> getal / 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);
        // resultaten
        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));
    }

}
  • regels 23-38: opbouw van een stroom van limite-getallen;
  • regel 47: de getallen worden gegroepeerd per honderd. De lambda-functie in de regels 39-44 wordt gebruikt om de uitvoeringsthread weer te geven;

De uitvoerresultaten zijn als volgt:

 

Als we in de code regel 38 vervangen door de volgende regel (Voorbeeld06g2):


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

krijgen we de volgende resultaten:

 

We zien dat de parallelle uitvoering van het groeperen de prestaties heeft verslechterd.

5.6.11. [distinct]


// verwijderen van dubbele elementen uit een lijst
affiche("distinct", Stream.of(1, 2, 1).distinct().collect(Collectors.toList()));

De methode [distinct] heeft de volgende handtekening:

 

Hiermee kunnen duplicaten uit een feed worden verwijderd. De methode [Stream.of] (regel 2) heeft de volgende handtekening:

 

Hiermee kan een Stream worden aangemaakt op basis van expliciet opgegeven waarden. De resultaten van de uitvoering zijn als volgt:

distinct ----
[1,2]

5.6.12. [flatMap]


// van een Stream<Stream<T>>, wordt een Stream<T> gemaakt
affiche("flatMap", Stream.of(1, 2, 3).flatMap(i -> Stream.of(i, i + 10)).collect(Collectors.toList()));

De methode [flatMap] heeft de volgende handtekening:

 

De methode [flatMap] accepteert als parameter een functie die:

  • als parameter een element van het type T van Stream accepteert;
  • een Stream<R>-stream als resultaat oplevert;

Als in plaats van de methode [flatMap] de in paragraaf 5.5 beschreven methode [map] was gebruikt, zou het resultaat een type Stream<Stream<R>> zijn, waarbij elk element van het type T van de oorspronkelijke stroom een element van het type Stream<R> zou hebben opgeleverd. Levert de methode [flatMap] dan een type Stream<R> op? Deze methode voegt de verschillende stromen Stream<R> samen tot één enkele stroom. Dit blijkt uit de resultaten van de uitvoering van de voorgaande code:

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

Er bestaan gespecialiseerde varianten van [flatMap]:


// van een Stream<IntStream> wordt een IntStream gemaakt, waarvan de som wordt berekend
affiche("flatMapToInt", Stream.of(1, 2, 3).flatMapToInt(i -> IntStream.of(i, i + 10)).sum());

De methode [flatMapToInt] heeft de volgende handtekening:

 

De methode [flatMapToInt] accepteert als parameter een functie die als resultaat het volgende type IntStream oplevert:

 

IntStream is een feed van int. Dit type verdient de voorkeur boven het type Stream<Integer>, omdat bij de verwerking ervan het 'boxing' en 'unboxing' tussen de typen Integer en int wordt vermeden. Deze interface neemt veel methoden van het type Stream<T> over en voegt er andere aan toe, waaronder de hierboven genoemde methode [sum], die de elementen van IntStream bij elkaar optelt.

De volgende code illustreert het gebruik van de analoge methode [flatMapToDouble]:


// van een Stream<DoubleStream> wordt een DoubleStream gemaakt en vervolgens een array
affiche("flatMapToDouble", Stream.of(1, 2, 3).flatMapToDouble(i -> DoubleStream.of(i, i * 1.2)).toArray());

Met de methode [DoubleStream.toArray] kunt u van een type DoubleStream naar een type double[] overschakelen.

De resultaten voor deze twee voorbeelden zijn als volgt:

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

Het volgende voorbeeld toont de prestatiewinst die wordt behaald door over te schakelen van het type Stream<Long> naar het type LongStream (Voorbeeld06i1):


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;
        // aantal processors
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
        // lijst met getallen
        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);
        // som van de getallen – sequentiële methode
        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);
    }
}
  • regel 22: berekening van de som van een reeks getallen van het type Long;

Dit levert de volgende resultaten op:

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

Laten we nu regel 22 vervangen door de volgende (Voorbeeld06i2):


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

Met de methode Stream<Integer>.mapToLong kunnen we een stroom van het type LongStream verkrijgen, bestaande uit elementen van het primitieve type long, die we vervolgens optellen met de functie sum. Dit levert de volgende resultaten op:

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

De prestatiewinst is duidelijk merkbaar.

5.6.13. methoden voor de verwerking van primitieve getallen


// maximum van een reeks integers
affiche("IntStream max", IntStream.of(1, 10, 8).max());
// minimum van een reeks double
affiche("DoubleStream min", DoubleStream.of(1.5, 10.4, 8.9).min());
// gemiddelde van een stroom van int
affiche("IntStream average", IntStream.of(1, 10, 8).average().getAsDouble());
// statistieken van een stroom van int
affiche("IntStream summaryStatistics", IntStream.of(1, 10, 8).summaryStatistics());

De gegevensstromen van primitieve typen (int, long, double) bieden methoden die op deze typen zijn afgestemd. Het resultaat van de uitvoering van de voorgaande code is als volgt:

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}
  • Het resultaat van regel 2 van de code is een type OptionalInt, vergelijkbaar met het type Optional<Integer>. De waarde die in dit object is opgeslagen, kan worden opgehaald met de methode [getAsInt()]. De aanwezigheid van een waarde kan worden getest met de methode [isPresent()]. Regel 2 van de resultaten betekent niet dat de klasse [OptionalInt] velden heeft met de naam [asInt, present]. Standaard gebruikt de bibliotheek jSON alle openbare methoden getX en isY van het object dat moet worden geserialiseerd naar jSON. En hier is er inderdaad een methode [getAsInt] en nog een methode [isPresent], terwijl de velden [asInt, present] zelf niet bestaan;
  • het resultaat van regel 4 van de code is een type OptionalDouble, vergelijkbaar met het type Optional<Double>;
  • het resultaat van regel 6 van de code is een type OptionalDouble waarvan de waarde kan worden verkregen met de methode [getAsDouble()]. De methode [average] berekent het gemiddelde van de reeks getallen;
  • het resultaat van regel 8 van de code is een type IntSummaryStatistics dat als volgt is gedefinieerd:
 

We zien dat het verkregen object IntSummaryStatistics verschillende gegevens over de reeks getallen weergeeft, zoals het aantal waarden, de som, het maximum, het minimum en het gemiddelde.