Skip to content

7. La libreria RxJava

La libreria RxJava si basa sul seguente concetto: un flusso di elementi di tipo T Observable<T> viene osservato da uno o più sottoscrittori (abbonati, osservatori, consumatori) Subscriber<T>. La libreria RxJava consente al flusso Observable<T> di essere eseguito in un thread T1 e al suo osservatore Subscriber<T> in un thread T2 senza che lo sviluppatoredebba preoccuparsi di gestire il ciclo di vita di questi thread e di problemi naturalmente complessi, come la condivisione dei dati tra thread e la loro sincronizzazione per l’esecuzione di un’attività globale. Facilita quindi la programmazione asincrona.

Un flusso Observable<T> produce elementi di tipo T, osservabili man mano che vengono generati. Se l’osservatore e l’osservabile (termine che indica impropriamente il tipo Observable<T>) si trovano nello stesso thread, allora l’osservabile può produrre l’elemento (i+1) solo quando l’osservatore ha consumato l’elemento i. Sono pochi i casi in cui questa architettura risulta vantaggiosa. Se l’osservatore e l’osservabile non si trovano nello stesso thread, allora l’osservabile e il suo osservatore hanno comportamenti autonomi: l’osservabile produce al proprio ritmo e l’osservatore consuma al proprio ritmo. È proprio qui che risiede l’interesse della libreria. Finora abbiamo sempre parlato di un unico osservatore. In realtà, un osservabile può avere un numero qualsiasi di osservatori.

La libreria RxJava è particolarmente adatta all’architettura descritta nel paragrafo 2 della sezione introduttiva e che riportiamo qui di seguito:

Image

  • in [1], un livello di servizio fornisce servizi, alcuni dei quali richiedono molto tempo per essere ottenuti (ad esempio le richieste di rete);
  • questo livello di servizi viene richiamato da un'interfaccia grafica [1] (Swing, Android, JavaFx). Se il livello di servizio viene eseguito nello stesso thread del metodo [swing] che lo utilizza, l’interfaccia grafica rimane bloccata (non reattiva) durante l’attesa del risultato del servizio;
  • in [2], un sottile livello di adattamento implementato con RxJava consente di presentare al livello grafico un’implementazione asincrona dello stesso servizio: quest’ultimo può essere eseguito in un thread diverso da quello del metodo del livello grafico che lo invoca. In questo caso, l’interfaccia grafica [3] rimane reattiva: l’utente può continuare a interagire con essa, ad esempio avviando una nuova richiesta di rete in parallelo alla prima e, soprattutto, è possibile offrirgli la possibilità di annullare elaborazioni troppo lunghe, cosa impossibile se l’interfaccia grafica fosse bloccata;
  • la chiamata [4] è sincrona, mentre la chiamata [5-6] è asincrona;

In questa architettura, il livello [2] offre servizi che restituiscono tipi Observable<T> a cui i metodi del livello grafico [3] possono abbonarsi. Un servizio del livello [2] fornisce quindi i propri risultati uno alla volta e il livello [3] può reagire a ciascuno di essi, ad esempio aggiornando uno o più componenti dell’interfaccia grafica.

La classe Observable<T> dispone di diverse decine di metodi. Questa è una delle difficoltà della libreria: è molto ricca ed è difficile comprenderne tutte le possibilità. Ne presenteremo alcune. La padronanza degli altri metodi verrà poi col tempo.

7.1. Creare osservabili e iscriversi ad essi

7.1.1. Esempio-01: il metodo [Observable.from]

  

Consideriamo il seguente codice:


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) {
    // osservabili di numeri interi
    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");
      }
    });
  }
}
  • riga 12: si crea un tipo Observable<Integer> a partire da un elenco di interi.

La classe Observable<T> è un flusso di elementi di tipo T che possono essere osservati, preferibilmente in modo asincrono ma non necessariamente, man mano che vengono generati. La sua definizione è la seguente:

 

Come già detto, la classe Observable<T> dispone di diverse decine di metodi. Alcuni sono simili a quelli della classe Stream<T> esaminata nel paragrafo 5. La documentazione di RxJava include dei «marble diagrams» [2] che illustrano il funzionamento di questi metodi:

  • la riga 3 illustra le emissioni dell’osservabile nel corso del tempo;
  • il metodo [4] viene applicato agli elementi emessi dall’osservabile. In genere produce un nuovo osservabile;
  • la riga 5 mostra il nuovo osservabile ottenuto;

Il metodo [Observable.from] ha la seguente firma:

 

Il metodo statico [Observable.from] consente di creare un Observable<T> a partire da una collezione di elementi di tipo T. Si tratta di un modo molto semplice per iniziare a utilizzare gli osservabili. La riga:


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

emetterà quindi tre elementi. Non li emette immediatamente, ma li emetterà per intero ogni volta che si registrerà un osservatore. Questo è ciò che viene definito un osservabile freddo. L’osservabile riemette i propri elementi per ogni nuovo sottoscrittore.

Si può considerare l’istruzione precedente come un’azione di configurazione dell’osservabile. Quest’ultimo viene configurato una volta ed eseguito n volte se si presentano n iscritti.

Come ci si abbona?

Un modo per farlo è utilizzare il metodo [Observable.subscribe], la cui definizione qui utilizzata è la seguente:

 
  • il primo parametro [Action1<T> onNext] (cfr. paragrafo 6.2) del metodo è il metodo da eseguire quando l'osservabile emette un nuovo elemento T;
  • il secondo parametro [Action1<Throwable> onError] del metodo è il metodo da eseguire quando l’osservabile genera un’eccezione;
  • il terzo parametro [Action0 onComplete] (cfr. paragrafo 6.1) del metodo è il metodo da eseguire quando l'osservabile genera un'eccezione;
  • il metodo restituisce un tipo [Subscription];

Il tipo [Subscription] rappresenta un abbonamento all’osservabile. La sua definizione è la seguente:

 

L'interesse di questa interfaccia [1] risiede nel suo metodo [2], che consente di annullare un abbonamento.

Nel nostro esempio, il codice dell'abbonamento all'observable è il seguente:


    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");
      }
});
  • riga 1: il risultato di tipo [Subscription] viene ignorato;
  • righe 1-15: i tre parametri sono istanze di classi anonime. Utilizzeremo anche delle lambda. Il vantaggio delle classi anonime è che si vedono chiaramente i tipi di dati previsti dall’unico metodo di queste classi;
  • righe 2-5: implementazione del primo parametro di tipo [Action1<Integer>];
  • righe 6-10: implementazione del secondo parametro di tipo [Action1<Throwable>];
  • righe 11-15: implementazione del terzo parametro di tipo [Action0];

Il codice completo è il seguente:


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) {
    // osservabili di interi
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    // abbonamento
    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");
      }
    });
  }
}

L'osservabile della riga 12 inizia a emettere i suoi 3 elementi non appena viene chiamato il metodo [subscribe] alla riga 14. Da quel momento in poi:

  • ad ogni elemento emesso, vengono eseguite le righe 15-18.
  • al termine dei 3 elementi, vengono eseguite le righe 24-29;
  • le righe 19-24 non verranno mai eseguite poiché l'osservabile non genera un'eccezione in questo caso;

Per impostazione predefinita, l’observable e l’osservatore vengono eseguiti nello stesso thread. Esistono alcuni observable predefiniti che vengono eseguiti in un thread diverso dal thread principale (in questo caso il thread del metodo main), ma per la maggior parte di essi non è così. Qui, quindi, tutto avviene nel thread del metodo [main]:

  • l'observable emette l'elemento 1;
  • le righe 15-18 vengono eseguite e visualizzano questo elemento;
  • l'observable emette l'elemento 2;
  • le righe 15-18 vengono eseguite e visualizzano questo elemento;
  • l'observable emette l'elemento 3;
  • le righe 15-18 vengono eseguite e visualizzano questo elemento;
  • l'observable emette la notifica [completed];
  • vengono eseguite le righe 24-29;

Ecco cosa mostrano i risultati ottenuti:

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

La classe [Exemple02] riprende [Exemple01] utilizzando questa volta funzioni lambda come parametri del metodo [Observable.subscribe]:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;

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

7.1.2. Esempio-03: la classe Observer

  

Il metodo [Observable.subscribe], che consente di sottoscrivere un osservabile, presenta diverse versioni, tra cui la seguente:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;
import rx.Observer;

public class Exemple03 {
    public static void main(String[] args) {
        // osservabili di numeri interi
        Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
        // abbonamento
        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);
            }
        });
    };
}

Riga 13: invece di passare tre parametri al metodo [subscribe], gli viene passato un tipo [Observer] come segue:

 

Il tipo [Observer] è un'interfaccia con tre metodi:

  • [onNext(T t)], che viene chiamato ogni volta che l'osservabile emette un elemento t;
  • [onError(Throwable th)], che viene chiamato quando l’osservabile genera un’eccezione th;
  • [onCompleted], che viene chiamato quando l'osservabile segnala di aver terminato l'emissione;

Il funzionamento del codice è analogo a quello spiegato in precedenza. Si ottengono i seguenti risultati:

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

7.1.3. Esempio-04: il metodo [Observable.create]

  

Il metodo statico Observable.create è definito come segue:

 
  • il metodo [create] restituisce un tipo Observable<T>;
  • il parametro del metodo [create] è una funzione di tipo [Observable.OnSubscribe<T>] definita come segue:
 

Il tipo [Observable.OnSubscribe<T>] è un'interfaccia funzionale che a sua volta estende l'interfaccia funzionale [Action1<Subscriber<? super T>>]. Il metodo [call] di questa interfaccia richiede un tipo [Subscriber] (abbonato, sottoscrittore, osservatore) definito come segue:

 

Si può notare in [1] che la classe [Subscriber<T>] implementa l’interfaccia [Observer<T>] presentata al paragrafo 7.1.2.

Infine, il metodo [<T> Observable.create]:

  • accetta come parametro un'istanza di tipo [Observable.OnSubscribe<T>] con un unico metodo dalla firma: void call(Subscriber<T> s). Il tipo [Subscriber<T>] estende il tipo [Observer<T>] e dispone quindi dei metodi onNext, onError, onCompleted;
  • restituisce un tipo Observable<T>;

Il metodo [<T> Observable.create] restituisce un osservabile configurato. Non è stata ancora emessa alcuna element. Quando un sottoscrittore [Subscriber<T> s] si abbona a questo osservabile, viene chiamato il metodo [void call(s)] della funzione passata come parametro al metodo [<T> Observable.create]. Il suo ruolo è quello di emettere elementi t di tipo T e di chiamare il metodo [s.onNext(t)] dell’osservatore ad ogni emissione. Una volta terminata quest’ultima, deve essere chiamato il metodo [s.onCompleted(t)] dell’osservatore e il metodo [call] deve terminare. Se il metodo [call] incontra un'eccezione th, deve essere chiamato il metodo [s.onError(th)] dell'osservatore e il metodo [call] deve terminare;

Per illustrare questo complesso funzionamento, utilizzeremo il seguente codice [Exemple04]:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;

import java.util.Random;

public class Exemple04 {
    public static void main(String[] args) {
        // configurazione osservabile di numeri reali
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                for (int i = 0; i < 3; i++) {
                    // emissione dell'elemento i
                    subscriber.onNext(new Random((i + 1)).nextDouble());
                }
                // fine trasmissione
                subscriber.onCompleted();
            }
        });
        // sottoscrizione e quindi emissione
        obs1.subscribe((d) -> System.out.printf("onNext %s%n", d), (th) -> System.out.printf("onError %s%n", th),
                () -> System.out.println("onCompleted"));
    }
}
  • riga 11: si crea un osservabile che emette tipi Double;
  • righe 11-21: il parametro del metodo [create] viene istanziato con una classe anonima che presenta l’unico metodo [call] delle righe 12-20. L'osservabile creato alla riga 11 è pronto a emettere, ma emetterà solo quando arriverà un osservatore;
  • righe 13-21: il metodo [call] riceve il riferimento di un osservatore;
  • righe 14-17: invio di 3 elementi all’osservatore;
  • riga 19: notifica di fine trasmissione all’osservatore;
  • righe 23-24: sottoscrizione all’osservabile della riga 11. Si implementano i tre parametri [onNext, onError, onCompleted] del metodo [subscribe] tramite tre lambda. Questa sottoscrizione creerà l'abbonato [Subscriber<Double>] che verrà passato al metodo [call] della riga 13. A questo punto inizierà l'emissione degli elementi;
  • tutto avviene nello stesso thread: osservabile e osservatore;

Si ottengono i seguenti risultati:

1
2
3
4
onNext 0.7308781907032909
onNext 0.7311469360199058
onNext 0.731057369148862
onCompleted

Il metodo [Observable.create] consente di creare un osservabile a partire da qualsiasi fenomeno. È questo il metodo che abbiamo utilizzato nel paragrafo 2 della sezione «Scoperta», per trasformare un’interfaccia sincrona in un’interfaccia asincrona.

7.1.4. Esempio-05: refactoring di [Exemple-04]

  

L’esempio seguente presenta una nuova versione del metodo statico [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) {
        // configurazione di un osservabile di numeri reali
        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 attesa
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // errore
                        subscriber.onError(e);
                    }
                    // azione
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // fine
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // un sottoscrittore
        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));
            }
        };

        // sottoscrizione
        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()));
    }
}
  • riga 56: la nuova versione del metodo statico [Observable.subscribe] accetta come parametro il tipo [Subscriber] che abbiamo presentato nel paragrafo precedente;
  • righe 37-52: il sottoscrittore (abbonato, osservatore). Esso implementa l'interfaccia Observer con i suoi tre metodi onNext, onError, onCompleted;
  • righe 61-64: d'ora in poi ci concentreremo sui thread in cui vengono eseguiti l'osservabile e il suo osservatore;
  • riga 62: il nome del thread;
  • riga 63: l'ora corrente espressa in secondi e millisecondi. Questo ci permetterà di vedere nel tempo l'emissione di elementi da parte dell'observable e la loro elaborazione da parte dell'osservatore;
  • questo codice ha la stessa funzionalità di quello precedente. Abbiamo semplicemente rifattorizzato quest'ultimo;

I risultati ottenuti sono i seguenti:

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]
  • riga 1 dei risultati: prima della riga 56 del codice, non è ancora successo nulla. L'osservabile è stato semplicemente configurato;
  • riga 2 dei risultati: la riga 56 del codice provoca la chiamata del metodo [call] della riga 15. Riga 3: il numero reale 80,39 viene inviato all'osservatore;
  • riga 4: l’osservatore riceve il numero inviato;
  • righe 5-8: il processo precedente si ripete due volte;
  • riga 9: l’osservabile invia la notifica di fine trasmissione;
  • riga 10: l'osservatore la riceve;
  • riga 11: visualizzata dalla riga 57 del codice;

Si nota quindi che la sola riga 56 di sottoscrizione ha provocato la visualizzazione delle righe 2-10 dei risultati. Quando si inizia a utilizzare la libreria RxJava ci si chiede come le cose si concatenino tra loro e, in particolare, quali siano i collegamenti che legano l’osservatore e l’osservabile. Si vede qui che la riga 56, la sottoscrizione all’osservabile,

  • ha provocato l'emissione di tutti gli elementi dell'osservabile;
  • che l’osservabile e l’osservatore vengono eseguiti nello stesso thread;
  • che, per questo motivo, si osserva la sequenza: emissione dell’elemento i, osservazione dell’elemento i, emissione dell’elemento (i+1), osservazione dell’elemento (i+1), ...

Ricordiamo che l’emittente attendeva prima di emettere i propri elementi:


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

dove i alla riga 3 rappresenta il numero di emissione (0<=i<3). Se si osservano gli orari di emissione degli elementi dell'osservabile:

  • righe 2, 3: l'elemento 0 è stato trasmesso circa 500 ms dopo l'inizio dell'abbonamento;
  • righe 3, 5: l'elemento 1 è stato trasmesso circa 400 ms dopo l'elemento 0;
  • righe 5, 7: l'elemento 2 è stato emesso circa 300 ms dopo l'elemento 1;

7.2. Thread di esecuzione, thread di osservazione

7.2.1. Esempio-06: osservabile e osservatore in un thread diverso da [main]

  

Rifattorizziamo l'esempio precedente nel modo seguente [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) {

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

        // configurazione di un osservabile di valori reali
        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 attesa
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // errore
                        subscriber.onError(e);
                    }
                    // azione
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // fine
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // un sottoscrittore
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // si abbassa la barriera
                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));
            }
        };

        // configurazione osservabile successiva
        obs1 = obs1.subscribeOn(Schedulers.computation());
        // sottoscrizione
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // attesa davanti alla barriera
        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()));
    }
}
  • riga 16: si crea un semaforo con un oggetto di tipo [CountDownLatch]. Questo oggetto serve a sincronizzare i thread tra loro. Qui viene inizializzato con il valore 1, che chiameremo valore del semaforo. Un thread si mette in attesa del semaforo tramite un'operazione:

latch.await();

Il thread rimane bloccato se il valore del guardrail è >0. Un thread può aumentare o diminuire il valore interno del guardrail. Alla riga 48, il valore del guardrail viene decrementato di 1.

  • riga 63: l’osservabile è configurato in modo da essere eseguito su un thread fornito dallo scheduler [Schedulers.computation()]. Questo scheduler può fornire un numero di thread pari al numero di core presenti sulla macchina di esecuzione. Il paragrafo dedicato all’applicazione di esempio ha illustrato l’utilizzo di altri scheduler (cfr. paragrafo 2.8);

Il principio del codice è il seguente:

  • il metodo [main] viene eseguito nel thread principale (main);
  • riga 66: avvia l’emissione degli elementi dell’osservabile. Questi saranno emessi su un thread diverso dal thread principale;
  • riga 70: il thread principale viene bloccato poiché il valore del gatekeeper è 1 (cfr. riga 16). Potrà proseguire solo quando tale valore passerà a 0. Ciò avviene alla riga 48. È l’osservatore che abbassa il gatekeeper quando riceve la notifica che l’osservabile ha terminato le emissioni;

L'esecuzione fornisce i seguenti risultati:

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]
  • riga 1: la sottoscrizione sta per avvenire;
  • riga 2: questa innesca l’esecuzione del metodo [call] sul thread [RxComputationThreadPool-1]. Ora abbiamo un’esecuzione parallela con due thread;
  • riga 3: per un motivo non chiarito, il thread [RxComputationThreadPool-1] ha ceduto il controllo. Il thread [main] ne prende quindi il controllo e viene bloccato dal guardrail (riga 70 del codice). Da questo momento in poi, solo il thread [RxComputationThreadPool-1] può operare;
  • righe 4-11: si osserva il comportamento già visto in precedenza tra l’osservabile e il suo osservatore, ma ora tutto avviene nel thread [RxComputationThreadPool-1];
  • righe 12-13: l’osservatore ha abbassato la barriera (riga 48 del codice) e il thread [RxComputationThreadPool-1] si è terminato. Il thread [main] prende il controllo e visualizza due messaggi;

7.2.2. Esempio-07: osservabile e osservatore in due thread diversi

  

Modifichiamo l’esempio precedente nel modo seguente:


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

        // guardiano della barriera
        CountDownLatch latch = new CountDownLatch(1);

        // configurazione di un osservabile di numeri reali
        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++) {
                    // attesa
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // errore
                        subscriber.onError(e);
                    }
                    // azione
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // fine
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // un sottoscrittore
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // si abbassa la barriera
                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));
            }
        };

        // configurazione osservabile successiva
        obs1 = obs1.subscribeOn(Schedulers.computation()).observeOn(Schedulers.computation());
        // sottoscrizione
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // in attesa che la barriera si alzi
        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()));
    }
}

Il codice è identico a quello dell'esempio precedente, tranne che per la riga 63:


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

che configura l'osservabile (subscribeOn) e l'osservatore (observeOn) affinché vengano eseguiti su uno dei thread forniti dallo scheduler [Schedulers.computation()].

I risultati ottenuti sono i seguenti:

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]

Si possono notare i seguenti punti:

  • l'osservabile viene eseguito nel thread [RxComputationThreadPool-4] (righe 3-4, 6, 8-9);
  • l’osservatore viene eseguito nel thread [RxComputationThreadPool-3] (righe 5, 7, 10-11);
  • che entrambi si eseguono in modo autonomo. Pertanto, alle righe 8-9, l’osservabile emette 2 notifiche (onNext, onCompleted) prima che l’osservatore recuperi la notifica [onNext] (riga 10);

La libreria RxJava si occupa del trasferimento dei dati (le emissioni) dal thread dell’osservabile a quello dell’osservatore. Lo sviluppatore non deve preoccuparsene.

Abbiamo visto come creare osservabili (Observable.from, Observable.create). Vediamo ora gli osservabili predefiniti della libreria RxJava.

7.3. Observable predefiniti

7.3.1. Esempio-08: il metodo [Observable.range]

 

D'ora in poi, utilizzeremo classi dedicate per i processi osservati e i relativi osservatori. L'idea è quella di poter registrare il loro nome, il thread di esecuzione e gli orari di esecuzione, in modo da poterli monitorare nel tempo.

La classe [Process] sarà semplicemente un Observable a cui è possibile assegnare un nome. Implementerà la seguente interfaccia [IProcess]:


package dvp.rxjava.observables.utils;

import rx.Observable;

public interface IProcess<T> {

    // nome dell'osservabile
    public String getName();

    // osservabile
    public Observable<T> getObservable();

}

Questa interfaccia potrà essere implementata dalla seguente classe [Process<T>]:


package dvp.rxjava.observables.utils;

import rx.Observable;
import rx.Scheduler;

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

    // nome dell'osservabile
    protected String name;
    // processo osservato
    protected Observable<T> observable;

    // costruttori
    public Process(String name, Observable<T> observable) {
        // inizializzazioni locali
        this.name = name;
        this.observable = observable;
    }

    // getter e setter
    public String getName() {
        return name;
    }

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

}
  • riga 9: il nome del processo;
  • riga 11: l'osservabile osservato;
  • righe 14-18: il costruttore;

L'osservatore sarà invece descritto dalla seguente classe [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> {

...
}
  • riga 11, la classe Observateur<T> estende la classe Subscriber<T> che abbiamo presentato brevemente nel paragrafo 7.1.3. La useremo come argomento del metodo [Observable.subscribe]:

// esecuzione osservabile (osservazione)
obs1.subscribe(observateur);

Il metodo [Observable.subscribe] utilizzato alla riga 2 sopra riportata ha la seguente definizione:

 

Il ruolo del [Subscriber] consiste principalmente nel gestire gli elementi emessi dall’osservabile a cui si è abbonato tramite i metodi dell’interfaccia [Observer]: onNext, onError, onCompleted. La classe [Subscriber] dispone dei seguenti metodi:

 

Nel codice della classe [Observateur], utilizzeremo il metodo [1] isUnsubscribed per verificare se l'abbonamento dell'abbonato è stato annullato o meno. Il codice completo della classe [Observateur<T>] è il seguente:


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

    // un semaforo
    private CountDownLatch latch;
    // un metodo di visualizzazione
    private Consumer<String> showInfos;
    // il nome dell'osservatore
    private String observerName;
    // il nome del processo osservato
    private String processName;

    // costruttori
    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;
    }

    // --------------------------- implementazione dell'interfaccia Observer<T>
    @Override
    public void onCompleted() {
        // fine delle trasmissioni
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber [%s,%s].onCompleted", observerName, processName));
        }
        // fine blocco del thread principale
        latch.countDown();
    }

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

    @Override
    public void onNext(T value) {
        // un'emissione aggiuntiva
        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));
            }
        }
    }
}
  • oltre alle caratteristiche di un Subscriber, l’osservatore Observateur includerà le seguenti informazioni:
    • riga 14: un guardrail o semaforo che servirà a bloccare il thread principale fino a quando l’osservatore non avrà ricevuto tutti gli elementi emessi dall’osservabile. Ciò avverrà alla riga 36 del codice quando l’osservatore riceverà dall’osservabile la notifica di fine emissione;
    • riga 16: un'istanza Consumer<String> che servirà a visualizzare un messaggio sulla console;
    • riga 18: il nome dell'osservatore per distinguerli l'uno dall'altro quando ce ne sono più di uno;
    • riga 20: il nome del processo osservato;
  • righe 36, 46, 54: i metodi [onCompleted, onError, onNext] dell'interfaccia [Observer<T>] implementata dalla classe astratta [Subscriber<T>]. Questa classe non li implementa. È quindi necessario farlo nelle classi figlie. Prima di eseguire qualsiasi operazione in questi metodi, si verifica se l’osservatore non sia stato disabbonato dall’osservabile che sta osservando;
  • riga 59: il metodo [onNext] dell’osservatore scrive la stringa jSON dell’elemento ricevuto. Ciò ci consentirà di visualizzare vari tipi di elementi;

Detto questo, esaminiamo un nuovo metodo della classe Observable, il metodo [range]:

 

L’osservabile Observable.range(n,m) emette (m) numeri interi compresi tra n e n+m-1. Lo analizziamo con il seguente codice [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 {

        // numero di osservatori
        final int nbObservateurs = 2;

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

        // configurazione osservabile
        Observable<Integer> obs1 = Observable.range(15, 3).subscribeOn(Schedulers.computation());
        // esecuzione osservabile (osservazione)
        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"));
        }
        // attesa
        showInfos.accept("main : attente fin observation");
        latch.await();
        // fine
        showInfos.accept("main : fin observation");
    }

    // visualizzazioni
    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()));
}
  • riga 16: useremo due osservatori;
  • riga 19: il semaforo viene inizializzato a due perché ogni osservatore verrà collocato su un thread diverso. Il thread principale dovrà quindi attendere il completamento di entrambi i thread di osservazione;
  • riga 22: si configura l’osservabile in modo tale che venga eseguito su un thread dello scheduler [Schedulers.computation()]. L’osservatore si troverà sullo stesso thread dell’osservabile;
  • righe 25-27: si sottoscrivono due osservatori all’osservabile. Ciò innescherà l’esecuzione completa di quest’ultimo per ciascuno degli osservatori: verranno emessi i numeri interi 15, 16 e 17;
  • riga 30: il thread principale attende il completamento degli osservatori;

I risultati ottenuti sono i seguenti:

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]
  • riga 2: il thread principale è bloccato in attesa che i 2 osservatori terminino;
  • righe 3-4: si vede che l'osservatore 0 si trova sul thread [RxComputationThreadPool-1] e l'osservatore 1 sul thread [RxComputationThreadPool-2];
  • righe 3-10: si vede che entrambi gli osservatori ricevono esattamente gli stessi elementi;

Utilizzeremo la classe Observateur così definita per illustrare il comportamento di altri tipi di osservabili.

7.3.2. Esempio-09: i metodi Observable.[interval, take, doNext]

  
 

Questo esempio illustra l’utilizzo dell’osservabile Observable.interval (intervallo lungo, unità TimeUnit) che emette numeri interi lunghi a intervalli di tempo regolari. Da notare il punto [1]: per impostazione predefinita, l'osservabile [Observable.interval] viene eseguito su uno dei thread dello scheduler [Schedulers.computation].

Il codice sarà il seguente:


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 {

        // numero di osservatori
        final int nbObservateurs = 2;

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

        // configurazione osservabile
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // esecuzione osservabile (osservazione)
        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"));
        }
        // attesa
        showInfos.accept("main : attente fin observation");
        latch.await();
        // fine
        showInfos.accept("main : fin observation");
    }

    // visualizzazioni
    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()));
}
  • riga 22: l'osservabile emette numeri interi lunghi ogni 500 millisecondi. La serie inizia con il numero 0;
  • riga 22: questo osservabile emette un numero infinito di valori. Il metodo [Observable.take(n)] crea un nuovo osservabile che conserva solo i primi n elementi emessi;
 

Torniamo al codice dell’osservabile:


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

Riga 2: il metodo [Observable.doOnNext] viene eseguito ogni volta che l'osservabile emette un nuovo elemento. Viene spesso utilizzato per registrare informazioni nel log. In questo caso, si desidera registrare la data di emissione degli elementi per verificare se l'intervallo di 500 millisecondi viene effettivamente rispettato. Il metodo [Observable.doOnNext] non modifica l’osservabile a cui si applica. La sua definizione è la seguente:

 

L'esecuzione fornisce i seguenti risultati:

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]
  • righe 3, 7 e 11: si nota che l'intervallo di emissione è approssimativamente vicino a 500 ms;
  • I due osservatori si trovano ovviamente su due thread diversi, nonostante l’osservabile non fosse stato configurato per essere eseguito con uno scheduler specifico. Quello che vediamo qui è il funzionamento predefinito dell’osservabile [Observable.interval];

7.3.3. Esempi-10/12: i metodi Observable.[error, empty, never]

 

D'ora in poi saremo più concisi nelle nostre illustrazioni dei metodi della classe [Observable]. Il codice precedente era il seguente:


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 {

        // numero di osservatori
        final int nbObservateurs = 2;

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

        // configurazione osservabile
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // esecuzione osservabile (osservazione)
        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"));
        }
        // attesa
        showInfos.accept("main : attente fin observation");
        latch.await();
        // fine
        showInfos.accept("main : fin observation");
    }

    // visualizzazioni
    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()));
}

Questo codice era già stato utilizzato per l’esempio precedente. Cambiavano solo le righe 21-22. Factorizzeremo quindi la maggior parte di questo codice nella seguente classe [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 {

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

        // esecuzione osservabile (osservazione)
        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()));
            }
        }
        // attesa
        showInfos.accept("main : attente fin observation");
        latch.await();
        // fine
        showInfos.accept("main : fin observation");
    }

    // visualizzazioni
    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()));
}
  • riga 13: il metodo accetta due parametri:
    • nbObservateurs: il numero di osservatori dei processi passati come secondo parametro;
    • processes: i processi (osservabili denominati) da osservare. Grazie alla notazione [IProcess<?>], i processi potranno emettere elementi di tipi diversi;
  • riga 16: il semaforo deve passare al verde quando tutti gli osservatori hanno completato tutte le loro osservazioni. Il valore iniziale del semaforo è quindi pari al numero di osservatori moltiplicato per il numero di osservazioni;
  • righe 20-25: si abbona ogni osservatore a tutti i processi che deve osservare;
  • riga 23: si recupera l’osservabile dal processo (cfr. paragrafo 7.3.1);
  • riga 23: si associa un osservatore all'osservabile. A quest'ultimo vengono trasmesse 4 informazioni:
    • il suo nome;
    • il semaforo che deve decrementare quando riceve la notifica di fine trasmissione dell'osservabile che sta osservando;
    • il metodo da utilizzare quando desidera registrare informazioni sulla console;
    • il nome del processo che osserverà;

Una volta definite queste classi, l’esempio 10 sarà il seguente:


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 {
        // configurazione osservabile
        Observable<?> obs = Observable.error(new RuntimeException("Erreur !!!")).subscribeOn(Schedulers.computation());
        // esecuzione (osservazione) osservabile
        ProcessUtils.subscribe(2,new Process<>("process1", obs));
    }
}

Riga 11, il metodo statico [Observable.error] è definito come segue:

 

La riga 8 configura quindi un osservabile che si limita a generare un'eccezione destinata al metodo [onError] dei suoi sottoscrittori. L'esecuzione fornisce i seguenti risultati:


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]

Righe 3 e 4: il metodo [onError] di entrambi gli osservatori ha ricevuto l’eccezione generata dall’osservabile.

Questa esecuzione presenta una particolarità: i metodi [onCompleted] di entrambi gli osservatori non sono stati chiamati. Di conseguenza, la barriera non è stata abbassata e il thread principale rimane bloccato nel metodo statico [ProcessUtils.subscribe] alla riga 3 seguente:


// attesa
showInfos.accept("main : attente fin observation");
latch.await();
// fine
showInfos.accept("main : fin observation");

Qui si nota che, in caso di errore dell’osservabile, il metodo [onCompleted] dei sottoscrittori non viene chiamato. Modifichiamo quindi il metodo [Observateur.onError] nel modo seguente:


    @Override
    public void onError(Throwable e) {
        // errore di trasmissione
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber[%s, %s].onError (%s)", observerName, processName, e));
        }
        // fine blocco thread principale
        latch.countDown();
}

Aggiungiamo le righe 7-8 per rimuovere il blocco in caso di errore dell’osservabile. Con questo nuovo codice, l’esecuzione fornisce i seguenti risultati:


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]

Otteniamo la riga 5 che prima non avevamo.

L’esempio 11 sarà il seguente:


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 {
        // configurazione osservabile
        Observable<?> obs1 = Observable.empty();
        // esecuzione (osservazione) osservabile
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

Riga 10: il metodo statico [Observable.empty] crea un osservabile che non emette alcun elemento. Emette solo la notifica di fine emissione;

 

L’esecuzione del codice dell’esempio sopra riportato fornisce i seguenti risultati:

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]
  • righe 2 e 3: si nota che entrambi gli osservatori ricevono la notifica di fine emissione senza aver ricevuto alcun elemento in precedenza.

Ci si potrebbe chiedere a cosa serva questo metodo. È possibile utilizzarlo in modo analogo a una collezione, inizialmente vuota, nella quale vengono poi aggiunti degli elementi:

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

Alla riga 3, si unisce l’osservabile iniziale obs (riga 1) con altri osservabili.

L’esempio 12 illustra il metodo statico [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 {
        // configurazione osservabile
        Observable<?> obs1 = Observable.never();
        // esecuzione (osservazione) osservabile
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

Il metodo statico [Observable.never] crea un osservabile che non emette mai:

 

L'esecuzione dell'esempio produce i seguenti risultati:

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

Riga 2: il thread principale rimane in attesa a tempo indeterminato. Infatti, nessun osservabile emette la notifica [onCompleted] che consente di far passare il semaforo (barriera) al verde (abbassare la barriera).

7.4. Multi-threading

7.4.1. Esempio 13: thread di azione, thread di osservazione

Nel paragrafo 7.1.3 abbiamo creato un osservabile con il metodo statico [Observable.create]:

 
  • il metodo [create] restituisce un tipo Observable<T>;
  • il parametro del metodo [create] è una funzione di tipo [Observable.OnSubscribe<T>] definita come segue:
 

Il tipo [Observable.OnSubscribe<T>] è un'interfaccia funzionale che a sua volta estende l'interfaccia funzionale [Action1<Subscriber<? super T>>]. Il metodo [call] di questa interfaccia richiede un tipo [Subscriber] (abbonato, sottoscrittore, osservatore). Nel prosieguo di questo documento, a volte chiameremo il tipo [Observable.OnSubscribe<T>] «azione». Creeremo delle azioni personalizzate a cui assegneremo un nome. Si tratterà di istanze della seguente interfaccia [IProcessAction]:

  

package dvp.rxjava.observables.utils;

import rx.Observable;

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

    // l'azione ha un nome
    public String getName();
}
  • riga 5: l’interfaccia [IProcessAction<T>] presenta tutte le caratteristiche dell’interfaccia [Observable.OnSubscribe<T>];
  • riga 8: dispone inoltre di un metodo [getName] che restituisce il nome dell'istanza che implementa l'interfaccia;

Utilizzeremo la seguente azione denominata [ProcessAction01]:


package dvp.rxjava.observables.utils;

import java.util.Random;

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

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

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

    // costruttori
    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++) {
            // attesa
            try {
                Thread.sleep(new Random().nextInt(500));
            } catch (InterruptedException e) {
                // errore
                ProcessUtils.showInfos.accept(String.format("Observable (%s) onError", getName()));
                subscriber.onError(e);
            }
            // invio di un elemento
            T value = func1.call(i);
            ProcessUtils.showInfos.accept(String.format("Observable (%s,%s) onNext (%s)", getName(), i, value));
            subscriber.onNext(value);
        }
        // fine
        ProcessUtils.showInfos.accept(String.format("Observable (%s) onCompleted", getName()));
        subscriber.onCompleted();
    }

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

}
  • riga 8: la classe [ProcessAction01<T>] implementa l'interfaccia [IProcessAction<T>] e quindi l'interfaccia [Observable.OnSubscribe<T>];
  • riga 11: il nome dell'azione;
  • riga 12: il numero di valori da emettere;
  • riga 13: un'istanza di tipo [Func1<Integer, T>] che, a partire da un numero intero, crea un tipo T che verrà emesso dall'osservabile (righe 35 e 37);
  • righe 16-20: al costruttore vengono passati il nome dell'azione, il numero di valori da emettere e la funzione di emissione;
  • righe 23-42: il codice del processo;
  • riga 23: il metodo [call] riceve come parametro l'abbonato all'osservabile associato al processo;
  • riga 28: il processo emette i propri elementi dopo un'attesa di durata casuale;
  • riga 32: l'emissione di un errore;
  • riga 37: un'emissione normale;
  • riga 41: emissione della notifica di fine emissione;
  • righe 25-38: l'azione emette valori reali nbValues dopo un tempo di attesa casuale (riga 30);
  • riga 35: il valore da emettere è fornito dalla funzione [func1] passata come parametro al costruttore (riga 16);

Rifattorizziamo la classe [Process] (cfr. paragrafo 7.3.1) in modo che possa essere istanziata anche con un'azione denominata. Aggiungiamo il seguente costruttore:


public Process(IProcessAction<T> na, Scheduler schedulerObserved, Scheduler schedulerObserver) {
        // nome processo=nome azione
        name = na.getName();
        // azione --> osservabile
        observable = Observable.create(na);
        // thread di esecuzione del processo osservato
        if (schedulerObserved != null) {
            observable = observable.subscribeOn(schedulerObserved);
        }
        // thread di osservazione dell'osservatore
        if (schedulerObserver != null) {
            observable = observable.observeOn(schedulerObserver);
        }
    }
  • alla riga 1, il costruttore accetta 3 parametri:
    1. l'azione denominata che servirà a costruire l'osservabile (riga 5);
    2. lo scheduler del processo osservato (può essere null);
    3. lo scheduler dell’osservatore (ad esempio null);
  • riga 5: l'osservabile viene creato a partire dall'azione passata come parametro;

Il codice seguente [Exemple13] osserva diversi osservabili:


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 {
        // processo 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 1, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        // processo 3
        Process<Integer> process3 = new Process<>(new ProcessAction01<Integer>("process3", 3, i -> i * 2), null,
                Schedulers.computation());
        // processo 4
        Process<Boolean> process4 = new Process<>(new ProcessAction01<Boolean>("process4", 4, i -> i % 2 == 0), null, null);
        // iscrizioni
        ProcessUtils.subscribe(1, process1);
        ProcessUtils.subscribe(1, process2);
        ProcessUtils.subscribe(1, process3);
        ProcessUtils.subscribe(1, process4);
    }
}
  • righe 13-15: il processo process1 produce un numero reale su un thread di calcolo che verrà osservato su un altro thread di calcolo;
  • righe 17-18: il processo process2 produce 2 stringhe di caratteri su un thread di calcolo e non viene fornita alcuna indicazione sul thread dell'osservatore. I risultati mostrano che l'osservazione avviene per impostazione predefinita sullo stesso thread di esecuzione del processo;
  • righe 20-21: il processo process3 genera 3 numeri interi su un thread non specificato che saranno osservati su un thread di calcolo. I risultati mostrano che l’esecuzione del processo avviene per impostazione predefinita sul thread principale;
  • riga 23: il processo process4 genera 4 valori booleani su un thread non specificato, che saranno osservati su un thread non specificato. I risultati mostrano che l’esecuzione del processo e la sua osservazione avvengono per impostazione predefinita sul thread principale;

Il risultato dell’esecuzione di questo codice è il seguente:

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]
  • il processo process1 genera 1 numero reale (riga 4) sul thread di calcolo [RxComputationThreadPool-4], che viene osservato sul thread di calcolo [RxComputationThreadPool-3] (riga 6);
  • il processo process2 genera 2 stringhe di caratteri (righe 12, 14) sul thread di calcolo [RxComputationThreadPool-5], che vengono osservate su quello stesso thread (righe 13, 15);
  • il processo process3 genera 3 numeri interi (righe 21, 23, 25) sul thread principale, che vengono osservati sul thread di calcolo [RxComputationThreadPool-6] (righe 22, 24, 28);
  • il processo process4 genera 4 valori booleani (righe 34, 36, 38, 40) sul thread principale, che vengono osservati su quello stesso thread principale (righe 33, 35, 37, 39);

Si invita il lettore a seguire quanto sopra:

  • il ciclo di vita del processo osservato e del suo thread;
  • il ciclo di vita del suo osservatore e del relativo thread;

Gran parte dell’interesse delle librerie Rx risiede proprio in questo multi-threading che lo sviluppatore non deve gestire autonomamente.

7.5. Combinazioni di più osservabili

7.5.1. Esempio 14: unire due osservabili con [Observable.merge]

Presentiamo ora i metodi statici della classe [Observable] che consentono di combinare più osservabili in un unico osservabile risultante.

Il primo esempio di questo tipo sarà il seguente:


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 {
        // processo 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        // unione
        Process<?> process12 = new Process<>("process12",
                Observable.merge(process1.getObservable(), process2.getObservable()));
        // sottoscrizioni
        ProcessUtils.subscribe(1, process12);
    }
}
  • righe 15-17: un processo denominato [process1] emetterà 3 numeri reali su un thread di calcolo. Sarà inoltre osservato su un thread di calcolo;
  • righe 19-20: un processo denominato [process2] emetterà 2 stringhe di caratteri su un thread di calcolo. Il thread di osservazione non è predefinito. Abbiamo visto in precedenza che in questo caso il thread di osservazione è il thread di calcolo;
  • riga 23: i due processi vengono fusi, ovvero si crea un osservabile i cui elementi provengono simultaneamente da entrambi i processi. A tal fine si utilizza il metodo statico [Observable.merge]:
 

Contrariamente a quanto potrebbe far supporre lo schema sopra riportato, durante la fusione gli elementi di un flusso 1 possono inserirsi tra gli elementi di un flusso 2. È quanto mostrano i risultati dell’esecuzione:

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]
  • riga 3: il processo [process1] viene eseguito sul thread di calcolo [RxComputationThreadPool-4];
  • riga 4: il processo [process2] viene eseguito sul thread di calcolo [RxComputationThreadPool-5];
  • riga 9: il processo [process12] viene osservato sul thread di calcolo [RxComputationThreadPool-3]. Non conosco la regola che ha portato a questa scelta;
  • righe 9-11: si nota che l'osservatore osserva elementi dei due processi [process1] (riga 5) e [process2] (righe 6, 7) mentre nessuno dei due è terminato (c'è una commistione);
  • il processo [process12] termina (riga 17) quando entrambi i processi process1 e process2 sono terminati;

7.5.2. Esempio 15: concatenare due osservabili con [Observable.concat]

Esaminiamo ora il codice seguente:


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 {
        // processo 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, Schedulers.computation());
        // concat
        Process<?> process12 = new Process<>("process12",
                Observable.concat(process1.getObservable(), process2.getObservable()));
        // sottoscrizioni
        ProcessUtils.subscribe(1, process12);
    }
}
  • righe 15-17: un processo denominato [process1] emetterà 3 numeri reali su un thread di calcolo. Sarà inoltre osservato su un thread di calcolo;
  • righe 19-20: un processo denominato [process2] emetterà 2 stringhe di caratteri su un thread non specificato, in questo caso il thread principale predefinito. Sarà osservato su un thread di calcolo;
  • riga 23: i due processi vengono concatenati, ovvero viene creato un osservabile i cui elementi provengono da entrambi i processi. Non vi è alcuna commistione dei valori emessi. Il processo [process12] emetterà innanzitutto tutti i valori del processo [process1] e successivamente quelli del processo [process2]. A tal fine si utilizza il metodo statico [Observable.concat]:
 

I risultati dell'esecuzione sono i seguenti:

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]
  • righe 3-10: il processo [process1] è in esecuzione e il processo [process12] emette i valori generati da [process1];
  • riga 9: il processo [process1] è terminato;
  • righe 11-17: il processo [process2] è in esecuzione e il processo [process12] emette i valori generati da [process2];

C'è una stranezza per il processo process2: non era stato specificato alcun thread di esecuzione. Ci si sarebbe quindi potuti aspettare che, per impostazione predefinita, fosse il thread principale. E invece non è così. Il thread di esecuzione è stato il thread di calcolo [RxComputationThreadPool-3] (riga 11). Pertanto, quando non si specifica un thread di esecuzione o di osservazione, non è possibile formulare ipotesi sul thread che verrà scelto.

7.5.3. Esempio 16: combinare due osservabili con [Observable.zip]

Esaminiamo ora il seguente codice:


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 {
        // processo 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, null);
        // funzione di combinazione dei 2 processi
        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");
                }
            }
        };
        // compressione dei 2 processi
        Process<String> process12 = new Process<>("process12",
                Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));
        // sottoscrizioni
        ProcessUtils.subscribe(1, process12);
    }
}
  • righe 16-18: un processo denominato [process1] emetterà 3 numeri reali su un thread di calcolo. Sarà inoltre osservato su un thread di calcolo;
  • righe 20-21: un processo denominato [process2] emetterà 2 stringhe di caratteri su un thread non imposto. Anche il thread di osservazione non è imposto;
  • righe 23-32: istanziamento di un tipo [FuncN<String>] con una classe anonima. FuncN è un'interfaccia funzionale:
 

Il metodo [FuncN.call] accetta un array di oggetti e restituisce un tipo R. La funzione [funcn] verrà utilizzata per combinare i processi process1 e process2 in questo ordine. Nel metodo [FuncN.call]:

  • args[0] sarà un Double;
  • args[1] sarà un String;

In questo caso, il risultato di [funcn.call] sarà la stringa di caratteri della riga 27. Per ottenere questo risultato non è necessario conoscere i tipi degli argomenti del metodo call.

I due processi vengono combinati nel modo seguente:


// compressione dei 2 processi
Process<String> process12 = new Process<>("process12",
Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));

Il metodo [Observable.zip] funziona come segue:

 

Si nota che:

  • il primo argomento di zip è un Iterable<Observable>. Nel nostro esempio, abbiamo un parametro effettivo di tipo List<Observable> formato dai nostri due osservabili;
  • il secondo argomento di zip è di tipo FuncN. Nel nostro esempio, il parametro effettivo è [funcn];

L'esecuzione fornisce i seguenti risultati:

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]
  • righe 7, 11: il processo process12 emette due elementi;
  • riga 8: l'elemento aggiuntivo emesso dal processo process1, che non ha un partner nel processo process2, non viene emesso dal processo di risultato process12;

Si nota che il processo process2, al quale non erano stati assegnati né thread di esecuzione né thread di osservazione, ha utilizzato il thread principale per entrambi.

7.5.4. Esempio 17: combinare due osservabili con [Observable.combineLatest]

Esaminiamo ora il codice seguente:


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 {
        // processo 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null,
                Schedulers.computation());
        // combinazione dei 2 processi
        Process<Double> process12 = new Process<>("process12",
                Observable.combineLatest(process1.getObservable(), process2.getObservable(), (d1, d2) -> d1 + d2));
        // sottoscrizioni
        ProcessUtils.subscribe(1, process12);
    }
}
  • righe 14-16: un processo denominato [process1] emetterà 3 numeri reali su un thread di calcolo. Sarà inoltre osservato su un thread di calcolo;
  • righe 18-20: un processo denominato [process2] emetterà 2 numeri reali su un thread non vincolato. Questi saranno osservati su un thread di calcolo;
  • riga 23: i due osservabili vengono combinati con il seguente metodo statico [Observable.combineLatest]:
 

L'osservabile [combineLatest] funziona nel modo seguente: quando uno dei due osservabili emette un elemento E1, tale elemento viene combinato da [combineFunction] con l'ultimo elemento emesso dall'altro osservabile.

L'esecuzione di questo codice produce il seguente risultato:

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]
  • riga 5: l'emissione di process2 (56) viene combinata con l'ultimo elemento emesso da process1 (54, riga 4) e produce il risultato della riga 7;
  • riga 6: l'emissione di process1 (51,6) viene combinata con l'ultimo elemento emesso da process2 (56, riga 5) e produce il risultato della riga 8;
  • riga 9: l'emissione di process2 (261,8) viene combinata con l'ultimo elemento emesso da process1 (51,6, riga 6) e produce il risultato della riga 12;
  • riga 13: l'emissione di process1 (80,39) viene combinata con l'ultimo elemento emesso da process2 (261,8, riga 9) e produce il risultato della riga 15;

Ci troviamo qui in una variante dell’osservabile [zip], dove questa volta gli elementi combinati non sono necessariamente quelli che occupano la stessa posizione nei flussi. Si noti che il processo process2, al quale non era stato assegnato alcun thread di esecuzione, è stato qui eseguito sul thread principale (riga 2).

7.5.5. Esempio 18: combinare due osservabili con [Observable.amb]

Esaminiamo ora il codice seguente:


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 {
        // processo 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null, null);
        // combinazione dei 2 processi
        Process<Double> process12 = new Process<>("process12",
                Observable.amb(process1.getObservable(), process2.getObservable()));
        // sottoscrizioni
        ProcessUtils.subscribe(1, process12);
    }
}
  • righe 14-16: un processo denominato [process1] emetterà 3 numeri reali su un thread di calcolo. Sarà inoltre osservato su un thread di calcolo;
  • righe 18-20: un processo denominato [process2] emetterà 2 numeri reali su un thread non vincolato. Questi saranno osservati su un thread non vincolato;
  • riga 22: i due osservabili vengono combinati con il seguente metodo statico [Observable.amb]:
 

Come mostra lo schema sopra riportato, l’osservabile [Observable.amb(Observable o1, Observable o2)] emette gli elementi dell’osservabile che emette per primo. Ciò è confermato dai risultati dell’esempio presentato:

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]
  • riga 4: è il processo process2 a emettere per primo;
  • righe 8, 12: il processo process12 emette tutti gli elementi emessi dal processo process2 (righe 4, 11);

7.6. Catena di elaborazione di un osservabile

7.6.1. Esempio 19: trasformare un osservabile con [Observable.map]

Negli esempi precedenti abbiamo esaminato diverse combinazioni di due osservabili in un terzo osservabile. Presentiamo ora i metodi statici della classe [Observable] che consentono operazioni di trasformazione, filtraggio e aggregazione su un osservabile. Qui ritroveremo metodi analoghi a quelli della classe [Stream] studiati nel paragrafo 5.

Il nostro primo esempio sarà il seguente:


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 {
        // processo 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<String> process2 = new Process<>("process2",
                process1.getObservable().map(d -> String.format("valeur-%s", d)));
        // abbonamenti
        ProcessUtils.subscribe(1, process2);
    }
}
  • righe 14-16: un processo denominato process1 emetterà 3 numeri reali su un thread di calcolo. Sarà inoltre osservato su un thread di calcolo;
  • righe 17-18: i numeri emessi da process1 verranno trasformati in stringhe di caratteri in un processo process2;
  • riga 20: si osserva process2;

Il metodo [Observable.map] della riga 18 è analogo al metodo [Stream.map] esaminato nel paragrafo 5.5:

 

I risultati dell'esempio sono i seguenti:

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]
  • righe 4, 5 e 8: i valori restituiti da process1. Si tratta di numeri reali;
  • righe 6, 7, 10: le emissioni di process2 osservate. Si tratta di stringhe di caratteri;

7.6.2. Esempio-20: filtrare un osservabile con [Observable.filter]

L'esempio sarà il seguente:


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 {
        // processo 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().filter(i -> i % 2 == 0));
        // abbonamenti
        ProcessUtils.subscribe(1, process2);
    }
}
  • righe 11-12: un processo denominato process1 emetterà i numeri interi da 0 a 2 su un thread di calcolo. Sarà inoltre osservato su un thread di calcolo;
  • riga 14: i numeri generati da process1 verranno filtrati per mantenere in process2 solo i numeri pari;
  • riga 20: si osserva process2;

Il metodo [Observable.filter] della riga 18 è analogo al metodo [Stream.filter] studiato nel paragrafo 5.4:

 

I risultati dell'esempio sono i seguenti:

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]
  • righe 4, 5 e 7: le trasmissioni di process1;
  • righe 6, 9: le emissioni di process2 osservate. Si tratta degli elementi di process1 che sono pari;

7.6.3. Esempio 21: trasformare un osservabile con [Observable.flatMap]

L'esempio sarà il seguente:


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 {
        // processo 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMap(i -> {
            int value = i * 10;
            return Observable.just(value, value + 1, value + 2);
        }));
        // abbonamenti
        ProcessUtils.subscribe(1, process2);
    }
}
  • righe 12-13: un processo denominato process1 emetterà i numeri interi da 0 a 2 su un thread di calcolo. Sarà inoltre osservato su un thread di calcolo;
  • righe 15-18: ogni numero n generato da process1 viene trasformato in un osservabile che genera i 3 numeri (10*n, 10*n+1, 10*n+2). Se alla riga 15 si utilizzasse il metodo [map], process2 emetterebbe un tipo Observable<Integer> e non un tipo Integer. Il metodo [flatMap] utilizzato consente di appiattire (flatten) questa sequenza di elementi di tipo Observable<Integer> in una sequenza di elementi di tipo Integer costituita da ciascuno degli elementi di ciascuno dei Observable<Integer>;
  • riga 20: si osserva process2;

Il metodo [Observable.flatMap] della riga 15 è analogo al metodo [Stream.flatMap] esaminato al paragrafo 5.6.12:

 

I risultati dell'esempio sono i seguenti:

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]
  • righe 5-7: le tre trasmissioni di process2 a seguito della trasmissione della riga 4 di process1;
  • righe 9-11: le tre trasmissioni di process2 a seguito della trasmissione della riga 8 di process1;
  • righe 14-16: le tre trasmissioni di process2 a seguito della trasmissione della riga 12 di process1;

Il codice seguente mostra come creare un tipo Observable<Integer[]> a partire da process1 e [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 {
        // processo 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<Integer[]> process2 = new Process<>("process2", process1.getObservable().map(i -> {
            int value = i * 10;
            return new Integer[] { value, value + 1, value + 2 };
        }));
        // abbonamenti
        ProcessUtils.subscribe(1, process2);
    }
}
  • riga 14: si utilizza il metodo [Observable.map];
  • riga 16: che restituisce un tipo Integer[];

I risultati sono i seguenti:

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]
  • righe 6, 7, 10: si vedono i risultati di map;

Tutte queste trasformazioni dell'osservabile possono essere concatenate poiché ogni trasformazione produce un nuovo osservabile. È quanto mostra il seguente esempio [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 {
        // processo 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // processo 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));
        // abbonamenti
        ProcessUtils.subscribe(1, process2);
    }
}
  • righe 15-18: il flatMap è seguito da un filter;

I risultati dell’esecuzione sono i seguenti:

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]
  • righe 8-13: process2 ha emesso solo gli elementi pari provenienti da flatMap;

Un metodo simile a [flatMap] è il metodo [flatMapIterable] illustrato dal seguente esempio [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 {
        // processo 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // processo 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));
        // abbonamenti
        ProcessUtils.subscribe(1, process2);
    }
}

Riga 16: invece di utilizzare il metodo [flatMap], si utilizza il metodo [flatMapIterable]. In questo caso, la funzione di trasformazione deve produrre un tipo Iterable<T> (riga 18) anziché un tipo Observable<T>.

Si ottengono gli stessi risultati di prima.

Torniamo alla definizione del metodo [flatMap]:

 

Come si vede sopra, un elemento blu [3] si è inserito tra i due elementi verdi [1-2]. Ciò significa che, nell’operazione di appiattimento dei Observable<T>, il metodo [flatMap] rispetta l’ordine di emissione di questi diversi osservabili interni. Ciò è illustrato dal seguente esempio [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 {
        // processo 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // processo 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().flatMap(i -> process2.getObservable()));
        // abbonamenti
        ProcessUtils.subscribe(1, process3);
    }
}
  • righe 11-12: il processo process1 genera i numeri interi [0,1];
  • righe 14-15: il processo process2 emette i numeri interi [10,11,12];
  • righe 17-18: a ciascun elemento emesso da process1 viene associata l'osservabile del processo process2. Ciò significa che:
    • all'elemento [0] di process1 sarà associata un'osservabile che emette i [10,11,12];
    • lo stesso vale per l’elemento 1;

Alla fine, verranno emessi i 6 numeri [10, 11, 12, 10, 11, 12]. Vogliamo vedere in quale ordine.

I risultati dell’esecuzione sono i seguenti:

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]

Si nota che l’ordine di emissione del processo process3 è stato: [10, 10, 11, 12, 11, 12] (righe 11, 12, 14, 17, 19, 22). Si è quindi verificata effettivamente una commistione degli elementi generati dal processo process2. È possibile evitare ciò utilizzando il metodo [concatMap] al posto del metodo [flatMap]. È quanto illustra il seguente codice [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 {
        // processo 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // processo 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().concatMap(i -> process2.getObservable()));
        // abbonamenti
        ProcessUtils.subscribe(1, process3);
    }
}

Alla riga 18, [flatMap] è stato sostituito con [concatMap]. I risultati dell’esecuzione sono i seguenti:

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]

Si nota che l'ordine di emissione del processo process3 è stato: [10, 11, 12, 10, 11, 12] (righe 12-14, 17, 19, 22). Gli elementi generati dal processo process2 non sono stati mescolati.

Un’altra variante del metodo [map] è il metodo [switchMap]:

 

Sopra, dall’osservabile [1] nascono altri 3 osservabili [2] composti da 2 elementi, che vengono poi appiattiti come in [flatMap] e [3]. Si può notare che il risultato ha 5 elementi e non 6. Ciò è dovuto al fatto che, prima che il secondo osservabile emetta il suo elemento n. 2 [6], il terzo osservabile emette a sua volta il suo primo elemento [5], con la conseguenza che il secondo osservabile viene scartato. Non si trova quindi l’elemento [6] nell’osservabile risultante [3].

Per illustrare [switchMap], utilizzeremo il seguente esempio [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 {
        // processo 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // processo 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().switchMap(i -> process2.getObservable()));
        // abbonamenti
        ProcessUtils.subscribe(1, process3);
    }
}

L'esecuzione dell'esempio fornisce i seguenti risultati:

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 emette 2 elementi che danno origine a 2 osservabili process2 composti da 3 elementi;
  • riga 14: l’osservatore riceve l’elemento n. 0 emesso dal primo osservabile process2 alla riga 6;
  • riga 15: l’osservatore riceve l’elemento n. 0 emesso dal secondo osservabile process2 alla riga 13. Non è chiaro perché l’osservatore non abbia ricevuto in precedenza gli elementi 1 e 2 trasmessi dal primo osservabile process2 alle righe 7 e 8. Resta il fatto che il primo osservabile process2 viene scartato;
  • alla fine, l’osservatore vede solo 4 elementi (righe 14, 15, 17, 20) invece dei 6 che sono stati trasmessi;

7.6.4. Esempi-22: altri metodi della classe [Observable]

La classe [Observable] riprende numerosi metodi della classe [Stream] con un funzionamento analogo. Eccone alcuni. Ci limitiamo a fornire il codice e i relativi risultati.

[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 {
        // processo
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).take(3));
        // abbonamenti
        ProcessUtils.subscribe(1, process);
    }
}

Risultati

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 {
        // processo
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).takeLast(2));
        // sottoscrizioni
        ProcessUtils.subscribe(1, process);
    }
}

risultati

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 {
        // processi
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).skip(5).take(2));
        // sottoscrizioni
        ProcessUtils.subscribe(1, process);
    }
}

risultati

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 {
        // processi
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).reduce(0, (i, a) -> i + a));
        // sottoscrizioni
        ProcessUtils.subscribe(1, process);
    }
}
  • riga 10: calcola la somma degli elementi dell'osservabile. Il risultato è un osservabile che emette tale somma;

risultati

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 {
        // processi
        Process<Boolean> process = new Process<>("process", Observable.range(1, 10).all(i -> i > 10));
        // sottoscrizioni
        ProcessUtils.subscribe(1, process);
    }
}
  • riga 10: restituisce un Observable<Boolean> che emette l'elemento true, se il predicato del metodo [all] è vero per tutti gli elementi, altrimenti false;

risultati

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 {
        // processi
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).count());
        // sottoscrizioni
        ProcessUtils.subscribe(1, process);
    }
}
  • riga 10: [Observable.count] crea un osservabile a 1 elemento che è la somma degli elementi osservati;

risultati

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 {
        // processi
        Process<Integer> process = new Process<>("process", Observable.just(1, 2, 1, 3).distinct());
        // sottoscrizioni
        ProcessUtils.subscribe(1, process);
    }
}

risultati

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 {
        // processi
        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()));
        // sottoscrizioni
        ProcessUtils.subscribe(1, process);
    }
}
  • riga 11: il metodo [groupBy] raggruppa i 10 elementi emessi in 2 gruppi, i numeri pari e i numeri dispari. Il risultato è un tipo Observable<GroupedObservable<Boolean, Integer>>, ovvero un osservabile i cui elementi sono di tipo GroupedObservable<Boolean, Integer>, dove Boolean è il tipo della chiave del gruppo (false, true in questo caso) ed è anche il tipo del risultato della lambda passata come parametro al metodo [groupBy], mentre Integer è il tipo degli elementi del gruppo;
  • riga 12: il tipo GroupedObservable dispone di un metodo [asObservable] che consente di creare un osservabile a partire da questo tipo. Avremo quindi due tipi Observable<Integer>, uno per i numeri pari e l’altro per i numeri dispari. Da questi due osservabili, il metodo [concatMap] ne creerà uno solo;

risultati

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 {
        // processo 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // processo 2
        Process<Timestamped<Integer>> process2 = new Process<>("process2", process1.getObservable().timestamp());
        // sottoscrizioni
        ProcessUtils.subscribe(1, process2);
    }
}
  • alla riga 15, il metodo [timestamp] associa un'ora a ciascun elemento dell'osservabile elaborato;

risultati

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 questo esempio, è difficile stabilire a cosa si riferisca l'informazione timestamp:

  • righe 4-5: si vede che l'elemento 1 di process1 è stato emesso 139 ms dopo l'elemento 0;
  • righe 6 e 7: si vede che l’elemento 1 di process2 è stato rilevato 234 ms dopo l’elemento 0;
  • righe 5, 8: si nota che l'elemento 2 di process1 è stato emesso 33 ms dopo l'elemento 1;
  • righe 7 e 10: si vede che l'elemento 2 di process2 è stato osservato 37 ms dopo l'elemento 1;

Questi sfasamenti sono dovuti al fatto che i thread di osservazione e di esecuzione degli osservabili non sono gli stessi. Se si sostituiscono le righe 12-13 con le seguenti (Esempio22j):


// processo 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
null);
  • righe 2-3: non si impone il thread di osservazione. Si sa che in questo caso l'osservabile viene osservato nel punto in cui viene eseguito;

Ciò fornisce i seguenti risultati:

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]
  • righe 4 e 6: il processo process1 emette il suo elemento n. 1 587 ms dopo il suo elemento n. 0;
  • righe 5 e 7: l'osservatore osserva questi due elementi con uno scarto di 586 ms;
  • righe 6 e 8: il processo process1 emette il suo elemento n. 2 396 ms dopo il suo elemento n. 1;
  • righe 7 e 9: l'osservatore rileva questi due elementi con un intervallo di 396 ms;

In questo caso, i valori di timestamp sono coerenti: rappresentano effettivamente la data di emissione dell'elemento.

7.7. Gli scheduler

7.7.1. Esempio-23: lo scheduler [Schedulers.computation]

Esaminiamo ora gli scheduler di esecuzione. L’analisi verterà sul thread di esecuzione.

L’argomento degli scheduler è un po’ oscuro. I diversi scheduler sono presentati in questa domanda sul sito di StackOverflow [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases]:

 

Cercheremo di illustrare l’uso di questi diversi scheduler con alcuni esempi. Il primo illustra lo 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 {
        // processi
        @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);
        }
        // abbonamenti
        ProcessUtils.subscribe(1, processes);
    }
}
  • righe 14-19: si crea un array di 10 processi in esecuzione su un thread di calcolo;
  • riga 17: ogni processo genera un numero reale casuale;
  • riga 21: ci si abbona a tutti questi processi;

I risultati sono i seguenti:

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]
  • righe 2-10: i primi 8 processi si avviano su 8 thread diversi (la macchina utilizzata ha 8 core). Si può notare che si avviano tutti all'incirca nello stesso momento;
  • righe 17-19: 3 processi terminano, liberando così 3 thread;
  • righe 23-24: gli ultimi due processi possono quindi avviarsi utilizzando 2 dei thread così liberati;

Si noti quindi che lo scheduler [Schedulers.computation] fornisce un pool di n thread, dove n è il numero di core della macchina. I thread vengono eseguiti in parallelo su questi core.

7.7.2. Esempio 24: lo scheduler [Schedulers.io]

Eseguiamo il codice precedente con lo 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 {
        // processi
        @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);
        }
        // sottoscrizioni
        ProcessUtils.subscribe(1, processes);
    }
}
  • riga 18: i processi vengono eseguiti con i thread dello scheduler [Schedulers.io];

Ciò produce i seguenti risultati:

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]
  • righe 2-10: i 10 processi vengono avviati ciascuno su un thread diverso. A differenza del caso precedente, tutti i processi sono stati avviati con successo. Si nota che l’avvio di questi processi richiede 6 ms, mentre in precedenza era stato di 1 ms;
  • righe 13-18: gli osservabili emettono i dati uno dopo l’altro e non in modo quasi parallelo come era avvenuto in precedenza;

Qual è la differenza tra gli scheduler [Schedulers.io] e [Schedulers.computation]? Una risposta si può trovare in URL [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases]:

 

7.7.3. Esempio 25: lo scheduler [Schedulers.newThread]

Eseguiamo il codice precedente con lo scheduler [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 {
        // processi
        @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);
        }
        // sottoscrizioni
        ProcessUtils.subscribe(1, processes);
    }
}

I risultati ottenuti sono gli stessi di quelli ottenuti con lo scheduler [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]

Nei documenti URL e [http://stackoverflow.com/questions/33415881/retrofit-with-rxjava-schedulers-newthread-vs-schedulers-io] si spiega che lo scheduler [Schedulers.io] fornisce un pool di thread, cosa che invece non fa lo scheduler [Schedulers.newThread]. Un pool di thread crea automaticamente un numero n di thread. Li assegna ai processi che ne hanno bisogno. Una volta terminati, i thread non vengono eliminati, ma tornano nel pool e possono quindi essere riutilizzati da un altro processo. Ciò è più efficiente rispetto alla creazione e all’eliminazione continua dei thread. Si può quindi ritenere preferibile utilizzare lo scheduler [Schedulers.io].

7.7.4. Esempio 26: gli scheduler [Schedulers.immediate, Schedulers.trampoline]

Torniamo alla spiegazione fornita per questi due scheduler:

 

La spiegazione è abbastanza semplice da comprendere, ma quando si vuole illustrarla ci si rende conto di non averla compresa. È stato il libro [Learning Reactive Programming With Java 8] a permettermi di creare un esempio che riprende un esempio trovato in quel libro, ma lo semplifica. È il seguente:


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 {

        // uno scheduler
        Scheduler scheduler = Schedulers.immediate();
        // un worker di questo scheduler
        Worker worker = scheduler.createWorker();
        // un tipo Action0 da eseguire sul worker
        Action0 action02 = new Action0() {
            @Override
            public void call() {
                // log azione02
                ProcessUtils.showInfos.accept("action02");
            }
        };

        // un tipo Action0 da eseguire sul worker
        Action0 action01 = new Action0() {
            @Override
            public void call() {
                // si programma una nuova azione sullo stesso worker
                worker.schedule(action02);
                // log dell'azione 01
                ProcessUtils.showInfos.accept("action01");
            }
        };
        // l'azione 01 è programmata sul worker
        worker.schedule(action01);
    }

    // visualizzazioni
    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()));

}
  • riga 17: uno scheduler. Sarà o [Schedulers.immediate] come in questo caso, oppure [Schedulers.trampoline] in seguito;
  • riga 19: è possibile far eseguire azioni del tipo Action0 (righe 21, 20) sui worker dello scheduler. Il metodo [Scheduler.createWorker] consente di creare un worker. Il metodo [Worker.schedule(Action0)] consente di far eseguire un’azione di tipo Action0 da un worker;
  • righe 21-27: una prima azione denominata [action02] che verrà eseguita (riga 40) dal worker della riga 19;
  • righe 30-38: una seconda azione denominata [action01]. Essa ha la particolarità di far eseguire l’azione action02 sullo stesso worker su cui si trova essa stessa (riga 34). È qui che risiede la differenza tra [Schedulers.immediate] e [Schedulers.trampoline]:
    • se lo scheduler è [Schedulers.immediate], allora alla riga 34 l’azione action02 verrà eseguita immediatamente (da cui il nome dello scheduler) e l’azione action01 in corso verrà interrotta. A questo punto verrà visualizzato il messaggio della riga 25. Una volta terminata l’azione action02, l’azione action01 riprenderà e verrà visualizzato il messaggio della riga 36;
    • se lo scheduler è [Schedulers.trampoline], allora alla riga 34 l’azione action02 viene messa in attesa. Verrà eseguita solo quando l’attività in corso action01 sarà terminata. A quel punto apparirà il messaggio della riga 36. Una volta terminata l’azione action01, verrà eseguita l’azione action02 e vedremo il messaggio della riga 25;

L’esecuzione del codice sopra riportato fornisce i seguenti risultati:

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

Se alla riga 17 si utilizza lo scheduler [Schedulers.trampoline], si ottengono i risultati opposti:

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

Detto questo, è difficile stabilire un collegamento con gli osservabili. Non ho trovato alcun esempio convincente che potesse dimostrare l’utilità di eseguire un osservabile su uno di questi due thread. Eccone comunque uno, che però non trovo affatto naturale:


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();
        // osservabile 1 sul worker
        worker.schedule(() -> Observable.range(1, 2).subscribe(new Action1<Integer>() {

            @Override
            public void call(Integer i) {
                ProcessUtils.showInfos.accept(String.valueOf(i));
                // osservabile 2 sullo stesso worker
                worker.schedule(() -> Observable.range(100, 2).subscribe(new Action1<Integer>() {
                    @Override
                    public void call(Integer i) {
                        ProcessUtils.showInfos.accept(String.valueOf(i));
                    }
                }));
            }
        }));
    }
}
  • righe 13-14: si crea un worker a partire da uno dei due scheduler [Schedulers.immediate] e [Schedulers.trampoline];
  • riga 16: un primo osservabile obs1 viene programmato su questo worker per generare i numeri [1,2]
  • riga 22: ogni volta che viene osservato un elemento di questo osservabile obs1, viene avviata l’osservazione di un secondo osservabile obs2 sullo stesso worker per generare i numeri [100,101];

Con lo scheduler [Schedulers.immediate], si ottengono i seguenti risultati:

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]

Mentre con lo scheduler [Schedulers.trampoline] si ottengono i seguenti risultati:

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

C'è ancora molto da fare. Per approfondire la libreria RxJava, il lettore è invitato a proseguire la propria formazione consultando i riferimenti indicati all'inizio di questo documento. Ciononostante, disponiamo delle basi per utilizzare RxJava negli ambienti Swing e Android. È proprio ciò che mostreremo ora.