Skip to content

7. De bibliotheek RxJava

De bibliotheek RxJava is gebaseerd op het volgende concept: een stroom van elementen van het type T Observable<T> wordt geobserveerd door een of meer abonnees (Subscriber<T>). De bibliotheek RxJava maakt het mogelijk dat de Observable<T>-stroom wordt uitgevoerd in een thread T1 en de bijbehorende waarnemer Subscriber<T> in een thread T2, zonder dat de ontwikkelaarzich zorgen hoeft te maken over het beheer van de levenscyclus van deze threads en over van nature lastige problemen, zoals het delen van gegevens tussen threads en de synchronisatie daarvan om een globale taak uit te voeren. Het vergemakkelijkt dus asynchroon programmeren.

Een Observable<T>-stream produceert elementen van het type T, die kunnen worden geobserveerd zodra ze worden geproduceerd. Als de waarnemer en de observable (een verkeerde benaming voor het type Observable<T>) zich in dezelfde thread bevinden, dan kan de observable het element (i+1) pas produceren wanneer de waarnemer het element i heeft verwerkt. Er zijn maar weinig gevallen waarin deze architectuur zinvol is. Als de waarnemer en het observeerbare object zich niet in dezelfde thread bevinden, dan gedragen het observeerbare object en zijn waarnemer zich autonoom: het observeerbare object produceert in zijn eigen tempo en de waarnemer verwerkt in zijn eigen tempo. Daarin ligt het nut van de bibliotheek. Tot nu toe hebben we het altijd over één waarnemer gehad. In werkelijkheid kan een observeerbaar object een willekeurig aantal waarnemers hebben.

De bibliotheek RxJava is bijzonder goed geschikt voor de architectuur die in paragraaf 2 van de inleiding is besproken en die we hier nogmaals weergeven:

Image

  • in [1] levert een servicelaag diensten, waarvan sommige tijdrovend zijn (bijvoorbeeld netwerkverzoeken);
  • deze servicelaag wordt aangeroepen door een grafische interface [1] (Swing, Android, JavaFx). Als de servicelaag in dezelfde thread wordt uitgevoerd als de methode [swing] die er gebruik van maakt, loopt de grafische interface vast (reageert niet) terwijl er op het resultaat van de service wordt gewacht;
  • in [2] maakt een dunne adaptatielaag, geïmplementeerd met RxJava, het mogelijk om aan de grafische laag een asynchrone implementatie van dezelfde service te presenteren: deze kan worden uitgevoerd in een andere thread dan die van de methode van de grafische laag die deze aanroept. In dit geval blijft de grafische interface [3] responsief: de gebruiker kan ermee blijven werken, bijvoorbeeld door parallel aan de eerste een nieuwe netwerkverzoek te starten, en bovenal kan men hem de mogelijkheid bieden om te lang durende bewerkingen te annuleren, iets wat onmogelijk is als de grafische interface vastloopt;
  • de aanroep [4] is synchroon, terwijl de aanroep [5-6] asynchroon is;

In deze architectuur biedt de laag [2] diensten aan die **Observable&lt;T&gt;**-typen retourneren, waarop de methoden van de grafische laag [3] zich kunnen abonneren. Een service van de laag [2] levert vervolgens zijn resultaten één voor één af en de laag [3] kan op elk resultaat reageren, bijvoorbeeld door een of meer componenten van de grafische interface bij te werken.

De klasse Observable<T> beschikt over tientallen methoden. Dit is een van de uitdagingen van de bibliotheek: ze is zeer uitgebreid en het is moeilijk om alle mogelijkheden te doorgronden. We zullen er enkele bespreken. De beheersing van de overige methoden komt dan vanzelf met de tijd.

7.1. Observables aanmaken en je erop abonneren

7.1.1. Voorbeeld-01: de methode [Observable.from]

  

Laten we de volgende code eens bekijken:


package dvp.rxjava.observables;

import rx.Observable;
import rx.functions.Action0;
import rx.functions.Action1;

import java.util.Arrays;

public class Exemple01 {
  public static void main(String[] args) {
    // waarneembare gehele getallen
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    obs1.subscribe(new Action1<Integer>() {
      @Override
      public void call(Integer integer) {
        System.out.printf("next : %s%n", integer);
      }
    }, new Action1<Throwable>() {
      @Override
      public void call(Throwable throwable) {
        System.out.println(throwable);
      }
    }, new Action0() {
      @Override
      public void call() {
        System.out.println("completed");
      }
    });
  }
}
  • regel 12: er wordt een type Observable<Integer> aangemaakt op basis van een lijst met gehele getallen.

De klasse Observable<T> is een stroom van elementen van het type T die kunnen worden geobserveerd, bij voorkeur asynchroon maar niet noodzakelijkerwijs, naarmate ze worden geproduceerd. De definitie ervan is als volgt:

 

Zoals reeds vermeld, beschikt de klasse Observable<T> over tientallen methoden. Sommige daarvan zijn vergelijkbaar met die van de klasse Stream<T> die in paragraaf 5 is besproken. De documentatie van RxJava bevat 'marble diagrams' [2] die de werking van deze methoden illustreren:

  • regel 3 illustreert de emissies van de observabele in de loop van de tijd;
  • de methode [4] wordt toegepast op de elementen die door de observabele worden uitgezonden. Dit levert doorgaans een nieuwe observabele op;
  • regel 5 toont de verkregen nieuwe observabele;

De methode [Observable.from] heeft de volgende signatuur:

 

Met de statische methode [Observable.from] kan een Observable<T> worden aangemaakt op basis van een verzameling elementen van het type T. Dit is een zeer eenvoudige manier om aan de slag te gaan met observables. De regel:


    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));

zal dus drie elementen verzenden. Deze worden niet onmiddellijk verzonden. Ze worden pas volledig verzonden telkens wanneer zich een waarnemer aanmeldt. Dit wordt een ‘koude’ observable genoemd. De observable verzendt zijn elementen opnieuw voor elke nieuwe abonnee.

We kunnen de vorige instructie beschouwen als een configuratieactie van de observable. Deze wordt één keer geconfigureerd en n keer uitgevoerd als er n abonnees zijn.

Hoe schrijf je je in?

Een manier om dit te doen is door gebruik te maken van de methode [Observable.subscribe], waarvan de hier gebruikte definitie als volgt luidt:

 
  • de eerste parameter [Action1<T> onNext] (zie paragraaf 6.2) van de methode is de methode die moet worden uitgevoerd wanneer de observable een nieuw T-element uitzendt;
  • de tweede parameter [Action1<Throwable> onError] van de methode is de methode die moet worden uitgevoerd wanneer de observable een uitzondering genereert;
  • de derde parameter [Action0 onComplete] (zie paragraaf 6.1) van de methode is de methode die moet worden uitgevoerd wanneer de observable een uitzondering genereert;
  • de methode retourneert een type [Subscription];

Het type [Subscription] vertegenwoordigt een abonnement op de observable. De definitie ervan is als volgt:

 

Het voordeel van deze interface [1] ligt in de methode [2], waarmee een abonnement kan worden opgezegd.

In ons voorbeeld is de code van het abonnement op de observable als volgt:


    obs1.subscribe(new Action1<Integer>() {
      @Override
      public void call(Integer integer) {
        System.out.printf("next : %s%n", integer);
      }
    }, new Action1<Throwable>() {
      @Override
      public void call(Throwable throwable) {
        System.out.println(throwable);
      }
    }, new Action0() {
      @Override
      public void call() {
        System.out.println("completed");
      }
});
  • regel 1: het resultaat van het type [Subscription] wordt genegeerd;
  • regels 1-15: de drie parameters zijn instanties van anonieme klassen. We zullen ook lambda's gebruiken. Het voordeel van anonieme klassen is dat duidelijk te zien is welke gegevenstypen de enige methode van deze klassen verwacht;
  • regels 2-5: implementatie van de eerste parameter van het type [Action1<Integer>];
  • regels 6-10: implementatie van de tweede parameter van het type [Action1<Throwable>];
  • regels 11-15: implementatie van de derde parameter van het type [Action0];

De volledige code is als volgt:


package dvp.rxjava.observables;

import rx.Observable;
import rx.functions.Action0;
import rx.functions.Action1;

import java.util.Arrays;

public class Exemple01 {
  public static void main(String[] args) {
    // waarnemingen van gehele getallen
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    // abonnement
    obs1.subscribe(new Action1<Integer>() {
      @Override
      public void call(Integer integer) {
        System.out.printf("next : %s%n", integer);
      }
    }, new Action1<Throwable>() {
      @Override
      public void call(Throwable throwable) {
        System.out.println(throwable);
      }
    }, new Action0() {
      @Override
      public void call() {
        System.out.println("completed");
      }
    });
  }
}

De observable in regel 12 begint zijn 3 elementen uit te zenden zodra de methode [subscribe] in regel 14 wordt aangeroepen. Vanaf dat moment:

  • worden bij elk uitgezonden element de regels 15-18 uitgevoerd.
  • aan het einde van de 3 elementen worden de regels 24-29 uitgevoerd;
  • de regels 19-24 worden nooit uitgevoerd omdat de observable hier geen uitzondering genereert;

Standaard worden de observable en de observer in dezelfde thread uitgevoerd. Er zijn enkele vooraf gedefinieerde observables die in een andere thread dan de hoofdthread (hier de thread van de methode main) worden uitgevoerd, maar voor de meeste is dit niet het geval. Hier gebeurt alles dus in de thread van de methode [main]:

  • de observable zendt element 1 uit;
  • de regels 15-18 worden uitgevoerd en geven dit element weer;
  • de observable zendt element 2 uit;
  • de regels 15-18 worden uitgevoerd en geven dit element weer;
  • de observable zendt element 3 uit;
  • de regels 15-18 worden uitgevoerd en geven dit element weer;
  • de observable verzendt de melding [completed];
  • de regels 24-29 worden uitgevoerd;

Dit blijkt uit de verkregen resultaten:

1
2
3
4
next : 1
next : 2
next : 3
completed

De klasse [Exemple02] is een afgeleide van [Exemple01], waarbij ditmaal lambda-functies worden gebruikt als parameters voor de methode [Observable.subscribe]:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;

public class Exemple02 {
  public static void main(String[] args) {
    // waarneembare grootheden van gehele getallen
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    // abonnement
    obs1.subscribe(
      (integer) -> System.out.printf("next : %s%n", integer),
      (th) -> System.out.println(th),
      () -> System.out.println("completed"));
  }
}

7.1.2. Voorbeeld-03: de klasse Observer

  

De methode [Observable.subscribe], waarmee men zich kan abonneren op een observable, kent verschillende versies, waaronder de volgende:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;
import rx.Observer;

public class Exemple03 {
    public static void main(String[] args) {
        // waarneembare gehele getallen
        Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
        // abonnement
        obs1.subscribe(new Observer<Integer>() {
            @Override
            public void onCompleted() {
                System.out.println("completed");
            }

            @Override
            public void onError(Throwable th) {
                System.out.printf("throwable %s", th);
            }

            @Override
            public void onNext(Integer integer) {
                System.out.printf("next : %s%n", integer);
            }
        });
    };
}

Regel 13: in plaats van drie parameters door te geven aan de methode [subscribe], wordt het volgende type [Observer] doorgegeven:

 

Het type [Observer] is een interface met drie methoden:

  • [onNext(T t)], die wordt aangeroepen telkens wanneer de observable een element t uitzendt;
  • [onError(Throwable th)], die wordt aangeroepen wanneer de observable een uitzondering th uitzendt;
  • [onCompleted], die wordt aangeroepen wanneer de observable aangeeft dat het is klaar met verzenden;

De werking van de code is vergelijkbaar met wat eerder is uitgelegd. Dit levert de volgende resultaten op:

1
2
3
4
next : 1
next : 2
next : 3
completed

7.1.3. Voorbeeld-04: de methode [Observable.create]

  

De statische methode Observable.create is als volgt gedefinieerd:

 
  • de methode [create] retourneert een type Observable<T>;
  • de parameter van de methode [create] is een functie van het type [Observable.OnSubscribe<T>], die als volgt is gedefinieerd:
 

Het type [Observable.OnSubscribe<T>] is een functionele interface die zelf de functionele interface [Action1<Subscriber<? super T>>] uitbreidt. De methode [call] van deze interface verwacht een type [Subscriber] (abonnee, inschrijver, waarnemer) dat als volgt is gedefinieerd:

 

In [1] zien we dat de klasse [Subscriber<T>] de in paragraaf 7.1.2 beschreven interface [Observer<T>] implementeert.

Uiteindelijk verwacht de methode [<T> Observable.create]:

  • een instantie van het type [Observable.OnSubscribe<T>] als parameter, met als enige methode: void call(Subscriber<T> s). Het type [Subscriber<T>] is een uitbreiding van het type [Observer<T>] en beschikt daarom over de methoden onNext, onError, onCompleted;
  • retourneert een type Observable<T>;

De methode [<T> Observable.create] retourneert een geconfigureerde observable. Er zijn nog geen elementen uitgezonden. Wanneer een abonnee [Subscriber<T> s] zich op deze observable abonneert, wordt de methode [void call(s)] van de functie die als parameter aan de methode [<T> Observable.create] is doorgegeven, aangeroepen. Deze methode heeft als taak elementen t van het type T te verzenden en bij elke verzending de methode [s.onNext(t)] van de observer aan te roepen. Zodra deze is voltooid, moet de methode [s.onCompleted(t)] van de waarnemer worden aangeroepen en moet de methode [call] worden beëindigd. Als de methode [call] een uitzondering th tegenkomt, moet de methode [s.onError(th)] van de waarnemer worden aangeroepen en moet de methode [call] worden beëindigd;

Om deze complexe werking te illustreren, gebruiken we de volgende code [Exemple04]:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;

import java.util.Random;

public class Exemple04 {
    public static void main(String[] args) {
        // configuratie van reële waarden
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                for (int i = 0; i < 3; i++) {
                    // uitzending element i
                    subscriber.onNext(new Random((i + 1)).nextDouble());
                }
                // einde uitzending
                subscriber.onCompleted();
            }
        });
        // abonnement en dus uitzending
        obs1.subscribe((d) -> System.out.printf("onNext %s%n", d), (th) -> System.out.printf("onError %s%n", th),
                () -> System.out.println("onCompleted"));
    }
}
  • regel 11: er wordt een observable aangemaakt die Double-typen uitzendt;
  • regels 11-21: de parameter van de methode [create] wordt geïnstantieerd met een anonieme klasse die de enige methode [call] uit de regels 12-20 bevat. De in regel 11 aangemaakte observable is klaar om te verzenden, maar zal pas verzenden wanneer er een observer arriveert;
  • regels 13-21: de methode [call] ontvangt de referentie van een waarnemer;
  • regels 14-17: verzending van 3 elementen naar de waarnemer;
  • regel 19: melding van het einde van de verzending aan de waarnemer;
  • regels 23-24: abonneren op de observable van regel 11. De drie parameters [onNext, onError, onCompleted] van de methode [subscribe] worden geïmplementeerd door middel van drie lambda's. Dit abonnement creëert de abonnee [Subscriber<Double>], die wordt doorgegeven aan de methode [call] in regel 13. Het verzenden van elementen begint dan;
  • alles vindt plaats in dezelfde thread: observable en observer;

We krijgen de volgende resultaten:

1
2
3
4
onNext 0.7308781907032909
onNext 0.7311469360199058
onNext 0.731057369148862
onCompleted

Met de methode [Observable.create] kan een observable worden aangemaakt op basis van elk willekeurig fenomeen. Dit is de methode die we in paragraaf 2 van de verkenning hebben gebruikt om een synchrone interface om te zetten in een asynchrone interface.

7.1.4. Voorbeeld-05: refactoring van [Exemple-04]

  

Het volgende voorbeeld toont een nieuwe versie van de statische methode [Observable.subscribe]:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Random;

public class Exemple05 {
    public static void main(String[] args) {
        // configuratie van een observabele reële variabele
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                showInfos("Observable.call start");
                for (int i = 0; i < 3; i++) {
                    // in afwachting
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // fout
                        subscriber.onError(e);
                    }
                    // actie
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // voltooid
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // een abonnee
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
            }

            @Override
            public void onError(Throwable e) {
                showInfos(String.format("Subscriber.onError (%s)", e));
            }

            @Override
            public void onNext(Double aDouble) {
                showInfos(String.format("Subscriber.onNext (%s)", aDouble));
            }
        };

        // inschrijving
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        showInfos("après souscription");

    }

    private static void showInfos(String message) {
        System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message, Thread.currentThread().getName(),
                new SimpleDateFormat("ss:SSS").format(new Date()));
    }
}
  • regel 56: de nieuwe versie van de statische methode [Observable.subscribe] accepteert als parameter het type [Subscriber] dat we in de vorige paragraaf hebben besproken;
  • regels 37-52: de abonnee (observer). Deze implementeert de interface Observer met zijn drie methoden onNext, onError en onCompleted;
  • regels 61-64: vanaf nu gaan we ons richten op de threads waarin de observable en de observer worden uitgevoerd;
  • regel 62: de naam van de thread;
  • regel 63: de huidige tijd uitgedrukt in seconden en milliseconden. Hierdoor kunnen we in de tijd zien wanneer de observable elementen uitzendt en hoe deze door de observer worden verwerkt;
  • deze code heeft dezelfde functionaliteit als de vorige code. We hebben de vorige code alleen geherstructureerd;

De verkregen resultaten zijn als volgt:

avant souscription ------Thread[main] ---- Time[31:685]
Observable.call start ------Thread[main] ---- Time[31:691]
Observable.call onNext(80.39999999999999) ------Thread[main] ---- Time[32:194]
Subscriber.onNext (80.39999999999999) ------Thread[main] ---- Time[32:195]
Observable.call onNext(73.2) ------Thread[main] ---- Time[32:595]
Subscriber.onNext (73.2) ------Thread[main] ---- Time[32:595]
Observable.call onNext(106.8) ------Thread[main] ---- Time[32:897]
Subscriber.onNext (106.8) ------Thread[main] ---- Time[32:897]
Observable.call onCompleted ------Thread[main] ---- Time[32:898]
Subscriber.onCompleted ------Thread[main] ---- Time[32:898]
après souscription ------Thread[main] ---- Time[32:899]
  • regel 1 van de resultaten: vóór regel 56 van de code is er nog niets gebeurd. De observable is alleen geconfigureerd;
  • regel 2 van de resultaten: regel 56 van de code roept de methode [call] van regel 15 aan. Regel 3: het reële getal 80,39 wordt naar de waarnemer verzonden;
  • regel 4: de waarnemer ontvangt het verzonden getal;
  • regels 5-8: het voorgaande proces herhaalt zich twee keer;
  • regel 9: de observable verstuurt de melding dat de verzending is voltooid;
  • regel 10: de waarnemer ontvangt deze;
  • regel 11: weergegeven door regel 57 van de code;

We zien dus dat alleen regel 56, de inschrijving, ervoor heeft gezorgd dat de regels 2-10 van de resultaten werden weergegeven. Wanneer men begint met de bibliotheek RxJava, vraagt men zich af hoe de zaken in elkaar grijpen en met name welke verbanden er bestaan tussen de waarnemer en de waarneembare. We zien hier dat regel 57, de inschrijving op de observeerbare,

  • ervoor heeft gezorgd dat alle elementen van de observeerbare zijn verzonden;
  • dat de observeerbare en de waarnemer in dezelfde thread worden uitgevoerd;
  • dat we daardoor de volgende reeks waarnemen: verzending element i, waarneming element i, verzending element (i+1), waarneming element (i+1), ...

We herinneren ons dat de zender wachtte voordat hij zijn elementen uitzond:


                    // in afwachting
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // fout
                        subscriber.onError(e);
}

waarbij i in regel 3 het verzendnummer vertegenwoordigt (0<=i<3). Als we de verzendtijden van de elementen van de observeerbare bekijken:

  • regels 2, 3: element 0 werd ongeveer 500 ms na het begin van het abonnement verzonden;
  • regels 3, 5: element 1 werd ongeveer 400 ms na element 0 verzonden;
  • regels 5, 7: element 2 is ongeveer 300 ms na element 1 verzonden;

7.2. Uitvoeringsthread, observatiethread

7.2.1. Voorbeeld-06: observabel en waarnemer in een andere thread dan [main]

  

We herstructureren het vorige voorbeeld als volgt [Exemple06]:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;
import rx.schedulers.Schedulers;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Random;
import java.util.concurrent.CountDownLatch;

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

        // barrièrebewaker
        CountDownLatch latch = new CountDownLatch(1);

        // configuratie van een reële observabele
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                showInfos("Observable.call start");
                for (int i = 0; i < 3; i++) {
                    // in afwachting
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // fout
                        subscriber.onError(e);
                    }
                    // actie
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // voltooid
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // een abonnee
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // de slagboom gaat omlaag
                latch.countDown();
            }

            @Override
            public void onError(Throwable e) {
                showInfos(String.format("Subscriber.onError (%s)", e));
            }

            @Override
            public void onNext(Double aDouble) {
                showInfos(String.format("Subscriber.onNext (%s)", aDouble));
            }
        };

        // vervolg configuratie waarneembaar
        obs1 = obs1.subscribeOn(Schedulers.computation());
        // inschrijving
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // wachten voor de slagboom
        try {
            showInfos("début attente barrière");
            latch.await();
            showInfos("fin attente barrière");
        } catch (InterruptedException e1) {
            System.out.println(e1);
        }
        showInfos("après souscription");

    }

    private static void showInfos(String message) {
        System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message, Thread.currentThread().getName(),
                new SimpleDateFormat("ss:SSS").format(new Date()));
    }
}
  • regel 16: we maken een barrière (semafoor) aan met een object van het type [CountDownLatch]. Dit object dient om threads onderling te synchroniseren. Het wordt hier geïnitialiseerd met de waarde 1, die we de waarde van de barrière (of semafoor) zullen noemen. Een thread wacht op de barrière door middel van een bewerking:

latch.await();

De thread wordt geblokkeerd als de waarde van de barrière >0 is. Een thread kan de interne waarde van de barrière verhogen of verlagen. In regel 48 wordt de waarde van de barrière met 1 verlaagd.

  • regel 63: de observable is zo geconfigureerd dat deze wordt uitgevoerd op een thread die wordt geleverd door de scheduler [Schedulers.computation()]. Deze scheduler kan evenveel threads leveren als er cores op de uitvoeringsmachine zijn. In de paragraaf over de voorbeeldtoepassing is het gebruik van andere schedulers getoond (zie paragraaf 2.8);

Het principe van de code is als volgt:

  • de methode [main] wordt uitgevoerd in de hoofdthread (main);
  • regel 66: start het verzenden van elementen van de observable. Deze worden verzonden via een andere thread dan de hoofdthread;
  • regel 70: de hoofdthread wordt geblokkeerd omdat de barrière de waarde 1 heeft (zie regel 16). De thread kan pas verdergaan wanneer deze waarde op 0 komt te staan. Dit gebeurt op regel 48. Het is de waarnemer die de barrière laat zakken wanneer hij de melding ontvangt dat de observable klaar is met verzenden;

De uitvoering levert de volgende resultaten op:

avant souscription ------Thread[main] ---- Time[09:268]
Observable.call start ------Thread[RxComputationThreadPool-1] ---- Time[09:278]
début attente barrière ------Thread[main] ---- Time[09:278]
Observable.call onNext(44.4) ------Thread[RxComputationThreadPool-1] ---- Time[09:783]
Subscriber.onNext (44.4) ------Thread[RxComputationThreadPool-1] ---- Time[09:783]
Observable.call onNext(18.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:183]
Subscriber.onNext (18.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:184]
Observable.call onNext(54.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:486]
Subscriber.onNext (54.0) ------Thread[RxComputationThreadPool-1] ---- Time[10:488]
Observable.call onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[10:489]
Subscriber.onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[10:490]
fin attente barrière ------Thread[main] ---- Time[10:491]
après souscription ------Thread[main] ---- Time[10:493]
  • regel 1: de inschrijving gaat plaatsvinden;
  • regel 2: dit activeert de uitvoering van de methode [call] op de thread [RxComputationThreadPool-1]. We hebben nu een parallelle uitvoering met twee threads;
  • regel 3: om een onduidelijke reden heeft de thread [RxComputationThreadPool-1] het stokje doorgegeven. De thread [main] neemt het vervolgens over en wordt geblokkeerd door de guardrail (regel 70 van de code). Vanaf dat moment kan alleen de thread [RxComputationThreadPool-1] nog opereren;
  • regels 4-11: we zien het eerder waargenomen gedrag tussen de observabele en zijn waarnemer, maar alles speelt zich nu af in de thread [RxComputationThreadPool-1];
  • regels 12-13: de waarnemer heeft de slagboom neergelaten (regel 48 van de code) en de thread [RxComputationThreadPool-1] is beëindigd. De thread [main] neemt het over en geeft twee berichten weer;

7.2.2. Voorbeeld-07: observable en observer in twee verschillende threads

  

We passen het vorige voorbeeld als volgt aan:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;
import rx.schedulers.Schedulers;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Random;
import java.util.concurrent.CountDownLatch;

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

        // slagboomwachter
        CountDownLatch latch = new CountDownLatch(1);

        // configuratie van een observabele van reële getallen
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                showInfos("Observable.call start");
                for (int i = 0; i < 3; i++) {
                    // wachten
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // fout
                        subscriber.onError(e);
                    }
                    // actie
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // voltooid
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // een abonnee
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // de slagboom gaat omlaag
                latch.countDown();
            }

            @Override
            public void onError(Throwable e) {
                showInfos(String.format("Subscriber.onError (%s)", e));
            }

            @Override
            public void onNext(Double aDouble) {
                showInfos(String.format("Subscriber.onNext (%s)", aDouble));
            }
        };

        // vervolg configuratie waarneembaar
        obs1 = obs1.subscribeOn(Schedulers.computation()).observeOn(Schedulers.computation());
        // inschrijving
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // in afwachting van het openen van de slagboom
        try {
            showInfos("début attente barrière");
            latch.await();
            showInfos("fin attente barrière");
        } catch (InterruptedException e1) {
            System.out.println(e1);
        }
        showInfos("après souscription");

    }

    private static void showInfos(String message) {
        System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message, Thread.currentThread().getName(),
                new SimpleDateFormat("ss:SSS").format(new Date()));
    }
}

De code is identiek aan die van het vorige voorbeeld, behalve regel 63:


obs1 = obs1.subscribeOn(Schedulers.computation()).observeOn(Schedulers.computation());

waarmee de observable (subscribeOn) en de observer (observeOn) worden geconfigureerd om te worden uitgevoerd op een van de threads die worden geleverd door de scheduler [Schedulers.computation()].

De verkregen resultaten zijn als volgt:

avant souscription ------Thread[main] ---- Time[09:643]
début attente barrière ------Thread[main] ---- Time[09:656]
Observable.call start ------Thread[RxComputationThreadPool-4] ---- Time[09:656]
Observable.call onNext(39.6) ------Thread[RxComputationThreadPool-4] ---- Time[10:162]
Subscriber.onNext (39.6) ------Thread[RxComputationThreadPool-3] ---- Time[10:163]
Observable.call onNext(98.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[10:562]
Subscriber.onNext (98.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[10:564]
Observable.call onNext(46.8) ------Thread[RxComputationThreadPool-4] ---- Time[10:864]
Observable.call onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[10:866]
Subscriber.onNext (46.8) ------Thread[RxComputationThreadPool-3] ---- Time[10:866]
Subscriber.onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[10:868]
fin attente barrière ------Thread[main] ---- Time[10:869]
après souscription ------Thread[main] ---- Time[10:870]

De volgende punten vallen op:

  • de observable wordt uitgevoerd in de thread [RxComputationThreadPool-4] (regels 3-4, 6, 8-9);
  • de observer wordt uitgevoerd in de thread [RxComputationThreadPool-3] (regels 5, 7, 10-11);
  • dat ze onafhankelijk van elkaar worden uitgevoerd. Zo verzendt de observable in de regels 8-9 twee meldingen (onNext, onCompleted) voordat de observer de melding [onNext] ophaalt (regel 10);

De bibliotheek RxJava zorgt voor de gegevensoverdracht (de verzendingen) van de thread van de observable naar de thread van de observer. De ontwikkelaar hoeft zich daar geen zorgen over te maken.

We hebben gezien hoe we observables kunnen aanmaken (Observable.from, Observable.create). Nu bekijken we de vooraf gedefinieerde observables van de bibliotheek RxJava.

7.3. Vooraf gedefinieerde observables

7.3.1. Voorbeeld-08: de methode [Observable.range]

 

Vanaf nu gaan we speciale klassen gebruiken voor de geobserveerde processen en hun waarnemers. Het idee is om hun naam, uitvoeringsthread en uitvoeringstijden te kunnen loggen, zodat we deze in de loop van de tijd kunnen volgen.

De klasse [Process] is simpelweg een Observable waaraan een naam kan worden toegekend. Deze klasse implementeert de volgende interface [IProcess]:


package dvp.rxjava.observables.utils;

import rx.Observable;

public interface IProcess<T> {

    // naam van de observabele
    public String getName();

    // waarneembare grootheid
    public Observable<T> getObservable();

}

Deze interface kan worden geïmplementeerd door de volgende klasse [Process<T>]:


package dvp.rxjava.observables.utils;

import rx.Observable;
import rx.Scheduler;

public class Process<T> implements IProcess<T>{

    // naam van de observabele
    protected String name;
    // geobserveerd proces
    protected Observable<T> observable;

    // constructors
    public Process(String name, Observable<T> observable) {
        // lokale initialisaties
        this.name = name;
        this.observable = observable;
    }

    // getters en setters
    public String getName() {
        return name;
    }

    public Observable<T> getObservable() {
        return observable;
    }

}
  • regel 9: de naam van het proces;
  • regel 11: de waargenomen grootheid;
  • regels 14-18: de constructor;

De waarnemer wordt beschreven door de volgende klasse [Observateur]:


package dvp.rxjava.observables.utils;

import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

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

import rx.Subscriber;

public class Observateur<T> extends Subscriber<T> {

...
}
  • regel 11: de klasse Observateur<T> is een uitbreiding van de klasse Subscriber<T> die we kort hebben besproken in paragraaf 7.1.3. We zullen deze gebruiken als argument voor de methode [Observable.subscribe]:

// waarneembare uitvoering (waarneming)
obs1.subscribe(observateur);

De methode [Observable.subscribe] die hierboven in regel 2 wordt gebruikt, heeft de volgende definitie:

 

De rol van [Subscriber] bestaat voornamelijk uit het beheren van de elementen die worden verzonden door de observable waarop het is geabonneerd, met behulp van de methoden van de interface [Observer]: onNext, onError, onCompleted. De klasse [Subscriber] beschikt over de volgende methoden:

 

In de code van de klasse [Observateur] zullen we de methode [1] isUnsubscribed gebruiken om te controleren of het abonnement van de abonnee al dan niet is opgezegd. De volledige klasse [Observateur<T>] is als volgt:


package dvp.rxjava.observables.utils;

import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

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

import rx.Subscriber;

public class Observateur<T> extends Subscriber<T> {

    // een semafoor
    private CountDownLatch latch;
    // een weergavemethode
    private Consumer<String> showInfos;
    // de naam van de waarnemer
    private String observerName;
    // de naam van het geobserveerde proces
    private String processName;

    // constructors
    public Observateur() {

    }

    public Observateur(String name, CountDownLatch latch, Consumer<String> showInfos, String observedName) {
        this.observerName = name;
        this.latch = latch;
        this.showInfos = showInfos;
        this.processName = observedName;
    }

    // --------------------------- implementatie van de interface Observer<T>
    @Override
    public void onCompleted() {
        // einde van de uitzendingen
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber [%s,%s].onCompleted", observerName, processName));
        }
        // einde blokkering hoofdthread
        latch.countDown();
    }

    @Override
    public void onError(Throwable e) {
        // uitzendfout
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber [%s, %s].onError (%s)", observerName, processName, e));
        }
    }

    @Override
    public void onNext(T value) {
        // een extra verzending
        if (!isUnsubscribed()) {
            try {
                showInfos.accept(String.format("Subscriber [%s,%s] : onNext (%s)", observerName, processName,
                        new ObjectMapper().writeValueAsString(value)));
            } catch (JsonProcessingException e) {
                showInfos.accept(String.format("Subscriber [%s,%s].onNext (%s)", observerName, processName, e));
            }
        }
    }
}
  • Naast de kenmerken van een Subscriber zal de waarnemer Observateur de volgende informatie meenemen:
    • regel 14: een barrière of semafoor die wordt gebruikt om de hoofdthread te blokkeren totdat de observer alle door de observable verzonden elementen heeft ontvangen. Dit gebeurt op regel 36 van de code wanneer de observer van de observable de melding ontvangt dat het verzenden is voltooid;
    • regel 16: een instantie van Consumer<String> die wordt gebruikt om een bericht op de console weer te geven;
    • regel 18: de naam van de observer om ze van elkaar te onderscheiden wanneer er meerdere zijn;
    • regel 20: de naam van het geobserveerde proces;
  • regels 36, 46, 54: de methoden [onCompleted, onError, onNext] van de interface [Observer<T>], geïmplementeerd door de abstracte klasse [Subscriber<T>]. Deze klasse implementeert ze niet. Dit moet dus in de afgeleide klassen gebeuren. Voordat er iets in deze methoden wordt gedaan, wordt gecontroleerd of de waarnemer zich niet heeft afgemeld voor het waarneembare object dat hij observeert;
  • regel 59: de methode [onNext] van de observer schrijft de tekenreeks jSON van het ontvangen element. Hierdoor kunnen we verschillende soorten elementen weergeven;

Laten we nu een nieuwe methode van de klasse Observable bekijken, namelijk de methode [range]:

 

De observabele Observable.range(n,m) genereert (m) gehele getallen in het bereik van n tot n+m-1. We bestuderen deze met de volgende code [Exemple08]:


package dvp.rxjava.observables.exemples;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

import dvp.rxjava.observables.utils.Observateur;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple08 {
    public static void main(String[] args) throws InterruptedException {

        // aantal waarnemers
        final int nbObservateurs = 2;

        // semafoor
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // waarneembare configuratie
        Observable<Integer> obs1 = Observable.range(15, 3).subscribeOn(Schedulers.computation());
        // waarneembare uitvoering (waarneming)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            obs1.subscribe(new Observateur<>(String.format("observateur[%d]", i), latch, showInfos,"obs1"));
        }
        // wachten
        showInfos.accept("main : attente fin observation");
        latch.await();
        // einde
        showInfos.accept("main : fin observation");
    }

    // weergaven
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}
  • regel 16: we gaan twee waarnemers gebruiken;
  • regel 19: de semafoor wordt geïnitialiseerd op twee, omdat we elke waarnemer op een andere thread gaan plaatsen. De hoofdthread moet dus wachten tot beide waarnemingsthreads zijn voltooid;
  • regel 22: we configureren de observable zodanig dat deze wordt uitgevoerd op een thread van de scheduler [Schedulers.computation()]. De observer bevindt zich op dezelfde thread als de observable;
  • regels 25-27: we abonneren twee observers op de observable. Dit zal ervoor zorgen dat deze voor elk van de observers volledig wordt uitgevoerd: de gehele getallen 15, 16 en 17 worden verzonden;
  • regel 30: de hoofdthread wacht tot de observers klaar zijn;

De verkregen resultaten zijn als volgt:

main : début observation ------Thread[main] ---- Time[27:875]
main : attente fin observation ------Thread[main] ---- Time[27:893]
Subscriber[observateur[1],obs1] : onNext (15) ------Thread[RxComputationThreadPool-2] ---- Time[28:245]
Subscriber[observateur[0],obs1] : onNext (15) ------Thread[RxComputationThreadPool-1] ---- Time[28:245]
Subscriber[observateur[1],obs1] : onNext (16) ------Thread[RxComputationThreadPool-2] ---- Time[28:247]
Subscriber[observateur[0],obs1] : onNext (16) ------Thread[RxComputationThreadPool-1] ---- Time[28:248]
Subscriber[observateur[1],obs1] : onNext (17) ------Thread[RxComputationThreadPool-2] ---- Time[28:249]
Subscriber[observateur[1],obs1].onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[28:250]
Subscriber[observateur[0],obs1] : onNext (17) ------Thread[RxComputationThreadPool-1] ---- Time[28:251]
Subscriber[observateur[0],obs1].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[28:252]
main : fin observation ------Thread[main] ---- Time[28:252]
  • regel 2: de hoofdthread is geblokkeerd in afwachting van het einde van de 2 waarnemers;
  • regels 3-4: we zien dat waarnemer 0 zich op de thread [RxComputationThreadPool-1] bevindt en waarnemer 1 op de thread [RxComputationThreadPool-2];
  • regels 3-10: we zien dat beide waarnemers precies dezelfde elementen ontvangen;

We gaan de aldus gedefinieerde klasse Observateur gebruiken om het gedrag van andere soorten observables te illustreren.

7.3.2. Voorbeeld-09: de methoden van Observable.[interval, take, doNext]

  
 

Dit voorbeeld illustreert het gebruik van de observable Observable.interval (lang interval, eenheid TimeUnit), die met regelmatige tussenpozen lange gehele getallen uitzendt. Let op het punt [1]: standaard wordt de observable [Observable.interval] uitgevoerd op een van de threads van de scheduler [Schedulers.computation].

De code ziet er als volgt uit:


package dvp.rxjava.observables.exemples;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;

import dvp.rxjava.observables.utils.Observateur;
import rx.Observable;

public class Exemple09 {
    public static void main(String[] args) throws InterruptedException {

        // aantal waarnemers
        final int nbObservateurs = 2;

        // semafoor
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // waarneembare configuratie
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // waarneembare uitvoering (waarneming)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            obs1.subscribe(new Observateur<>(String.format("observateur [%d]", i), latch, showInfos,
                    "obs1"));
        }
        // wachten
        showInfos.accept("main : attente fin observation");
        latch.await();
        // einde
        showInfos.accept("main : fin observation");
    }

    // weergaven
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}
  • regel 22: de observable genereert elke 500 milliseconden long-getallen. De reeks begint met het getal 0;
  • regel 22: deze observable genereert een oneindig aantal waarden. De methode [Observable.take(n)] maakt een nieuwe observable aan die alleen de eerste n gegenereerde elementen bewaart;
 

Laten we nog eens terugkomen op de code van de observable:


Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
.doOnNext(l -> showInfos.accept(l.toString()));

Regel 2: de methode [Observable.doOnNext] wordt uitgevoerd telkens wanneer de observable een nieuw element uitzendt. Dit wordt vaak gebruikt om informatie te loggen. Hier willen we de datum van uitzending van de elementen loggen om te controleren of het interval van 500 milliseconden inderdaad wordt nageleefd. De methode [Observable.doOnNext] brengt geen wijzigingen aan in de observable waarop deze wordt toegepast. De definitie ervan is als volgt:

 

De uitvoering levert de volgende resultaten op:

main : début observation ------Thread[main] ---- Time[55:892]
main : attente fin observation ------Thread[main] ---- Time[55:911]
0 ------Thread[RxComputationThreadPool-1] ---- Time[56:412]
0 ------Thread[RxComputationThreadPool-2] ---- Time[56:413]
Subscriber[observateur [1],obs1] : onNext (0) ------Thread[RxComputationThreadPool-2] ---- Time[56:723]
Subscriber[observateur [0],obs1] : onNext (0) ------Thread[RxComputationThreadPool-1] ---- Time[56:723]
1 ------Thread[RxComputationThreadPool-1] ---- Time[56:906]
Subscriber[observateur [0],obs1] : onNext (1) ------Thread[RxComputationThreadPool-1] ---- Time[56:908]
1 ------Thread[RxComputationThreadPool-2] ---- Time[56:912]
Subscriber[observateur [1],obs1] : onNext (1) ------Thread[RxComputationThreadPool-2] ---- Time[56:914]
2 ------Thread[RxComputationThreadPool-1] ---- Time[57:405]
Subscriber[observateur [0],obs1] : onNext (2) ------Thread[RxComputationThreadPool-1] ---- Time[57:407]
Subscriber[observateur [0],obs1].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[57:408]
2 ------Thread[RxComputationThreadPool-2] ---- Time[57:412]
Subscriber[observateur [1],obs1] : onNext (2) ------Thread[RxComputationThreadPool-2] ---- Time[57:414]
Subscriber[observateur [1],obs1].onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[57:415]
main : fin observation ------Thread[main] ---- Time[57:416]
  • regels 3, 7 en 11: we zien dat het zendinterval ongeveer 500 ms bedraagt;
  • De twee observers bevinden zich uiteraard in twee verschillende threads, terwijl de observable niet was geconfigureerd om met een specifieke scheduler te worden uitgevoerd. Dit is de standaardwerking van de observable [Observable.interval] die we hier zien;

7.3.3. Voorbeelden-10/12: de methoden van Observable.[error, empty, never]

 

Vanaf nu zullen we de methoden van de klasse [Observable] beknopter illustreren. De vorige code zag er als volgt uit:


package dvp.rxjava.observables;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;

import rx.Observable;

public class Exemple09 {
    public static void main(String[] args) throws InterruptedException {

        // aantal waarnemers
        final int nbObservateurs = 2;

        // semafoor
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // waarneembare configuratie
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // waarneembare uitvoering (waarneming)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            obs1.subscribe(new Observateur<>(String.format("observateur [%d]", i), latch, showInfos,
                    "obs1"));
        }
        // wachten
        showInfos.accept("main : attente fin observation");
        latch.await();
        // einde
        showInfos.accept("main : fin observation");
    }

    // weergaven
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}

Deze code was al gebruikt voor het vorige voorbeeld. Alleen de regels 21-22 waren anders. We gaan daarom het grootste deel van deze code samenvoegen in de volgende klasse [ProcessUtils]:


package dvp.rxjava.observables.utils;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
import java.util.function.Consumer;

import rx.Observable;

public class ProcessUtils {

    @SafeVarargs
    public static void subscribe(int nbObservateurs, IProcess<?>... processes) throws InterruptedException {

        // semafoor
        CountDownLatch latch = new CountDownLatch(nbObservateurs * processes.length);

        // waarneembare uitvoering (waarneming)
        showInfos.accept("main : début observation");
        for (int i = 0; i < nbObservateurs; i++) {
            for (IProcess<?> process : processes) {
                Observable<?> obs = process.getObservable();
                obs.subscribe(new Observateur<>(String.format("observateur[%d]", i), latch, showInfos, process.getName()));
            }
        }
        // wachten
        showInfos.accept("main : attente fin observation");
        latch.await();
        // einde
        showInfos.accept("main : fin observation");
    }

    // weergaven
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));
}
  • regel 13: de methode accepteert twee parameters:
    • nbObservateurs: het aantal waarnemers van de processen die als tweede parameter worden doorgegeven;
    • processes: de (benoemde) processen die moeten worden geobserveerd. Dankzij de notatie [IProcess<?>] kunnen de processen elementen van verschillende typen verzenden;
  • regel 16: de semafoor moet op groen springen wanneer alle waarnemers al hun waarnemingen hebben voltooid. De beginwaarde van de semafoor is dus het aantal waarnemers vermenigvuldigd met het aantal waarnemingen;
  • regels 20-25: elke waarnemer wordt geabonneerd op alle processen die moeten worden geobserveerd;
  • regel 23: de observeerbare wordt opgehaald uit het proces (zie paragraaf 7.3.1);
  • regel 23: er wordt een waarnemer aan toegewezen. Aan deze waarnemer worden vier gegevens doorgegeven:
    • zijn naam;
    • de semafoor die hij moet verlagen wanneer hij de melding ontvangt dat de uitzending van de door hem geobserveerde observabele is beëindigd;
    • de methode die moet worden gebruikt wanneer hij informatie op de console wil loggen;
    • de naam van het proces dat hij gaat observeren;

Nu deze klassen zijn gedefinieerd, ziet voorbeeld 10 er als volgt uit:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple10 {
    public static void main(String[] args) throws InterruptedException {
        // waarneembare configuratie
        Observable<?> obs = Observable.error(new RuntimeException("Erreur !!!")).subscribeOn(Schedulers.computation());
        // waarneembare uitvoering (waarneming)
        ProcessUtils.subscribe(2,new Process<>("process1", obs));
    }
}

Op regel 11 wordt de statische methode [Observable.error] als volgt gedefinieerd:

 

Regel 8 configureert dus een observable die zich beperkt tot het genereren van een uitzondering naar de methode [onError] van zijn abonnees. De uitvoering levert de volgende resultaten op:


main : début observation ------Thread[main] ---- Time[22:618]
main : attente fin observation ------Thread[main] ---- Time[22:636]
Subscriber[observateur[1], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-2] ---- Time[22:638]
Subscriber[observateur[0], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-1] ---- Time[22:638]

In regel 3 en 4 heeft de methode [onError] van beide abonnees de door de observable gegenereerde uitzondering ontvangen.

Deze uitvoering heeft een bijzonderheid: de methoden [onCompleted] van beide waarnemers zijn niet aangeroepen. Daardoor is de barrière niet verlaagd en blijft de hoofdthread geblokkeerd in de statische methode [ProcessUtils.subscribe] op de volgende regel 3:


// wachten
showInfos.accept("main : attente fin observation");
latch.await();
// einde
showInfos.accept("main : fin observation");

Hier zien we dat bij een fout in de observable de methode [onCompleted] van de subscribers niet wordt aangeroepen. We passen de methode [Observateur.onError] daarom als volgt aan:


    @Override
    public void onError(Throwable e) {
        // verzendfout
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber[%s, %s].onError (%s)", observerName, processName, e));
        }
        // einde blokkering hoofdthread
        latch.countDown();
}

We voegen de regels 7-8 toe om de blokkering op te heffen in geval van een fout in de observable. Met deze nieuwe code levert de uitvoering de volgende resultaten op:


main : début observation ------Thread[main] ---- Time[40:750]
main : attente fin observation ------Thread[main] ---- Time[40:764]
Subscriber[observateur[0], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-1] ---- Time[40:766]
Subscriber[observateur[1], process1].onError (java.lang.RuntimeException: Erreur !!!) ------Thread[RxComputationThreadPool-2] ---- Time[40:766]
main : fin observation ------Thread[main] ---- Time[40:767]

We krijgen regel 5, die we eerder niet hadden.

Voorbeeld 11 ziet er als volgt uit:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple11 {
    public static void main(String[] args) throws InterruptedException {
        // waarneembare configuratie
        Observable<?> obs1 = Observable.empty();
        // uitvoering (observatie) waarneembaar
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

In regel 10 creëert de statische methode [Observable.empty] een observable die geen elementen uitzendt. Deze zendt alleen de melding van het einde van de uitzending uit;

 

De uitvoering van de code uit het bovenstaande voorbeeld levert de volgende resultaten op:

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[37:073]
Subscriber[observateur[0],process1].onCompleted ------Thread[main] ---- Time[37:086]
Subscriber[observateur[1],process1].onCompleted ------Thread[main] ---- Time[37:086]
main : attente fin observation ------Thread[main] ---- Time[37:087]
main : fin observation ------Thread[main] ---- Time[37:087]
  • regels 2 en 3: we zien dat beide observers de melding van het einde van de uitzending ontvangen zonder dat ze eerder elementen hebben ontvangen.

Men kan zich afvragen waarvoor deze methode precies dient. Men kan deze gebruiken op een manier die vergelijkbaar is met een verzameling, die aanvankelijk leeg is en waarin vervolgens elementen worden toegevoegd:

1
2
3
4
Observable obs=Observable.empty() ;
for(Observable o : observables){
    obs=obs.mergeWith(o) ;
}

In regel 3 wordt de oorspronkelijke observabele obs (regel 1) samengevoegd met andere observabelen.

Voorbeeld 12 illustreert de statische methode [Observable.never]:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple12 {
    public static void main(String[] args) throws InterruptedException {
        // waarneembare configuratie
        Observable<?> obs1 = Observable.never();
        // uitvoering (observatie) observeerbaar
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

De statische methode [Observable.never] creëert een observable die nooit een waarde uitzendt:

 

De uitvoering van het voorbeeld levert de volgende resultaten op:

main : début observation ------Thread[main] ---- Time[27:018]
main : attente fin observation ------Thread[main] ---- Time[27:030]

Regel 2: de hoofdthread wacht oneindig lang. Er is namelijk geen observable die de melding [onCompleted] verstuurt, waarmee de semafoor (slagboom) op groen kan worden gezet (de slagboom omlaag).

7.4. Multi-threading

7.4.1. Voorbeeld 13: actiethread, observatiethread

In paragraaf 7.1.3 hebben we een observable aangemaakt met de statische methode [Observable.create]:

 
  • de methode [create] retourneert een type Observable<T>;
  • de parameter van de methode [create] is een functie van het type [Observable.OnSubscribe<T>], die als volgt is gedefinieerd:
 

Het type [Observable.OnSubscribe<T>] is een functionele interface die zelf de functionele interface [Action1<Subscriber<? super T>>] uitbreidt. De methode [call] van deze interface verwacht een type [Subscriber] (abonnee, inschrijver, waarnemer). In het verdere verloop van dit document zullen we het type [Observable.OnSubscribe<T>] soms een actie noemen. We gaan aangepaste acties maken die een naam krijgen. Dit zijn instanties van de volgende interface [IProcessAction]:

  

package dvp.rxjava.observables.utils;

import rx.Observable;

public interface IProcessAction<T> extends Observable.OnSubscribe<T> {

    // de actie heeft een naam
    public String getName();
}
  • regel 5: de interface [IProcessAction<T>] heeft alle kenmerken van de interface [Observable.OnSubscribe<T>];
  • regel 8: deze heeft bovendien een methode [getName] die de naam van de instantie die de interface implementeert, retourneert;

We gaan de volgende actie met de naam [ProcessAction01] gebruiken:


package dvp.rxjava.observables.utils;

import java.util.Random;

import rx.Subscriber;
import rx.functions.Func1;

public class ProcessAction01<T> implements IProcessAction<T> {

    // gegevens
    private String name;
    private int nbValues;
    private Func1<Integer, T> func1;

    // constructors
    public ProcessAction01(String name, int nbValues, Func1<Integer, T> func1) {
        this.name = name;
        this.nbValues = nbValues;
        this.func1 = func1;
    }

    @Override
    public void call(Subscriber<? super T> subscriber) {
        ProcessUtils.showInfos.accept(String.format("Observable (%s) call start", getName()));
        for (int i = 0; i < nbValues; i++) {
            // verwachting
            try {
                Thread.sleep(new Random().nextInt(500));
            } catch (InterruptedException e) {
                // fout
                ProcessUtils.showInfos.accept(String.format("Observable (%s) onError", getName()));
                subscriber.onError(e);
            }
            // uitzending van een element
            T value = func1.call(i);
            ProcessUtils.showInfos.accept(String.format("Observable (%s,%s) onNext (%s)", getName(), i, value));
            subscriber.onNext(value);
        }
        // voltooid
        ProcessUtils.showInfos.accept(String.format("Observable (%s) onCompleted", getName()));
        subscriber.onCompleted();
    }

    @Override
    public String getName() {
        return name;
    }

}
  • regel 8: de klasse [ProcessAction01<T>] implementeert de interface [IProcessAction<T>] en dus ook de interface [Observable.OnSubscribe<T>];
  • regel 11: de naam van de actie;
  • regel 12: het aantal uit te zenden waarden;
  • regel 13: een instantie van het type [Func1<Integer, T>] die op basis van een geheel getal een type T genereert dat door de observable wordt uitgezonden (regels 35 en 37);
  • regels 16-20: we geven de constructor de naam van de actie, het aantal uit te zenden waarden en de uitzendfunctie door;
  • regels 23-42: de procescode;
  • regel 23: de methode [call] ontvangt als parameter de abonnee van de observable die aan het proces is gekoppeld;
  • regel 28: het proces zendt zijn elementen uit na een willekeurige wachttijd;
  • regel 32: het verzenden van een fout;
  • regel 37: een normale verzending;
  • regel 41: verzending van de melding dat de verzending is voltooid;
  • regels 25-38: de actie verzendt reële waarden naar nbValues na een willekeurige wachttijd (regel 30);
  • regel 35: de uit te zenden waarde wordt geleverd door de functie [func1] die als parameter aan de constructor wordt doorgegeven (regel 16);

We herstructureren de klasse [Process] (zie paragraaf 7.3.1) zodat deze ook met een benoemde actie kan worden geconstrueerd. We voegen de volgende constructor toe:


public Process(IProcessAction<T> na, Scheduler schedulerObserved, Scheduler schedulerObserver) {
        // procesnaam=actienaam
        name = na.getName();
        // actie --> waarneembaar
        observable = Observable.create(na);
        // uitvoeringsthread van het geobserveerde proces
        if (schedulerObserved != null) {
            observable = observable.subscribeOn(schedulerObserved);
        }
        // observatiethread van de waarnemer
        if (schedulerObserver != null) {
            observable = observable.observeOn(schedulerObserver);
        }
    }
  • regel 1: de constructor accepteert 3 parameters:
    1. de benoemde actie die zal worden gebruikt om de observable te construeren (regel 5);
    2. de scheduler van het geobserveerde proces (kan null zijn);
    3. de scheduler van de waarnemer (bijvoorbeeld null);
  • regel 5: de observable wordt aangemaakt op basis van de als parameter doorgegeven actie;

De volgende code [Exemple13] observeert verschillende observables:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple13 {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 1, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        // proces 3
        Process<Integer> process3 = new Process<>(new ProcessAction01<Integer>("process3", 3, i -> i * 2), null,
                Schedulers.computation());
        // proces 4
        Process<Boolean> process4 = new Process<>(new ProcessAction01<Boolean>("process4", 4, i -> i % 2 == 0), null, null);
        // inschrijvingen
        ProcessUtils.subscribe(1, process1);
        ProcessUtils.subscribe(1, process2);
        ProcessUtils.subscribe(1, process3);
        ProcessUtils.subscribe(1, process4);
    }
}
  • regels 13-15: het proces process1 genereert 1 reëel getal op een rekenthread dat op een andere rekenthread wordt geobserveerd;
  • regels 17-18: het proces process2 genereert 2 tekenreeksen op een rekenthread en er wordt geen indicatie gegeven over de thread van de waarnemer. De resultaten tonen aan dat de waarneming standaard plaatsvindt op dezelfde thread als die waarop het proces wordt uitgevoerd;
  • regels 20-21: het proces process3 genereert 3 gehele getallen op een niet-voorgeschreven thread, die zullen worden geobserveerd op een rekenthread. De resultaten tonen aan dat de uitvoering van het proces standaard plaatsvindt op de hoofdthread;
  • regel 23: het proces process4 genereert 4 booleaanse waarden op een niet-voorgeschreven thread, die worden geobserveerd op een niet-voorgeschreven thread. Uit de resultaten blijkt dat zowel de uitvoering van het proces als de observatie ervan standaard plaatsvinden op de hoofdthread;

Het resultaat van de uitvoering van deze code is als volgt:

main : début observation ------Thread[main] ---- Time[18:642]
main : attente fin observation ------Thread[main] ---- Time[18:660]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[18:660]
Observable (process1,0) onNext (68.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[19:093]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[19:094]
Subscriber[observateur[0],process1] : onNext (68.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[19:396]
Subscriber[observateur[0],process1].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[19:397]
main : fin observation ------Thread[main] ---- Time[19:397]
main : début observation ------Thread[main] ---- Time[19:398]
main : attente fin observation ------Thread[main] ---- Time[19:399]
Observable (process2) call start ------Thread[RxComputationThreadPool-5] ---- Time[19:399]
Observable (process2,0) onNext (valeur-0) ------Thread[RxComputationThreadPool-5] ---- Time[19:630]
Subscriber[observateur[0],process2] : onNext ("valeur-0") ------Thread[RxComputationThreadPool-5] ---- Time[19:631]
Observable (process2,1) onNext (valeur-1) ------Thread[RxComputationThreadPool-5] ---- Time[20:094]
Subscriber[observateur[0],process2] : onNext ("valeur-1") ------Thread[RxComputationThreadPool-5] ---- Time[20:095]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[20:096]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[20:096]
main : fin observation ------Thread[main] ---- Time[20:097]
main : début observation ------Thread[main] ---- Time[20:097]
Observable (process3) call start ------Thread[main] ---- Time[20:098]
Observable (process3,0) onNext (0) ------Thread[main] ---- Time[20:188]
Subscriber[observateur[0],process3] : onNext (0) ------Thread[RxComputationThreadPool-6] ---- Time[20:213]
Observable (process3,1) onNext (2) ------Thread[main] ---- Time[20:336]
Subscriber[observateur[0],process3] : onNext (2) ------Thread[RxComputationThreadPool-6] ---- Time[20:338]
Observable (process3,2) onNext (4) ------Thread[main] ---- Time[20:676]
Observable (process3) onCompleted ------Thread[main] ---- Time[20:677]
main : attente fin observation ------Thread[main] ---- Time[20:677]
Subscriber[observateur[0],process3] : onNext (4) ------Thread[RxComputationThreadPool-6] ---- Time[20:678]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[20:679]
main : fin observation ------Thread[main] ---- Time[20:679]
main : début observation ------Thread[main] ---- Time[20:680]
Observable (process4) call start ------Thread[main] ---- Time[20:680]
Observable (process4,0) onNext (true) ------Thread[main] ---- Time[21:065]
Subscriber[observateur[0],process4] : onNext (true) ------Thread[main] ---- Time[21:067]
Observable (process4,1) onNext (false) ------Thread[main] ---- Time[21:187]
Subscriber[observateur[0],process4] : onNext (false) ------Thread[main] ---- Time[21:188]
Observable (process4,2) onNext (true) ------Thread[main] ---- Time[21:624]
Subscriber[observateur[0],process4] : onNext (true) ------Thread[main] ---- Time[21:625]
Observable (process4,3) onNext (false) ------Thread[main] ---- Time[21:765]
Subscriber[observateur[0],process4] : onNext (false) ------Thread[main] ---- Time[21:766]
Observable (process4) onCompleted ------Thread[main] ---- Time[21:767]
Subscriber[observateur[0],process4].onCompleted ------Thread[main] ---- Time[21:767]
main : attente fin observation ------Thread[main] ---- Time[21:767]
main : fin observation ------Thread[main] ---- Time[21:768]
  • het proces process1 genereert 1 reëel getal (regel 4) op de reken-thread [RxComputationThreadPool-4], dat wordt geobserveerd op de reken-thread [RxComputationThreadPool-3] (regel 6);
  • het proces process2 genereert 2 tekenreeksen (regels 12, 14) op de rekenthread [RxComputationThreadPool-5], die op diezelfde thread worden waargenomen (regels 13, 15);
  • het proces process3 genereert 3 gehele getallen (regels 21, 23, 25) op de hoofdthread, die worden waargenomen op de rekenthread [RxComputationThreadPool-6] (regels 22, 24, 28);
  • het proces process4 genereert 4 booleaanse waarden (regels 34, 36, 38, 40) in de hoofdthread, die in diezelfde hoofdthread worden waargenomen (regels 33, 35, 37, 39);

De lezer wordt uitgenodigd om hierboven het volgende te volgen:

  • de levenscyclus van het geobserveerde proces en de bijbehorende thread;
  • de levenscyclus van de waarnemer en de bijbehorende thread;

Een groot deel van het belang van Rx-bibliotheken berust op deze multithreading, die de ontwikkelaar niet zelf hoeft te beheren.

7.5. Combinaties van meerdere observables

7.5.1. Voorbeeld 14: twee observables samenvoegen met [Observable.merge]

We presenteren nu statische methoden van de klasse [Observable] waarmee meerdere observables kunnen worden gecombineerd tot één resulterende observable.

Het eerste voorbeeld hiervan is het volgende:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.ProcessAction01;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple14 {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        // samenvoeging
        Process<?> process12 = new Process<>("process12",
                Observable.merge(process1.getObservable(), process2.getObservable()));
        // inschrijvingen
        ProcessUtils.subscribe(1, process12);
    }
}
  • regels 15-17: een proces met de naam [process1] zal 3 reële getallen uitvoeren op een rekenthread. Het zal ook worden geobserveerd op een rekenthread;
  • regels 19-20: een proces met de naam [process2] zal 2 tekenreeksen naar een reken-thread sturen. De observatiethread is niet vastgelegd. We hebben eerder gezien dat in dit geval de observatiethread de reken-thread is;
  • regel 23: de twee processen worden samengevoegd, d.w.z. er wordt een observable aangemaakt waarvan de elementen gelijktijdig uit beide processen afkomstig zijn. Hiervoor wordt de statische methode [Observable.merge] gebruikt:
 

In tegenstelling tot wat het bovenstaande schema zou doen vermoeden, kunnen de elementen van stroom 1 tijdens het samenvoegen tussen de elementen van stroom 2 worden ingevoegd. Dit blijkt uit de resultaten van de uitvoering:

main : début observation ------Thread[main] ---- Time[56:053]
main : attente fin observation ------Thread[main] ---- Time[56:073]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[56:073]
Observable (process2) call start ------Thread[RxComputationThreadPool-5] ---- Time[56:074]
Observable (process1,0) onNext (64.8) ------Thread[RxComputationThreadPool-4] ---- Time[56:263]
Observable (process2,0) onNext (valeur-0) ------Thread[RxComputationThreadPool-5] ---- Time[56:403]
Observable (process2,1) onNext (valeur-1) ------Thread[RxComputationThreadPool-5] ---- Time[56:515]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[56:516]
Subscriber[observateur[0],process12] : onNext (64.8) ------Thread[RxComputationThreadPool-3] ---- Time[56:552]
Subscriber[observateur[0],process12] : onNext ("valeur-0") ------Thread[RxComputationThreadPool-3] ---- Time[56:553]
Subscriber[observateur[0],process12] : onNext ("valeur-1") ------Thread[RxComputationThreadPool-3] ---- Time[56:553]
Observable (process1,1) onNext (56.4) ------Thread[RxComputationThreadPool-4] ---- Time[56:716]
Subscriber[observateur[0],process12] : onNext (56.4) ------Thread[RxComputationThreadPool-3] ---- Time[56:718]
Observable (process1,2) onNext (22.8) ------Thread[RxComputationThreadPool-4] ---- Time[57:082]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[57:083]
Subscriber[observateur[0],process12] : onNext (22.8) ------Thread[RxComputationThreadPool-3] ---- Time[57:084]
Subscriber[observateur[0],process12].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[57:085]
main : fin observation ------Thread[main] ---- Time[57:085]
  • regel 3: het proces [process1] wordt uitgevoerd op de rekenthread [RxComputationThreadPool-4];
  • regel 4: het proces [process2] wordt uitgevoerd op de rekenthread [RxComputationThreadPool-5];
  • regel 9: het proces [process12] wordt waargenomen op de rekenthread [RxComputationThreadPool-3]. Ik weet niet welke regel tot deze keuze heeft geleid;
  • regels 9-11: we zien dat de waarnemer elementen waarneemt van zowel proces [process1] (regel 5) als proces [process2] (regels 6, 7), terwijl geen van beide is voltooid (er is sprake van vermenging);
  • het proces [process12] wordt beëindigd (regel 17) zodra de twee processen process1 en process2 zijn beëindigd;

7.5.2. Voorbeeld 15: twee observables samenvoegen met [Observable.concat]

We bekijken nu de volgende code:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.ProcessAction01;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple15 {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, Schedulers.computation());
        // samenvoegen
        Process<?> process12 = new Process<>("process12",
                Observable.concat(process1.getObservable(), process2.getObservable()));
        // inschrijvingen
        ProcessUtils.subscribe(1, process12);
    }
}
  • regels 15-17: een proces met de naam [process1] zal 3 reële getallen naar een rekenthread sturen. Het zal ook op een rekenthread worden geobserveerd;
  • regels 19-20: een proces met de naam [process2] zal 2 tekenreeksen uitzenden naar een niet-voorgeschreven thread, in dit geval de standaard hoofdthread. Het zal worden geobserveerd op een rekenthread;
  • regel 23: de twee processen worden samengevoegd, d.w.z. er wordt een observable aangemaakt waarvan de elementen afkomstig zijn uit beide processen. De uitgezonden waarden worden niet gemengd. Het proces [process12] zal eerst alle waarden van het proces [process1] uitzenden en vervolgens die van het proces [process2]. Hiervoor wordt de statische methode [Observable.concat] gebruikt:
 

De resultaten van de uitvoering zijn als volgt:

main : début observation ------Thread[main] ---- Time[30:162]
main : attente fin observation ------Thread[main] ---- Time[30:189]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[30:190]
Observable (process1,0) onNext (79.2) ------Thread[RxComputationThreadPool-4] ---- Time[30:681]
Observable (process1,1) onNext (98.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[30:792]
Subscriber[observateur[0],process12] : onNext (79.2) ------Thread[RxComputationThreadPool-3] ---- Time[30:975]
Subscriber[observateur[0],process12] : onNext (98.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[30:976]
Observable (process1,2) onNext (84.0) ------Thread[RxComputationThreadPool-4] ---- Time[31:084]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[31:085]
Subscriber[observateur[0],process12] : onNext (84.0) ------Thread[RxComputationThreadPool-3] ---- Time[31:086]
Observable (process2) call start ------Thread[RxComputationThreadPool-3] ---- Time[31:087]
Observable (process2,0) onNext (valeur-0) ------Thread[RxComputationThreadPool-3] ---- Time[31:556]
Subscriber[observateur[0],process12] : onNext ("valeur-0") ------Thread[RxComputationThreadPool-5] ---- Time[31:557]
Observable (process2,1) onNext (valeur-1) ------Thread[RxComputationThreadPool-3] ---- Time[31:608]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[31:609]
Subscriber[observateur[0],process12] : onNext ("valeur-1") ------Thread[RxComputationThreadPool-5] ---- Time[31:609]
Subscriber[observateur[0],process12].onCompleted ------Thread[RxComputationThreadPool-5] ---- Time[31:610]
main : fin observation ------Thread[main] ---- Time[31:611]
  • regels 3-10: het proces [process1] wordt uitgevoerd en het proces [process12] geeft de waarden door die door [process1] zijn gegenereerd;
  • regel 9: het proces [process1] is voltooid;
  • regels 11-17: het proces [process2] wordt uitgevoerd en het proces [process12] geeft de waarden door die door [process2] zijn gegenereerd;

Er is iets vreemds aan de hand met het proces process2: er was geen uitvoeringsthread opgegeven. Men zou dan kunnen verwachten dat dit standaard de hoofdthread zou zijn. Dat is echter niet het geval. De uitvoeringsthread was de berekeningsthread [RxComputationThreadPool-3] (regel 11). Wanneer er dus geen uitvoerings- of observatiethread wordt opgegeven, kan er geen aanname worden gedaan over welke thread er zal worden gekozen.

7.5.3. Voorbeeld 16: twee observabelen combineren met [Observable.zip]

We bekijken nu de volgende code:


package dvp.rxjava.observables.exemples;

import java.util.Arrays;
import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.functions.FuncN;
import rx.schedulers.Schedulers;

public class Exemple16 {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, null);
        // functie voor het combineren van de 2 processen
        FuncN<String> funcn = new FuncN<String>() {
            @Override
            public String call(Object... args) {
                if (args.length == 2) {
                    return String.format("double=%s, string=%s", args[0], args[1]);
                } else {
                    throw new RuntimeException("la fonction attend 2 paramètres exactement");
                }
            }
        };
        // de twee processen samenvoegen
        Process<String> process12 = new Process<>("process12",
                Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));
        // inschrijvingen
        ProcessUtils.subscribe(1, process12);
    }
}
  • regels 16-18: een proces met de naam [process1] zal 3 reële getallen uitvoeren op een rekenthread. Het zal ook worden geobserveerd op een rekenthread;
  • regels 20-21: een proces met de naam [process2] zal 2 tekenreeksen naar een niet-opgelegde thread verzenden. De observatiethread is evenmin opgelegd;
  • regels 23-32: instantiëring van een type [FuncN<String>] met een anonieme klasse. FuncN is een functionele interface:
 

De methode [FuncN.call] verwacht een array van objecten en retourneert een type R. De functie [funcn] zal worden gebruikt om de processen process1 en process2 in deze volgorde te combineren. In de methode [FuncN.call]:

  • zal args[0] een Double zijn;
  • args[1] is een String;

Hier zal het resultaat van [funcn.call] de tekenreeks van regel 27 zijn. Voor het samenstellen van dit resultaat is het niet nodig om de typen van de argumenten van de methode call te kennen.

De twee processen worden als volgt gecombineerd:


// de twee processen samenvoegen
Process<String> process12 = new Process<>("process12",
Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));

De methode [Observable.zip] werkt als volgt:

 

We zien dat:

  • het eerste argument van zip een Iterable<Observable> is. In ons voorbeeld hebben we een daadwerkelijke parameter van het type List<Observable>, gevormd door onze twee observables;
  • het tweede argument van zip is van het type FuncN. In ons voorbeeld is de effectieve parameter [funcn];

De uitvoering levert de volgende resultaten op:

main : début observation ------Thread[main] ---- Time[55:636]
Observable (process2) call start ------Thread[main] ---- Time[55:666]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[55:666]
Observable (process1,0) onNext (69.6) ------Thread[RxComputationThreadPool-4] ---- Time[55:902]
Observable (process2,0) onNext (valeur-0) ------Thread[main] ---- Time[56:076]
Observable (process1,1) onNext (82.8) ------Thread[RxComputationThreadPool-4] ---- Time[56:271]
Subscriber[observateur[0],process12] : onNext ("double=69.6, string=valeur-0") ------Thread[main] ---- Time[56:352]
Observable (process1,2) onNext (14.399999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[56:641]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[56:642]
Observable (process2,1) onNext (valeur-1) ------Thread[main] ---- Time[56:778]
Subscriber[observateur[0],process12] : onNext ("double=82.8, string=valeur-1") ------Thread[main] ---- Time[56:779]
Observable (process2) onCompleted ------Thread[main] ---- Time[56:779]
Subscriber[observateur[0],process12].onCompleted ------Thread[main] ---- Time[56:780]
main : attente fin observation ------Thread[main] ---- Time[56:781]
main : fin observation ------Thread[main] ---- Time[56:781]
  • regels 7, 11: het proces process12 zendt twee elementen uit;
  • regel 8: het extra element dat wordt verzonden door het proces process1, dat geen partner heeft in het proces process2, wordt niet verzonden door het resultaatproces process12;

We zien dat het proces process2, waaraan noch een uitvoeringsthread noch een observatiethread was toegewezen, de hoofdthread voor beide heeft gebruikt.

7.5.4. Voorbeeld 17: twee observables combineren met [Observable.combineLatest]

We bekijken nu de volgende code:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple17 {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null,
                Schedulers.computation());
        // combinatie van de 2 processen
        Process<Double> process12 = new Process<>("process12",
                Observable.combineLatest(process1.getObservable(), process2.getObservable(), (d1, d2) -> d1 + d2));
        // inschrijvingen
        ProcessUtils.subscribe(1, process12);
    }
}
  • regels 14-16: een proces met de naam [process1] zal 3 reële getallen uitzenden via een reken-thread. Het zal ook worden geobserveerd via een reken-thread;
  • regels 18-20: een proces met de naam [process2] zal 2 reële getallen uitzenden op een niet-opgelegde thread. Deze worden waargenomen op een reken-thread;
  • regel 23: de twee observabelen worden gecombineerd met de volgende statische methode [Observable.combineLatest]:
 

De observabele [combineLatest] werkt als volgt: wanneer een van de twee observabelen een element E1 uitzendt, wordt dit element door [combineFunction] gecombineerd met het laatste element dat door de andere observabele is uitgezonden.

De uitvoering van deze code levert het volgende resultaat op:

main : début observation ------Thread[main] ---- Time[01:768]
Observable (process2) call start ------Thread[main] ---- Time[01:791]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[01:791]
Observable (process1,0) onNext (54.0) ------Thread[RxComputationThreadPool-4] ---- Time[01:991]
Observable (process2,0) onNext (56.0) ------Thread[main] ---- Time[02:245]
Observable (process1,1) onNext (51.6) ------Thread[RxComputationThreadPool-4] ---- Time[02:358]
Subscriber[observateur[0],process12] : onNext (110.0) ------Thread[RxComputationThreadPool-5] ---- Time[02:521]
Subscriber[observateur[0],process12] : onNext (107.6) ------Thread[RxComputationThreadPool-5] ---- Time[02:522]
Observable (process2,1) onNext (261.8) ------Thread[main] ---- Time[02:595]
Observable (process2) onCompleted ------Thread[main] ---- Time[02:596]
main : attente fin observation ------Thread[main] ---- Time[02:596]
Subscriber[observateur[0],process12] : onNext (313.40000000000003) ------Thread[RxComputationThreadPool-5] ---- Time[02:597]
Observable (process1,2) onNext (80.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[02:790]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[02:791]
Subscriber[observateur[0],process12] : onNext (342.2) ------Thread[RxComputationThreadPool-3] ---- Time[02:792]
Subscriber[observateur[0],process12].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[02:792]
main : fin observation ------Thread[main] ---- Time[02:793]
  • regel 5: de uitzending van process2 (56) wordt gecombineerd met het laatste element dat is uitgezonden door process1 (54, regel 4) en levert het resultaat van regel 7 op;
  • regel 6: de uitzending van process1 (51,6) wordt gecombineerd met het laatste element dat is uitgezonden door process2 (56, regel 5) en levert het resultaat van regel 8 op;
  • regel 9: de uitzending van process2 (261,8) wordt gecombineerd met het laatste element dat is uitgezonden door process1 (51,6, regel 6) en levert het resultaat van regel 12 op;
  • regel 13: de uitzending van process1 (80,39) wordt gecombineerd met het laatste element dat is uitgezonden door process2 (261,8, regel 9) en levert het resultaat van regel 15 op;

We hebben hier te maken met een variant van de observable [zip], waarbij de gecombineerde elementen ditmaal niet noodzakelijkerwijs de elementen op dezelfde positie in de stromen zijn. We merken hier op dat het proces process2, waaraan geen uitvoeringsthread was toegewezen, hier op de hoofdthread is uitgevoerd (regel 2).

7.5.5. Voorbeeld 18: twee observables combineren met [Observable.amb]

We bekijken nu de volgende code:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple18 {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null, null);
        // combinatie van de 2 processen
        Process<Double> process12 = new Process<>("process12",
                Observable.amb(process1.getObservable(), process2.getObservable()));
        // inschrijvingen
        ProcessUtils.subscribe(1, process12);
    }
}
  • regels 14-16: een proces met de naam [process1] zal 3 reële getallen uitzenden op een rekenthread. Het zal ook worden geobserveerd op een rekenthread;
  • regels 18-20: een proces met de naam [process2] zal 2 reële getallen uitzenden naar een niet-opgelegde thread. Deze worden waargenomen op een niet-opgelegde thread;
  • regel 22: de twee observabelen worden gecombineerd met de volgende statische methode [Observable.amb]:
 

Zoals het bovenstaande schema laat zien, zendt de observable [Observable.amb(Observable o1, Observable o2)] de elementen uit van de observable die als eerste uitzendt. Dit wordt bevestigd door de resultaten van het gepresenteerde voorbeeld:

main : début observation ------Thread[main] ---- Time[21:594]
Observable (process2) call start ------Thread[main] ---- Time[21:612]
Observable (process1) call start ------Thread[RxComputationThreadPool-3] ---- Time[21:612]
Observable (process2,0) onNext (155.39999999999998) ------Thread[main] ---- Time[21:817]
Observable (process1) onError ------Thread[RxComputationThreadPool-3] ---- Time[21:820]
Observable (process1,0) onNext (90.0) ------Thread[RxComputationThreadPool-3] ---- Time[21:820]
Observable (process1,1) onNext (104.39999999999999) ------Thread[RxComputationThreadPool-3] ---- Time[21:877]
Subscriber[observateur[0],process12] : onNext (155.39999999999998) ------Thread[main] ---- Time[22:105]
Observable (process1,2) onNext (44.4) ------Thread[RxComputationThreadPool-3] ---- Time[22:122]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[22:123]
Observable (process2,1) onNext (201.6) ------Thread[main] ---- Time[22:581]
Subscriber[observateur[0],process12] : onNext (201.6) ------Thread[main] ---- Time[22:583]
Observable (process2) onCompleted ------Thread[main] ---- Time[22:583]
Subscriber[observateur[0],process12].onCompleted ------Thread[main] ---- Time[22:584]
main : attente fin observation ------Thread[main] ---- Time[22:585]
main : fin observation ------Thread[main] ---- Time[22:586]
  • regel 4: het is het proces process2 dat als eerste elementen uitzendt;
  • regels 8, 12: het proces process12 verzendt alle elementen die door het proces process2 zijn verzonden (regels 4, 11);

7.6. Verwerkingsketen van een observable

7.6.1. Voorbeeld 19: een observable transformeren met [Observable.map]

In de voorgaande voorbeelden hebben we verschillende combinaties van twee observables tot een derde observable bekeken. We presenteren nu statische methoden van de klasse [Observable] die transformatie-, filter- en aggregatiebewerkingen op een observable mogelijk maken. We zullen hier methoden aantreffen die analoog zijn aan die van de klasse [Stream] die in paragraaf 5 zijn besproken.

Ons eerste voorbeeld is het volgende:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple19 {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<String> process2 = new Process<>("process2",
                process1.getObservable().map(d -> String.format("valeur-%s", d)));
        // inschrijvingen
        ProcessUtils.subscribe(1, process2);
    }
}
  • regels 14-16: een proces met de naam process1 zal 3 reële getallen genereren op een rekenthread. Het zal ook worden geobserveerd op een rekenthread;
  • regels 17-18: de door process1 uitgevoerde getallen worden omgezet in tekenreeksen in een proces process2;
  • regel 20: process2 wordt geobserveerd;

De methode [Observable.map] uit regel 18 is vergelijkbaar met de methode [Stream.map] die in paragraaf 5.5 is besproken:

 

De resultaten van het voorbeeld zijn als volgt:

main : début observation ------Thread[main] ---- Time[55:328]
main : attente fin observation ------Thread[main] ---- Time[55:346]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[55:347]
Observable (process1,0) onNext (21.599999999999998) ------Thread[RxComputationThreadPool-4] ---- Time[55:354]
Observable (process1,1) onNext (97.2) ------Thread[RxComputationThreadPool-4] ---- Time[55:512]
Subscriber[observateur[0],process2] : onNext ("valeur-21.599999999999998") ------Thread[RxComputationThreadPool-3] ---- Time[55:615]
Subscriber[observateur[0],process2] : onNext ("valeur-97.2") ------Thread[RxComputationThreadPool-3] ---- Time[55:616]
Observable (process1,2) onNext (98.39999999999999) ------Thread[RxComputationThreadPool-4] ---- Time[55:803]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[55:804]
Subscriber[observateur[0],process2] : onNext ("valeur-98.39999999999999") ------Thread[RxComputationThreadPool-3] ---- Time[55:804]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[55:805]
main : fin observation ------Thread[main] ---- Time[55:805]
  • regels 4, 5 en 8: de uitvoer van process1. Dit zijn reële getallen;
  • regels 6, 7, 10: de waargenomen emissies van process2. Dit zijn tekenreeksen;

7.6.2. Voorbeeld-20: een observabele filteren met [Observable.filter]

Het voorbeeld ziet er als volgt uit:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple20 {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().filter(i -> i % 2 == 0));
        // abonnementen
        ProcessUtils.subscribe(1, process2);
    }
}
  • regels 11-12: een proces met de naam process1 zal de gehele getallen van 0 tot 2 uitzenden op een rekenthread. Het zal ook worden geobserveerd op een rekenthread;
  • regel 14: de getallen die door process1 worden gegenereerd, worden gefilterd, zodat alleen de even getallen in process2 overblijven;
  • regel 20: process2 wordt geobserveerd;

De methode [Observable.filter] uit regel 18 is analoog aan de methode [Stream.filter] die in paragraaf 5.4 is besproken:

 

De resultaten van het voorbeeld zijn als volgt:

main : début observation ------Thread[main] ---- Time[30:319]
main : attente fin observation ------Thread[main] ---- Time[30:335]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[30:336]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[30:388]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[30:625]
Subscriber[observateur[0],process2] : onNext (0) ------Thread[RxComputationThreadPool-3] ---- Time[30:703]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[30:704]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[30:705]
Subscriber[observateur[0],process2] : onNext (2) ------Thread[RxComputationThreadPool-3] ---- Time[30:706]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[30:707]
main : fin observation ------Thread[main] ---- Time[30:707]
  • regels 4, 5 en 7: de uitzendingen van process1;
  • regels 6, 9: de waargenomen uitzendingen van process2. Dit zijn de elementen van process1 die even zijn;

7.6.3. Voorbeeld-21: een observabele transformeren met [Observable.flatMap]

Het voorbeeld is als volgt:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple21 {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMap(i -> {
            int value = i * 10;
            return Observable.just(value, value + 1, value + 2);
        }));
        // abonnementen
        ProcessUtils.subscribe(1, process2);
    }
}
  • regels 12-13: een proces met de naam process1 zal de gehele getallen van 0 tot 2 uitzenden op een rekenthread. Het zal ook worden geobserveerd op een rekenthread;
  • regels 15-18: elk getal n dat door process1 wordt uitgezonden, wordt omgezet in een observable die de 3 getallen (10*n, 10*n+1, 10*n+2) uitzendt. Als in regel 15 de methode [map] zou worden gebruikt, zou process2 een type Observable<Integer> genereren en niet een type Integer. Met de gebruikte methode [flatMap] kan (flatten) deze reeks elementen van het type Observable<Integer> te vereenvoudigen tot een reeks elementen van het type Integer, bestaande uit elk van de elementen van elk van de Observable<Integer>;
  • regel 20: we zien process2;

De methode [Observable.flatMap] uit regel 15 is analoog aan de methode [Stream.flatMap] die in paragraaf 5.6.12 is besproken:

 

De resultaten van het voorbeeld zijn als volgt:

main : début observation ------Thread[main] ---- Time[31:466]
main : attente fin observation ------Thread[main] ---- Time[31:486]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[31:486]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[31:777]
Subscriber[observateur[0],process2] : onNext (0) ------Thread[RxComputationThreadPool-3] ---- Time[32:082]
Subscriber[observateur[0],process2] : onNext (1) ------Thread[RxComputationThreadPool-3] ---- Time[32:085]
Subscriber[observateur[0],process2] : onNext (2) ------Thread[RxComputationThreadPool-3] ---- Time[32:087]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[32:192]
Subscriber[observateur[0],process2] : onNext (10) ------Thread[RxComputationThreadPool-3] ---- Time[32:194]
Subscriber[observateur[0],process2] : onNext (11) ------Thread[RxComputationThreadPool-3] ---- Time[32:196]
Subscriber[observateur[0],process2] : onNext (12) ------Thread[RxComputationThreadPool-3] ---- Time[32:197]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[32:686]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[32:687]
Subscriber[observateur[0],process2] : onNext (20) ------Thread[RxComputationThreadPool-3] ---- Time[32:688]
Subscriber[observateur[0],process2] : onNext (21) ------Thread[RxComputationThreadPool-3] ---- Time[32:690]
Subscriber[observateur[0],process2] : onNext (22) ------Thread[RxComputationThreadPool-3] ---- Time[32:692]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[32:693]
main : fin observation ------Thread[main] ---- Time[32:693]
  • regels 5-7: de drie uitzendingen van process2 na de uitzending van regel 4 van process1;
  • regels 9-11: de drie uitzendingen van process2 na de uitzending van regel 8 van process1;
  • regels 14-16: de drie uitzendingen van process2 naar aanleiding van de uitzending van regel 12 van process1;

De volgende code laat zien hoe je een type Observable<Integer[]> kunt aanmaken op basis van process1 en [Exemple21b]:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21b {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<Integer[]> process2 = new Process<>("process2", process1.getObservable().map(i -> {
            int value = i * 10;
            return new Integer[] { value, value + 1, value + 2 };
        }));
        // abonnementen
        ProcessUtils.subscribe(1, process2);
    }
}
  • regel 14: we gebruiken de methode [Observable.map];
  • regel 16: die een type Integer[] retourneert;

De resultaten zijn als volgt:

main : début observation ------Thread[main] ---- Time[58:089]
main : attente fin observation ------Thread[main] ---- Time[58:107]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[58:108]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[58:503]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[58:762]
Subscriber[observateur[0],process2] : onNext ([0,1,2]) ------Thread[RxComputationThreadPool-3] ---- Time[58:792]
Subscriber[observateur[0],process2] : onNext ([10,11,12]) ------Thread[RxComputationThreadPool-3] ---- Time[58:795]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[58:851]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[58:852]
Subscriber[observateur[0],process2] : onNext ([20,21,22]) ------Thread[RxComputationThreadPool-3] ---- Time[58:853]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[58:854]
main : fin observation ------Thread[main] ---- Time[58:854]
  • regels 6, 7, 10: hier zien we de resultaten van de map;

Al deze transformaties van de observable kunnen worden gekoppeld, aangezien elke transformatie een nieuwe observable oplevert. Dit wordt geïllustreerd in het volgende voorbeeld [Exemple21c]:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.schedulers.Schedulers;

public class Exemple21c {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMap(i -> {
            int value = i * 10;
            return Observable.just(value, value + 1, value + 2);
        }).filter(i -> i % 2 == 0));
        // abonnementen
        ProcessUtils.subscribe(1, process2);
    }
}
  • regels 15-18: de flatMap wordt gevolgd door een filter;

De uitvoerresultaten zijn als volgt:

main : début observation ------Thread[main] ---- Time[37:993]
main : attente fin observation ------Thread[main] ---- Time[38:016]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[38:017]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[38:124]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[38:366]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[38:380]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[38:381]
Subscriber[observateur[0],process2] : onNext (0) ------Thread[RxComputationThreadPool-3] ---- Time[38:436]
Subscriber[observateur[0],process2] : onNext (2) ------Thread[RxComputationThreadPool-3] ---- Time[38:439]
Subscriber[observateur[0],process2] : onNext (10) ------Thread[RxComputationThreadPool-3] ---- Time[38:441]
Subscriber[observateur[0],process2] : onNext (12) ------Thread[RxComputationThreadPool-3] ---- Time[38:443]
Subscriber[observateur[0],process2] : onNext (20) ------Thread[RxComputationThreadPool-3] ---- Time[38:445]
Subscriber[observateur[0],process2] : onNext (22) ------Thread[RxComputationThreadPool-3] ---- Time[38:446]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[38:447]
main : fin observation ------Thread[main] ---- Time[38:447]
  • regels 8-13: process2 heeft alleen de even elementen uit flatMap weergegeven;

Een methode die vergelijkbaar is met [flatMap] is de methode [flatMapIterable], geïllustreerd door het volgende voorbeeld [Exemple21d]:


package dvp.rxjava.observables.exemples;

import java.util.Arrays;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21d {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMapIterable(i -> {
            int value = i * 10;
            return Arrays.asList(value, value + 1, value + 2);
        }).filter(i -> i % 2 == 0));
        // abonnementen
        ProcessUtils.subscribe(1, process2);
    }
}

Regel 16: in plaats van de methode [flatMap] wordt de methode [flatMapIterable] gebruikt. In dit geval moet de transformatiefunctie een type Iterable<T> (regel 18) opleveren in plaats van een type Observable<T>.

We krijgen dezelfde resultaten als eerder.

Laten we teruggaan naar de definitie van de methode [flatMap]:

 

Hierboven zien we dat er een blauw element [3] is ingevoegd tussen de twee groene elementen [1-2]. Dit betekent dat de methode [flatMap] bij het samenvoegen van de Observable<T> de volgorde van uitzending van deze verschillende interne observabelen respecteert. Dit wordt geïllustreerd door het volgende voorbeeld [Exemple21e]:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21e {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // proces 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().flatMap(i -> process2.getObservable()));
        // abonnementen
        ProcessUtils.subscribe(1, process3);
    }
}
  • regels 11-12: het proces process1 genereert de gehele getallen [0,1];
  • regels 14-15: het proces process2 genereert de gehele getallen [10,11,12];
  • regels 17-18: aan elk element dat door process1 wordt uitgezonden, wordt de observabele van het proces process2 gekoppeld. Dit betekent dat:
    • aan het element [0] van proces1 wordt een observabel gekoppeld die de [10,11,12] uitzendt;
    • hetzelfde geldt voor element 1;

Uiteindelijk zullen de 6 getallen [10, 11, 12, 10, 11, 12] worden uitgezonden. We willen zien in welke volgorde.

De resultaten van de uitvoering zijn als volgt:

main : début observation ------Thread[main] ---- Time[22:540]
main : attente fin observation ------Thread[main] ---- Time[22:566]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[22:566]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[22:949]
Observable (process2) call start ------Thread[RxComputationThreadPool-6] ---- Time[22:951]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[23:159]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[23:160]
Observable (process2) call start ------Thread[RxComputationThreadPool-8] ---- Time[23:160]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-6] ---- Time[23:286]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-8] ---- Time[23:513]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[23:597]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[23:599]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-6] ---- Time[23:645]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-5] ---- Time[23:647]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-6] ---- Time[23:789]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[23:790]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-5] ---- Time[23:791]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-8] ---- Time[23:976]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-7] ---- Time[23:978]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-8] ---- Time[24:184]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-8] ---- Time[24:184]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-7] ---- Time[24:186]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-7] ---- Time[24:187]
main : fin observation ------Thread[main] ---- Time[24:187]

We zien dat de volgorde van uitzending van het proces process3 als volgt was: [10, 10, 11, 12, 11, 12] (regels 11, 12, 14, 17, 19, 22). Er is dus inderdaad een vermenging opgetreden van de elementen die door het proces process2 zijn verzonden. Dit kan worden voorkomen door de methode [concatMap] te gebruiken in plaats van de methode [flatMap]. Dit wordt geïllustreerd door de volgende code [Exemple21ef]:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21ef {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // proces 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().concatMap(i -> process2.getObservable()));
        // abonnementen
        ProcessUtils.subscribe(1, process3);
    }
}

Op regel 18 is [flatMap] vervangen door [concatMap]. De resultaten van de uitvoering zijn als volgt:

main : début observation ------Thread[main] ---- Time[45:507]
main : attente fin observation ------Thread[main] ---- Time[45:530]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[45:530]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[45:775]
Observable (process2) call start ------Thread[RxComputationThreadPool-6] ---- Time[45:778]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-6] ---- Time[45:846]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-6] ---- Time[45:890]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[45:947]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[45:948]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-6] ---- Time[46:096]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[46:097]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[46:144]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-5] ---- Time[46:147]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-5] ---- Time[46:148]
Observable (process2) call start ------Thread[RxComputationThreadPool-8] ---- Time[46:149]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-8] ---- Time[46:364]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-7] ---- Time[46:366]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-8] ---- Time[46:529]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-7] ---- Time[46:531]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-8] ---- Time[46:558]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-8] ---- Time[46:559]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-7] ---- Time[46:560]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-7] ---- Time[46:562]
main : fin observation ------Thread[main] ---- Time[46:562]

We zien dat de volgorde van uitvoering van het proces process3 als volgt was: [10, 11, 12, 10, 11, 12] (regels 12-14, 17, 19, 22). De elementen die door het proces process2 zijn gegenereerd, zijn niet gemengd.

Een andere variant van de methode [map] is de methode [switchMap]:

 

Hierboven ontstaan uit de observable [1] nog 3 andere observables [2] met 2 elementen, die vervolgens worden afgevlakt zoals in [flatMap] en [3]. Opvallend is dat het resultaat 5 elementen telt en geen 6. Dit komt doordat, voordat de tweede observable zijn element nr. 2 [6] uitzendt, de derde observable zijn eerste element [5] uitzendt, waardoor de tweede observable wordt genegeerd. Het element [6] komt dus niet voor in de resulterende observable [3].

Om [switchMap] te illustreren, gebruiken we het volgende voorbeeld [Exemple21eg]:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple21eg {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // proces 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().switchMap(i -> process2.getObservable()));
        // abonnementen
        ProcessUtils.subscribe(1, process3);
    }
}

De uitvoering van het voorbeeld levert de volgende resultaten op:

main : début observation ------Thread[main] ---- Time[02:388]
main : attente fin observation ------Thread[main] ---- Time[02:419]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[02:419]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[02:641]
Observable (process2) call start ------Thread[RxComputationThreadPool-6] ---- Time[02:643]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-6] ---- Time[02:802]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-6] ---- Time[02:888]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-6] ---- Time[02:957]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-6] ---- Time[02:958]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[03:005]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[03:007]
Observable (process2) call start ------Thread[RxComputationThreadPool-8] ---- Time[03:007]
Observable (process2,0) onNext (10) ------Thread[RxComputationThreadPool-8] ---- Time[03:106]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[03:106]
Subscriber[observateur[0],process3] : onNext (10) ------Thread[RxComputationThreadPool-5] ---- Time[03:108]
Observable (process2,1) onNext (11) ------Thread[RxComputationThreadPool-8] ---- Time[03:236]
Subscriber[observateur[0],process3] : onNext (11) ------Thread[RxComputationThreadPool-7] ---- Time[03:238]
Observable (process2,2) onNext (12) ------Thread[RxComputationThreadPool-8] ---- Time[03:716]
Observable (process2) onCompleted ------Thread[RxComputationThreadPool-8] ---- Time[03:717]
Subscriber[observateur[0],process3] : onNext (12) ------Thread[RxComputationThreadPool-7] ---- Time[03:718]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-7] ---- Time[03:718]
main : fin observation ------Thread[main] ---- Time[03:719]
  • process1 zendt 2 elementen uit die aanleiding geven tot 2 waarneembare process2 met 3 elementen;
  • regel 14: de waarnemer ontvangt element nr. 0, uitgezonden door de eerste waarnemelijke grootheid process2 op regel 6;
  • regel 15: de waarnemer ontvangt element nr. 0, uitgezonden door de tweede waarnemelijke grootheid process2 op regel 13. Het is onduidelijk waarom hij de elementen 1 en 2, verzonden door het eerste observeerbare object process2 op de regels 7 en 8, niet eerder heeft ontvangen. Feit is in ieder geval dat het eerste observeerbare object process2 wordt genegeerd;
  • uiteindelijk ziet de waarnemer slechts 4 elementen (regels 14, 15, 17, 20) in plaats van de 6 die zijn verzonden;

7.6.4. Voorbeelden-22: andere methoden van de klasse [Observable]

De klasse [Observable] bevat veel methoden uit de klasse [Stream] die op dezelfde manier werken. Hier volgen er enkele. We geven alleen de code en de resultaten weer.

[Exemple22a - take=limit]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22a {
    public static void main(String[] args) throws InterruptedException {
        // proces
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).take(3));
        // abonnementen
        ProcessUtils.subscribe(1, process);
    }
}

resultaten

1
2
3
4
5
6
7
main : début observation ------Thread[main] ---- Time[25:071]
Subscriber[observateur[0],process] : onNext (1) ------Thread[main] ---- Time[25:399]
Subscriber[observateur[0],process] : onNext (2) ------Thread[main] ---- Time[25:402]
Subscriber[observateur[0],process] : onNext (3) ------Thread[main] ---- Time[25:404]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[25:404]
main : attente fin observation ------Thread[main] ---- Time[25:406]
main : fin observation ------Thread[main] ---- Time[25:406]

[Exemple22b - takeLast]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22b {
    public static void main(String[] args) throws InterruptedException {
        // proces
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).takeLast(2));
        // inschrijvingen
        ProcessUtils.subscribe(1, process);
    }
}

resultaten

1
2
3
4
5
6
main : début observation ------Thread[main] ---- Time[19:440]
Subscriber[observateur[0],process] : onNext (9) ------Thread[main] ---- Time[19:726]
Subscriber[observateur[0],process] : onNext (10) ------Thread[main] ---- Time[19:728]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[19:728]
main : attente fin observation ------Thread[main] ---- Time[19:729]
main : fin observation ------Thread[main] ---- Time[19:730]

[Exemple22c - skip]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22c {
    public static void main(String[] args) throws InterruptedException {
        // processen
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).skip(5).take(2));
        // inschrijvingen
        ProcessUtils.subscribe(1, process);
    }
}

resultaten

1
2
3
4
5
6
main : début observation ------Thread[main] ---- Time[16:685]
Subscriber[observateur[0],process] : onNext (6) ------Thread[main] ---- Time[17:002]
Subscriber[observateur[0],process] : onNext (7) ------Thread[main] ---- Time[17:004]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[17:005]
main : attente fin observation ------Thread[main] ---- Time[17:006]
main : fin observation ------Thread[main] ---- Time[17:006]

[Exemple22d - reduce]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22d {
    public static void main(String[] args) throws InterruptedException {
        // processen
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).reduce(0, (i, a) -> i + a));
        // inschrijvingen
        ProcessUtils.subscribe(1, process);
    }
}
  • regel 10: berekent de som van de elementen van de observable. Het resultaat is een observable die deze som weergeeft;

resultaten

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[52:412]
Subscriber[observateur[0],process] : onNext (55) ------Thread[main] ---- Time[52:640]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[52:640]
main : attente fin observation ------Thread[main] ---- Time[52:642]
main : fin observation ------Thread[main] ---- Time[52:642]

[Exemple22e - all]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22e {
    public static void main(String[] args) throws InterruptedException {
        // processen
        Process<Boolean> process = new Process<>("process", Observable.range(1, 10).all(i -> i > 10));
        // inschrijvingen
        ProcessUtils.subscribe(1, process);
    }
}
  • regel 10: retourneert een Observable<Boolean> die het element true uitzendt, als het predikaat van de methode [all] voor alle elementen waar is, anders false;

resultaten

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[59:866]
Subscriber[observateur[0],process] : onNext (false) ------Thread[main] ---- Time[00:069]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[00:070]
main : attente fin observation ------Thread[main] ---- Time[00:071]
main : fin observation ------Thread[main] ---- Time[00:071]

[Exemple22f - count]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22f {
    public static void main(String[] args) throws InterruptedException {
        // processen
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).count());
        // inschrijvingen
        ProcessUtils.subscribe(1, process);
    }
}
  • regel 10: [Observable.count] creëert een observable met 1 element dat de som is van de geobserveerde elementen;

resultaten

1
2
3
4
5
main : début observation ------Thread[main] ---- Time[16:409]
Subscriber[observateur[0],process] : onNext (10) ------Thread[main] ---- Time[16:634]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[16:634]
main : attente fin observation ------Thread[main] ---- Time[16:635]
main : fin observation ------Thread[main] ---- Time[16:635]

[Exemple22g - distinct]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;

public class Exemple22g {
    public static void main(String[] args) throws InterruptedException {
        // processen
        Process<Integer> process = new Process<>("process", Observable.just(1, 2, 1, 3).distinct());
        // inschrijvingen
        ProcessUtils.subscribe(1, process);
    }
}

resultaten

1
2
3
4
5
6
7
main : début observation ------Thread[main] ---- Time[05:373]
Subscriber[observateur[0],process] : onNext (1) ------Thread[main] ---- Time[05:594]
Subscriber[observateur[0],process] : onNext (2) ------Thread[main] ---- Time[05:595]
Subscriber[observateur[0],process] : onNext (3) ------Thread[main] ---- Time[05:596]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[05:597]
main : attente fin observation ------Thread[main] ---- Time[05:597]
main : fin observation ------Thread[main] ---- Time[05:597]

[Exemple22h - groupBy, asObservable]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.observables.GroupedObservable;

public class Exemple22h {
    public static void main(String[] args) throws InterruptedException {
        // processen
        Observable<GroupedObservable<Boolean, Integer>> obs = Observable.range(1, 10).groupBy(i -> i % 2 == 0);
        Process<Integer> process = new Process<>("process", obs.concatMap(g -> g.asObservable()));
        // inschrijvingen
        ProcessUtils.subscribe(1, process);
    }
}
  • regel 11: de methode [groupBy] groepeert de 10 gegenereerde elementen in 2 groepen: de even getallen en de oneven getallen. Het resultaat is een type Observable<GroupedObservable<Boolean, Integer>>, d.w.z. een observable waarvan de elementen van het type GroupedObservable<Boolean, Integer> zijn, waarbij Boolean het type is van de sleutel van de groep (hier false, true) en dat tevens het type is van het resultaat van de lambda die als parameter wordt doorgegeven aan de methode [groupBy], en Integer het type van de elementen van de groep;
  • regel 12: het type GroupedObservable heeft een methode [asObservable] waarmee een observable kan worden aangemaakt op basis van dit type. We krijgen dus twee typen Observable<Integer>, één voor even getallen en één voor oneven getallen. Van deze twee observables zal de methode [concatMap] er één maken;

resultaten

main : début observation ------Thread[main] ---- Time[23:809]
Subscriber[observateur[0],process] : onNext (1) ------Thread[main] ---- Time[24:034]
Subscriber[observateur[0],process] : onNext (3) ------Thread[main] ---- Time[24:036]
Subscriber[observateur[0],process] : onNext (5) ------Thread[main] ---- Time[24:037]
Subscriber[observateur[0],process] : onNext (7) ------Thread[main] ---- Time[24:038]
Subscriber[observateur[0],process] : onNext (9) ------Thread[main] ---- Time[24:039]
Subscriber[observateur[0],process] : onNext (2) ------Thread[main] ---- Time[24:041]
Subscriber[observateur[0],process] : onNext (4) ------Thread[main] ---- Time[24:043]
Subscriber[observateur[0],process] : onNext (6) ------Thread[main] ---- Time[24:044]
Subscriber[observateur[0],process] : onNext (8) ------Thread[main] ---- Time[24:045]
Subscriber[observateur[0],process] : onNext (10) ------Thread[main] ---- Time[24:046]
Subscriber[observateur[0],process].onCompleted ------Thread[main] ---- Time[24:047]
main : attente fin observation ------Thread[main] ---- Time[24:047]
main : fin observation ------Thread[main] ---- Time[24:048]

[Exemple22i - timestamp]


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;
import rx.schedulers.Timestamped;

public class Exemple22i {
    public static void main(String[] args) throws InterruptedException {
        // proces 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // proces 2
        Process<Timestamped<Integer>> process2 = new Process<>("process2", process1.getObservable().timestamp());
        // inschrijvingen
        ProcessUtils.subscribe(1, process2);
    }
}
  • regel 15: de methode [timestamp] koppelt een tijdstip aan elk element van de verwerkte observabele;

resultaten

main : début observation ------Thread[main] ---- Time[59:362]
main : attente fin observation ------Thread[main] ---- Time[59:377]
Observable (process1) call start ------Thread[RxComputationThreadPool-4] ---- Time[59:378]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-4] ---- Time[59:553]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-4] ---- Time[59:692]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462975259555,"value":0}) ------Thread[RxComputationThreadPool-3] ---- Time[59:789]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462975259789,"value":1}) ------Thread[RxComputationThreadPool-3] ---- Time[59:791]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-4] ---- Time[00:025]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[00:027]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462975260026,"value":2}) ------Thread[RxComputationThreadPool-3] ---- Time[00:031]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-3] ---- Time[00:033]
main : fin observation ------Thread[main] ---- Time[00:034]

In dit voorbeeld is het moeilijk te zeggen wat de informatie timestamp voorstelt:

  • regels 4-5: we zien dat element 1 van process1 139 ms na element 0 is verzonden;
  • regels 6 en 7: we zien dat element 1 van process2 234 ms na element 0 is waargenomen;
  • regels 5 en 8: hieruit blijkt dat element 2 van process1 33 ms na element 1 is verzonden;
  • regels 7 en 10: hieruit blijkt dat element 2 van process2 37 ms na element 1 is waargenomen;

Deze vertragingen zijn te wijten aan het feit dat de threads voor het observeren en uitvoeren van de observables niet dezelfde zijn. Als we de regels 12-13 vervangen door de volgende regels (Voorbeeld22j):


// proces 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
null);
  • regels 2-3: de observatiethread wordt niet opgelegd. We weten dat in dit geval de observable wordt geobserveerd op de plaats waar deze wordt uitgevoerd;

Dit levert de volgende resultaten op:

main : début observation ------Thread[main] ---- Time[43:834]
main : attente fin observation ------Thread[main] ---- Time[43:845]
Observable (process1) call start ------Thread[RxComputationThreadPool-1] ---- Time[43:846]
Observable (process1,0) onNext (0) ------Thread[RxComputationThreadPool-1] ---- Time[44:291]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462976384293,"value":0}) ------Thread[RxComputationThreadPool-1] ---- Time[44:552]
Observable (process1,1) onNext (1) ------Thread[RxComputationThreadPool-1] ---- Time[44:878]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462976384879,"value":1}) ------Thread[RxComputationThreadPool-1] ---- Time[44:884]
Observable (process1,2) onNext (2) ------Thread[RxComputationThreadPool-1] ---- Time[45:274]
Subscriber[observateur[0],process2] : onNext ({"timestampMillis":1462976385275,"value":2}) ------Thread[RxComputationThreadPool-1] ---- Time[45:280]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[45:281]
Subscriber[observateur[0],process2].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[45:283]
main : fin observation ------Thread[main] ---- Time[45:284]
  • regels 4 en 6: het proces process1 zendt zijn element nr. 1 587 ms na zijn element nr. 0 uit;
  • regels 5 en 7: de waarnemer waarneemt deze twee elementen met een tussenpoos van 586 ms;
  • regels 6 en 8: het proces process1 verzendt zijn element nr. 2 396 ms na zijn element nr. 1;
  • regels 7 en 9: de waarnemer waarneemt deze twee elementen met een tijdsverschil van 396 ms;

Hier zijn de waarden van timestamp consistent: ze geven inderdaad de verzenddatum van het element weer.

7.7. De schedulers

7.7.1. Voorbeeld-23: de scheduler [Schedulers.computation]

We gaan nu de uitvoeringsplanners bekijken. De observatie vindt plaats op de uitvoeringsthread.

Het onderwerp van de schedulers is enigszins onduidelijk. De verschillende schedulers worden in deze vraag op de website van StackOverflow [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases] gepresenteerd:

 

We zullen proberen het gebruik van deze verschillende schedulers aan de hand van voorbeelden te illustreren. Het eerste voorbeeld illustreert de scheduler [Schedulers.computation]:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple23 {
    public static void main(String[] args) throws InterruptedException {
        // processen
        @SuppressWarnings("unchecked")
        Process<Double> processes[] = new Process[10];
        for (int i = 0; i < processes.length; i++) {
            processes[i] = new Process<>(
                    new ProcessAction01<Double>(String.format("process%s", i), 1, value -> new Random().nextInt(100) * 1.2),
                    Schedulers.computation(), null);
        }
        // abonnementen
        ProcessUtils.subscribe(1, processes);
    }
}
  • regels 14-19: er wordt een array aangemaakt met 10 processen die op een reken-thread worden uitgevoerd;
  • regel 17: elk proces genereert een willekeurig reëel getal;
  • regel 21: we abonneren ons op al deze processen;

De resultaten zijn als volgt:

main : début observation ------Thread[main] ---- Time[01:034]
Observable (process0) call start ------Thread[RxComputationThreadPool-1] ---- Time[01:042]
Observable (process2) call start ------Thread[RxComputationThreadPool-3] ---- Time[01:042]
Observable (process1) call start ------Thread[RxComputationThreadPool-2] ---- Time[01:042]
Observable (process5) call start ------Thread[RxComputationThreadPool-6] ---- Time[01:043]
Observable (process7) call start ------Thread[RxComputationThreadPool-8] ---- Time[01:043]
Observable (process4) call start ------Thread[RxComputationThreadPool-5] ---- Time[01:042]
Observable (process3) call start ------Thread[RxComputationThreadPool-4] ---- Time[01:042]
main : attente fin observation ------Thread[main] ---- Time[01:043]
Observable (process6) call start ------Thread[RxComputationThreadPool-7] ---- Time[01:043]
Observable (process3,0) onNext (70.8) ------Thread[RxComputationThreadPool-4] ---- Time[01:115]
Observable (process1,0) onNext (13.2) ------Thread[RxComputationThreadPool-2] ---- Time[01:153]
Observable (process0,0) onNext (63.599999999999994) ------Thread[RxComputationThreadPool-1] ---- Time[01:215]
Subscriber[observateur[0],process0] : onNext (63.599999999999994) ------Thread[RxComputationThreadPool-1] ---- Time[01:326]
Subscriber[observateur[0],process3] : onNext (70.8) ------Thread[RxComputationThreadPool-4] ---- Time[01:326]
Subscriber[observateur[0],process1] : onNext (13.2) ------Thread[RxComputationThreadPool-2] ---- Time[01:326]
Observable (process3) onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[01:326]
Observable (process0) onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[01:326]
Observable (process1) onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[01:327]
Subscriber[observateur[0],process0].onCompleted ------Thread[RxComputationThreadPool-1] ---- Time[01:327]
Subscriber[observateur[0],process3].onCompleted ------Thread[RxComputationThreadPool-4] ---- Time[01:327]
Subscriber[observateur[0],process1].onCompleted ------Thread[RxComputationThreadPool-2] ---- Time[01:327]
Observable (process8) call start ------Thread[RxComputationThreadPool-1] ---- Time[01:329]
Observable (process9) call start ------Thread[RxComputationThreadPool-2] ---- Time[01:329]
...
main : fin observation ------Thread[main] ---- Time[01:610]
  • regels 2-10: de eerste 8 processen starten op 8 verschillende threads (de gebruikte machine heeft 8 cores). Je kunt zien dat ze allemaal ongeveer tegelijkertijd starten;
  • regels 17-19: 3 processen worden beëindigd en maken zo 3 threads vrij;
  • regels 23-24: de laatste twee processen kunnen vervolgens starten door 2 van de vrijgekomen threads te gebruiken;

We onthouden dus dat de scheduler [Schedulers.computation] een pool van n threads ter beschikking stelt, waarbij n het aantal cores van de machine is. De threads worden parallel op deze cores uitgevoerd.

7.7.2. Voorbeeld 24: de scheduler [Schedulers.io]

We laten de vorige code uitvoeren met de scheduler [Schedulers.io]:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple24 {
    public static void main(String[] args) throws InterruptedException {
        // processen
        @SuppressWarnings("unchecked")
        Process<Double> processes[] = new Process[10];
        for (int i = 0; i < processes.length; i++) {
            processes[i] = new Process<>(
                    new ProcessAction01<Double>(String.format("process%s", i), 1, value -> new Random().nextInt(100) * 1.2),
                    Schedulers.io(), null);
        }
        // inschrijvingen
        ProcessUtils.subscribe(1, processes);
    }
}
  • regel 18: de processen worden uitgevoerd met de threads van de scheduler [Schedulers.io];

Dit levert de volgende resultaten op:

main : début observation ------Thread[main] ---- Time[03:451]
Observable (process0) call start ------Thread[RxCachedThreadScheduler-1] ---- Time[03:459]
Observable (process1) call start ------Thread[RxCachedThreadScheduler-2] ---- Time[03:459]
Observable (process2) call start ------Thread[RxCachedThreadScheduler-3] ---- Time[03:460]
Observable (process3) call start ------Thread[RxCachedThreadScheduler-4] ---- Time[03:460]
Observable (process4) call start ------Thread[RxCachedThreadScheduler-5] ---- Time[03:464]
Observable (process5) call start ------Thread[RxCachedThreadScheduler-6] ---- Time[03:464]
Observable (process6) call start ------Thread[RxCachedThreadScheduler-7] ---- Time[03:465]
Observable (process8) call start ------Thread[RxCachedThreadScheduler-9] ---- Time[03:465]
Observable (process9) call start ------Thread[RxCachedThreadScheduler-10] ---- Time[03:465]
main : attente fin observation ------Thread[main] ---- Time[03:465]
Observable (process7) call start ------Thread[RxCachedThreadScheduler-8] ---- Time[03:465]
Observable (process7,0) onNext (54.0) ------Thread[RxCachedThreadScheduler-8] ---- Time[03:473]
Observable (process8,0) onNext (116.39999999999999) ------Thread[RxCachedThreadScheduler-9] ---- Time[03:500]
Observable (process6,0) onNext (105.6) ------Thread[RxCachedThreadScheduler-7] ---- Time[03:506]
Observable (process0,0) onNext (96.0) ------Thread[RxCachedThreadScheduler-1] ---- Time[03:509]
Observable (process5,0) onNext (25.2) ------Thread[RxCachedThreadScheduler-6] ---- Time[03:583]
Observable (process3,0) onNext (97.2) ------Thread[RxCachedThreadScheduler-4] ---- Time[03:684]
Subscriber[observateur[0],process7] : onNext (54.0) ------Thread[RxCachedThreadScheduler-8] ---- Time[03:685]
Subscriber[observateur[0],process6] : onNext (105.6) ------Thread[RxCachedThreadScheduler-7] ---- Time[03:685]
Subscriber[observateur[0],process0] : onNext (96.0) ------Thread[RxCachedThreadScheduler-1] ---- Time[03:685]
Subscriber[observateur[0],process8] : onNext (116.39999999999999) ------Thread[RxCachedThreadScheduler-9] ---- Time[03:685]
Observable (process0) onCompleted ------Thread[RxCachedThreadScheduler-1] ---- Time[03:686]
Observable (process6) onCompleted ------Thread[RxCachedThreadScheduler-7] ---- Time[03:686]
Observable (process7) onCompleted ------Thread[RxCachedThreadScheduler-8] ---- Time[03:685]
...
main : fin observation ------Thread[main] ---- Time[03:933]
  • regels 2-10: de 10 processen starten elk op een andere thread. In tegenstelling tot het vorige geval konden alle processen worden gestart. We zien dat het starten 6 ms duurt, terwijl dat eerder 1 ms was;
  • regels 13-18: de observables zenden de ene na de andere uit en niet vrijwel parallel zoals eerder het geval was;

Wat is het verschil tussen de schedulers [Schedulers.io] en [Schedulers.computation]? Een antwoord is te vinden in URL [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases]:

 

7.7.3. Voorbeeld-25: de scheduler [Schedulers.newThread]

We voeren de vorige code uit met de planner [Schedulers.newThread]:


package dvp.rxjava.observables.exemples;

import java.util.Random;

import dvp.rxjava.observables.utils.Process;
import dvp.rxjava.observables.utils.ProcessAction01;
import dvp.rxjava.observables.utils.ProcessUtils;
import rx.schedulers.Schedulers;

public class Exemple25 {
    public static void main(String[] args) throws InterruptedException {
        // processen
        @SuppressWarnings("unchecked")
        Process<Double> processes[] = new Process[10];
        for (int i = 0; i < processes.length; i++) {
            processes[i] = new Process<>(
                    new ProcessAction01<Double>(String.format("process%s", i), 1, value -> new Random().nextInt(100) * 1.2),
                    Schedulers.newThread(), null);
        }
        // inschrijvingen
        ProcessUtils.subscribe(1, processes);
    }
}

De verkregen resultaten zijn dezelfde als bij de planner [Schedulers.io]:

main : début observation ------Thread[main] ---- Time[17:058]
Observable (process0) call start ------Thread[RxNewThreadScheduler-1] ---- Time[17:065]
Observable (process1) call start ------Thread[RxNewThreadScheduler-2] ---- Time[17:065]
Observable (process2) call start ------Thread[RxNewThreadScheduler-3] ---- Time[17:066]
Observable (process3) call start ------Thread[RxNewThreadScheduler-4] ---- Time[17:066]
Observable (process4) call start ------Thread[RxNewThreadScheduler-5] ---- Time[17:068]
Observable (process5) call start ------Thread[RxNewThreadScheduler-6] ---- Time[17:069]
Observable (process6) call start ------Thread[RxNewThreadScheduler-7] ---- Time[17:069]
Observable (process8) call start ------Thread[RxNewThreadScheduler-9] ---- Time[17:069]
Observable (process7) call start ------Thread[RxNewThreadScheduler-8] ---- Time[17:069]
Observable (process9) call start ------Thread[RxNewThreadScheduler-10] ---- Time[17:069]
main : attente fin observation ------Thread[main] ---- Time[17:069]
Observable (process6,0) onNext (25.2) ------Thread[RxNewThreadScheduler-7] ---- Time[17:120]
Observable (process3,0) onNext (39.6) ------Thread[RxNewThreadScheduler-4] ---- Time[17:193]
Observable (process5,0) onNext (21.599999999999998) ------Thread[RxNewThreadScheduler-6] ---- Time[17:212]
Observable (process0,0) onNext (19.2) ------Thread[RxNewThreadScheduler-1] ---- Time[17:273]
Observable (process8,0) onNext (81.6) ------Thread[RxNewThreadScheduler-9] ---- Time[17:308]
Subscriber[observateur[0],process3] : onNext (39.6) ------Thread[RxNewThreadScheduler-4] ---- Time[17:331]
Subscriber[observateur[0],process0] : onNext (19.2) ------Thread[RxNewThreadScheduler-1] ---- Time[17:331]
Subscriber[observateur[0],process6] : onNext (25.2) ------Thread[RxNewThreadScheduler-7] ---- Time[17:331]
Subscriber[observateur[0],process8] : onNext (81.6) ------Thread[RxNewThreadScheduler-9] ---- Time[17:331]
Subscriber[observateur[0],process5] : onNext (21.599999999999998) ------Thread[RxNewThreadScheduler-6] ---- Time[17:331]
Observable (process8) onCompleted ------Thread[RxNewThreadScheduler-9] ---- Time[17:333]
Observable (process5) onCompleted ------Thread[RxNewThreadScheduler-6] ---- Time[17:333]
Observable (process6) onCompleted ------Thread[RxNewThreadScheduler-7] ---- Time[17:332]
Observable (process0) onCompleted ------Thread[RxNewThreadScheduler-1] ---- Time[17:332]
Observable (process3) onCompleted ------Thread[RxNewThreadScheduler-4] ---- Time[17:332]
...
main : fin observation ------Thread[main] ---- Time[17:571]

In URL en [http://stackoverflow.com/questions/33415881/retrofit-with-rxjava-schedulers-newthread-vs-schedulers-io] wordt uitgelegd dat de scheduler [Schedulers.io] een threadpool ter beschikking stelt, wat de scheduler [Schedulers.newThread] niet doet. Een threadpool maakt automatisch een aantal n threads aan. Deze worden toegewezen aan de processen die ze nodig hebben. Wanneer deze processen zijn voltooid, worden hun threads niet verwijderd, maar keren ze terug naar de pool en kunnen ze vervolgens door een ander proces worden hergebruikt. Dit is efficiënter dan voortdurend threads aan te maken en te verwijderen. We kunnen dus aannemen dat het beter is om de scheduler [Schedulers.io] te gebruiken.

7.7.4. Voorbeeld 26: de schedulers [Schedulers.immediate, Schedulers.trampoline]

Laten we terugkomen op de uitleg die voor deze twee schedulers is gegeven:

 

De uitleg is vrij eenvoudig te begrijpen, maar als je het wilt illustreren, merk je dat je het niet helemaal begrepen hebt. Dankzij het boek [Learning Reactive Programming With Java 8] kon ik een voorbeeld maken dat is gebaseerd op een voorbeeld uit dit boek, maar dan vereenvoudigd. Het ziet er als volgt uit:


package dvp.rxjava.observables.exemples;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.function.Consumer;

import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Scheduler;
import rx.Scheduler.Worker;
import rx.functions.Action0;
import rx.schedulers.Schedulers;

public class Exemple26 {
    public static void main(String[] args) throws InterruptedException {

        // een planner
        Scheduler scheduler = Schedulers.immediate();
        // een worker van deze planner
        Worker worker = scheduler.createWorker();
        // een type Action0 dat op de worker moet worden uitgevoerd
        Action0 action02 = new Action0() {
            @Override
            public void call() {
                // logboek van actie02
                ProcessUtils.showInfos.accept("action02");
            }
        };

        // een Action0-type dat op de worker moet worden uitgevoerd
        Action0 action01 = new Action0() {
            @Override
            public void call() {
                // er wordt een nieuwe actie op dezelfde worker geprogrammeerd
                worker.schedule(action02);
                // logboek van actie01
                ProcessUtils.showInfos.accept("action01");
            }
        };
        // actie01 is ingepland op de worker
        worker.schedule(action01);
    }

    // weergaven
    static Consumer<String> showInfos = message -> System.out.printf("%s ------Thread[%s] ---- Time[%s]%n", message,
            Thread.currentThread().getName(), new SimpleDateFormat("ss:SSS").format(new Date()));

}
  • regel 17: een planner. Dit is ofwel [Schedulers.immediate] zoals hier, ofwel later [Schedulers.trampoline];
  • regel 19: er kunnen acties van het type Action0 (regels 21, 20) worden uitgevoerd op workers van de planner. Met de methode [Scheduler.createWorker] kan een worker worden aangemaakt. Met de methode [Worker.schedule(Action0)] kan een actie van het type Action0 door een worker worden uitgevoerd;
  • regels 21-27: een eerste actie met de naam [action02] die (regel 40) door de worker van regel 19 zal worden uitgevoerd;
  • regels 30-38: een tweede actie met de naam [action01]. Het bijzondere hieraan is dat de actie action02 op dezelfde worker wordt uitgevoerd als deze actie zelf (regel 34). Hierin ligt het verschil tussen [Schedulers.immediate] en [Schedulers.trampoline]:
    • als de scheduler [Schedulers.immediate] is, dan wordt in regel 34 de actie action02 onmiddellijk uitgevoerd (vandaar de naam van de scheduler) en wordt de lopende actie action01 onderbroken. Vervolgens verschijnt het bericht op regel 25. Zodra de actie action02 is voltooid, wordt de actie action01 hervat en verschijnt het bericht op regel 36;
    • als de planner [Schedulers.trampoline] is, dan wordt in regel 34 de actie action02 in de wachtrij geplaatst. Deze wordt pas uitgevoerd wanneer de lopende taak action01 is voltooid. Vervolgens verschijnt het bericht op regel 36. Zodra de actie action01 is voltooid, wordt de actie action02 uitgevoerd en zien we het bericht op regel 25;

De uitvoering van de bovenstaande code levert de volgende resultaten op:

action02 ------Thread[main] ---- Time[38:480]
action01 ------Thread[main] ---- Time[38:485]

Als we op regel 17 de planner [Schedulers.trampoline] gebruiken, krijgen we de omgekeerde resultaten:

action01 ------Thread[main] ---- Time[42:972]
action02 ------Thread[main] ---- Time[42:976]

Dat gezegd hebbende, is het moeilijk om een verband te leggen met de observables. Ik heb geen overtuigend voorbeeld gevonden dat het nut zou kunnen aantonen van het uitvoeren van een observable op een van deze twee threads. Hier is er echter wel een, maar ik vind het helemaal niet natuurlijk:


package dvp.rxjava.observables.exemples;

import dvp.rxjava.observables.utils.ProcessUtils;
import rx.Observable;
import rx.Scheduler.Worker;
import rx.functions.Action1;
import rx.schedulers.Schedulers;

public class Exemple27 {
    public static void main(String[] args) throws InterruptedException {

        // Worker
        Worker worker = Schedulers.immediate().createWorker();
        // Worker worker = Schedulers.trampoline().createWorker();
        // observable 1 op worker
        worker.schedule(() -> Observable.range(1, 2).subscribe(new Action1<Integer>() {

            @Override
            public void call(Integer i) {
                ProcessUtils.showInfos.accept(String.valueOf(i));
                // observable 2 op dezelfde worker
                worker.schedule(() -> Observable.range(100, 2).subscribe(new Action1<Integer>() {
                    @Override
                    public void call(Integer i) {
                        ProcessUtils.showInfos.accept(String.valueOf(i));
                    }
                }));
            }
        }));
    }
}
  • regels 13-14: er wordt een worker aangemaakt op basis van een van de twee schedulers [Schedulers.immediate] en [Schedulers.trampoline];
  • regel 16: een eerste observable obs1 wordt op deze worker ingepland om de getallen [1,2] te genereren
  • regel 22: telkens wanneer een element van deze observable obs1 wordt waargenomen, wordt de waarneming van een tweede observable obs2 op dezelfde worker gestart om de getallen [100,101] te genereren;

Met de scheduler [Schedulers.immediate] worden de volgende resultaten verkregen:

1
2
3
4
5
6
1 ------Thread[main] ---- Time[44:604]
100 ------Thread[main] ---- Time[44:610]
101 ------Thread[main] ---- Time[44:610]
2 ------Thread[main] ---- Time[44:612]
100 ------Thread[main] ---- Time[44:612]
101 ------Thread[main] ---- Time[44:612]

Terwijl met de scheduler [Schedulers.trampoline] de volgende resultaten worden verkregen:

1
2
3
4
5
6
1 ------Thread[main] ---- Time[14:107]
2 ------Thread[main] ---- Time[14:114]
100 ------Thread[main] ---- Time[14:115]
101 ------Thread[main] ---- Time[14:115]
100 ------Thread[main] ---- Time[14:115]
101 ------Thread[main] ---- Time[14:116]

7.8. Conclusion

Er valt nog veel te doen. Om de bibliotheek RxJava verder te verdiepen, wordt de lezer uitgenodigd om zijn opleiding voort te zetten aan de hand van de referenties die aan het begin van dit document zijn vermeld. Desondanks beschikken we over de basis om RxJava te gebruiken in de Swing- en Android-omgevingen. Dat gaan we nu laten zien.