Skip to content

7. Бібліотека RxJava

Бібліотека RxJava базується на такій концепції: потік елементів типу T Observable<T> спостерігається одним або кількома підписниками (абонентами, спостерігачами, споживачами) Subscriber<T>. Бібліотека RxJava дозволяє потоку Observable<T> виконуватися в потоці T1, а його спостерігачеві Subscriber<T> — у потоці T2, при цьому розробнику нетурбуватися про управління життєвим циклом цих потоків та про природно складні проблеми, такі як обмін даними між потоками та їх синхронізація для виконання загального завдання. Таким чином, вона спрощує асинхронне програмування.

Потік Observable<T> генерує елементи типу T, які можна спостерігати по мірі їх генерації. Якщо спостерігач та об’єкт спостереження (термін «Observable<T>» вживається тут неточно) знаходяться в одному потоці, то об’єкт спостереження може генерувати елемент (i+1) лише тоді, коли спостерігач спожив елемент i. Випадків, коли така архітектура є доцільною, небагато. Якщо спостерігач і об’єкт спостереження не знаходяться в одному потоці, то об’єкт спостереження та його спостерігач поводяться автономно: об’єкт спостереження генерує елементи у своєму темпі, а спостерігач споживає їх у своєму темпі. Саме в цьому полягає цінність бібліотеки. До цього моменту ми завжди говорили про одного спостерігача. Насправді об’єкт спостереження може мати довільну кількість спостерігачів.

Бібліотека RxJava особливо добре підходить для архітектури, розглянутої в розділі 2, яку нагадаємо тут:

Image

  • у [1] сервісний рівень надає послуги, деякі з яких вимагають тривалого часу на виконання (наприклад, мережеві запити);
  • цей сервісний рівень викликається графічним інтерфейсом [1] (Swing, Android, JavaFx). Якщо сервісний рівень виконується в тому самому потоці, що й метод [swing], який його використовує, графічний інтерфейс зависає (не реагує) під час очікування результату виконання сервісу;
  • у [2] тонкий адаптаційний шар, реалізований за допомогою RxJava, дозволяє надати графічному шару асинхронну реалізацію того самого сервісу: він може виконуватися в потоці, відмінному від потоку методу графічного шару, який його викликає. У цьому випадку графічний інтерфейс [3] залишається реактивним: користувач може продовжувати взаємодіяти з ним, наприклад, запускати новий мережевий запит паралельно з першим, і, що найголовніше, йому можна надати можливість скасувати занадто тривалі обробки, що неможливо, якщо графічний інтерфейс застиг;
  • виклик [4] є синхронним, тоді як виклик [5-6] є асинхронним;

У цій архітектурі рівень [2] надає сервіси, що повертають типи Observable<T>, на які можуть підписуватися методи графічного рівня [3]. Послуга шару [2] передає свої результати по одному, і шар [3] може реагувати на кожен із них, наприклад, оновлюючи один або кілька компонентів графічного інтерфейсу.

Клас Observable<T> має кілька десятків методів. Це одна з складнощів бібліотеки: вона дуже насичена, і важко охопити всі її можливості. Ми розглянемо деякі з них. Опанування інших методів прийде з часом.

7.1. Створення об’єктів-спостерігачів та підписка на них

7.1.1. Приклад-01: метод [Observable.from]

  

Розглянемо такий код:


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) {
    // спостережувані цілі числа
    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");
      }
    });
  }
}
  • рядок 12: створюється тип Observable<Integer> на основі списку цілих чисел.

Клас Observable<T> — це потік елементів типу T, які можна спостерігати (бажано асинхронно, але не обов’язково) у міру їхнього генерування. Його визначення таке:

 

Як уже зазначалося, клас Observable<T> має кілька десятків методів. Деякі з них подібні до методів класу Stream<T>, розглянутого в розділі 5. Документація до RxJava містить «марбл-діаграми» [2], які ілюструють роботу цих методів:

  • рядок 3 ілюструє випромінювання спостережуваної величини в часі;
  • метод [4] застосовується до елементів, випромінюваних спостережуваною величиною. Зазвичай він утворює нову спостережувану величину;
  • рядок 5 показує отриману нову спостережувану величину;

Метод [Observable.from] має таку сигнатуру:

 

Статичний метод [Observable.from] дозволяє створити Observable<T> на основі колекції елементів типу T. Це дуже простий спосіб розпочати роботу з спостережуваними величинами. Рядок:


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

отже, видасть три елементи. Вона не видає їх одразу. Вона видаватиме їх повністю щоразу, коли з’явиться новий спостерігач. Це називається «холодною» спостережуваною величиною. Спостережувана величина повторно видає свої елементи для кожного нового підписника.

Попередню інструкцію можна розглядати як дію з налаштування спостережуваного об’єкта. Він налаштовується один раз і виконується n разів, якщо з’являється n підписників.

Як підписатися?

Один із способів — використати метод [Observable.subscribe], визначення якого в даному випадку таке:

 
  • перший параметр [Action1<T> onNext] (див. параграф 6.2) методу — це метод, який слід виконати, коли спостережуваний об’єкт генерує новий елемент T;
  • другий параметр [Action1<Throwable> onError] методу — це метод, який слід виконати, коли об’єкт спостереження генерує виняток;
  • третій параметр [Action0 onComplete] (див. параграф 6.1) методу — це метод, який слід виконати, коли об’єкт спостереження генерує виняток;
  • метод повертає тип [Subscription];

Тип [Subscription] представляє підписку на об’єкт спостереження. Його визначення таке:

 

Цінність цього інтерфейсу [1] полягає в його методі [2], який дозволяє скасувати підписку.

У нашому прикладі код підписки на об’єкт спостереження виглядає так:


    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");
      }
});
  • рядок 1: результат типу [Subscription] ігнорується;
  • рядки 1–15: три параметри є екземплярами анонімних класів. Ми також будемо використовувати лямбди. Перевага анонімних класів полягає в тому, що чітко видно типи даних, які очікує єдиний метод цих класів;
  • рядки 2–5: реалізація першого параметра типу [Action1<Integer>];
  • рядки 6–10: реалізація другого параметра типу [Action1<Throwable>];
  • рядки 11–15: реалізація третього параметра типу [Action0];

Повний код виглядає наступним чином:


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) {
    // спостережувані цілі числа
    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");
      }
    });
  }
}

Обсерватор у рядку 12 починає випускати свої 3 елементи, щойно у рядку 14 викликається метод [subscribe]. З цього моменту:

  • з кожним випущеним елементом виконуються рядки 15–18.
  • після випуску всіх 3 елементів виконуються рядки 24–29;
  • рядки 19–24 ніколи не будуть виконані, оскільки обсервабел тут не генерує винятку;

За замовчуванням об’єкт, що спостерігається, та спостерігач виконуються в одному потоці. Існує кілька попередньо визначених об’єктів, що спостерігаються, які виконуються в потоці, відмінному від головного (у даному випадку — потоку методу main), але для більшості з них це не так. Отже, тут усе відбувається у потоці методу [main]:

  • спостережуваний об’єкт видає елемент 1;
  • виконуються рядки 15–18, які виводять цей елемент;
  • обсервабел випускає елемент 2;
  • виконуються рядки 15–18, які виводять цей елемент;
  • обсервабел випускає елемент 3;
  • виконуються рядки 15–18 і відображається цей елемент;
  • обсервабел надсилає сповіщення [completed];
  • виконуються рядки 24–29;

Саме це показують отримані результати:

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

Клас [Exemple02] повторює [Exemple01], але цього разу використовує лямбда-функції як параметри методу [Observable.subscribe]:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;

public class Exemple02 {
  public static void main(String[] args) {
    // спостережувані цілі числа
    Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
    // підписка
    obs1.subscribe(
      (integer) -> System.out.printf("next : %s%n", integer),
      (th) -> System.out.println(th),
      () -> System.out.println("completed"));
  }
}

7.1.2. Приклад-03: клас Observer

  

Метод [Observable.subscribe], що дозволяє підписатися на об’єкт спостереження, має різні версії, серед яких така:


package dvp.rxjava.observables;

import java.util.Arrays;

import rx.Observable;
import rx.Observer;

public class Exemple03 {
    public static void main(String[] args) {
        // спостережувані цілі числа
        Observable<Integer> obs1 = Observable.from(Arrays.asList(1, 2, 3));
        // підписка
        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);
            }
        });
    };
}

У рядку 13 замість передачі трьох параметрів методу [subscribe] йому передається такий тип [Observer]:

 

Тип [Observer] є інтерфейсом із трьома методами:

  • [onNext(T t)], який викликається щоразу, коли об’єкт спостереження генерує елемент t;
  • [onError(Throwable th)], який викликається, коли об’єкт спостереження генерує виняток th;
  • [onCompleted], що викликається, коли об’єкт спостереження повідомляє про завершення випуску;

Принцип роботи коду аналогічний описаному раніше. Отримуємо такі результати:

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

7.1.3. Приклад-04: метод [Observable.create]

  

Статичний метод Observable.create визначено таким чином:

 
  • метод [create] повертає тип Observable<T>;
  • параметром методу [create] є функція типу [Observable.OnSubscribe<T>], визначена наступним чином:
 

Тип [Observable.OnSubscribe<T>] є функціональним інтерфейсом, який, у свою чергу, розширює функціональний інтерфейс [Action1<Subscriber<? super T>>]. Метод [call] цього інтерфейсу очікує тип [Subscriber] (абонент, підписник, спостерігач), який визначено таким чином:

 

У [1] видно, що клас [Subscriber<T>] реалізує інтерфейс [Observer<T>], представлений у розділі 7.1.2.

Зрештою, метод [<T> Observable.create]:

  • очікує як параметр екземпляр типу [Observable.OnSubscribe<T>], що має єдиний метод із сигнатурою: void call(Subscriber<T> s). Тип [Subscriber<T>] успадковує тип [Observer<T>] і, отже, має методи onNext, onError, onCompleted;
  • повертає тип Observable<T>;

Метод [<T> Observable.create] повертає налаштований об’єкт Observable. Поки що не відбувалося жодної емісії елементів. Коли підписник [Subscriber<T> s] підписується на цей об’єкт, що спостерігається, викликається метод [void call(s)] функції, переданої як параметр методу [<T> Observable.create]. Її завданням є надсилання елементів t типу T та виклик методу [s.onNext(t)] спостерігача під час кожного надсилання. Коли це завершується, має бути викликаний метод [s.onCompleted(t)] спостерігача, а метод [call] має завершитися. Якщо метод [call] зустрічає виняток th, має бути викликаний метод [s.onError(th)] спостерігача, а метод [call] має завершитися;

Щоб проілюструвати цей складний механізм роботи, ми використаємо наступний код [Exemple04]:


package dvp.rxjava.observables;

import rx.Observable;
import rx.Subscriber;

import java.util.Random;

public class Exemple04 {
    public static void main(String[] args) {
        // спостережувана конфігурація дійсних чисел
        Observable<Double> obs1 = Observable.create(new Observable.OnSubscribe<Double>() {
            @Override
            public void call(Subscriber<? super Double> subscriber) {
                for (int i = 0; i < 3; i++) {
                    // видача елемента i
                    subscriber.onNext(new Random((i + 1)).nextDouble());
                }
                // завершення передачі
                subscriber.onCompleted();
            }
        });
        // підписка і, отже, передача
        obs1.subscribe((d) -> System.out.printf("onNext %s%n", d), (th) -> System.out.printf("onError %s%n", th),
                () -> System.out.println("onCompleted"));
    }
}
  • рядок 11: створюється об’єкт, що спостерігається, який генерує типи Double;
  • рядки 11–21: параметр методу [create] інстанціюється за допомогою анонімного класу, що має єдиний метод [call] у рядках 12–20. Спостерігабельний об’єкт, створений у рядку 11, готовий до випромінювання, але він почне випромінювати лише тоді, коли з’явиться спостерігач;
  • рядки 13–21: метод [call] отримує посилання на спостерігача;
  • рядки 14–17: передача 3 елементів спостерігачеві;
  • рядок 19: повідомлення спостерігачеві про завершення передачі;
  • рядки 23–24: підписка на спостережуваний об’єкт із рядка 11. Три параметри [onNext, onError, onCompleted] методу [subscribe] реалізуються за допомогою трьох лямбда-виразів. Ця підписка створить підписника [Subscriber<Double>], який буде передано до методу [call] у рядку 13. Після цього почнеться передача елементів;
  • все відбувається в одному потоці: спостережуваний об’єкт і спостерігач;

Отримуємо такі результати:

1
2
3
4
onNext 0.7308781907032909
onNext 0.7311469360199058
onNext 0.731057369148862
onCompleted

Метод [Observable.create] дозволяє створити об’єкт спостереження на основі будь-якого явища. Саме цей метод ми використовували в розділі 2 «Відкриття», щоб перетворити синхронний інтерфейс на асинхронний.

7.1.4. Приклад-05: рефакторинг [Exemple-04]

  

У наступному прикладі представлено нову версію статичного методу [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) {
        // конфігурація спостережуваної величини з числами дійсних
        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++) {
                    // очікування
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // помилка
                        subscriber.onError(e);
                    }
                    // дія
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // завершено
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // підписник
        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));
            }
        };

        // підписка
        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()));
    }
}
  • рядок 56: нова версія статичного методу [Observable.subscribe] приймає як параметр тип [Subscriber], який ми розглянули в попередньому параграфі;
  • рядки 37–52: підписник (абонент, спостерігач). Він реалізує інтерфейс Observer із трьома методами: onNext, onError, onCompleted;
  • рядки 61–64: відтепер ми зосередимося на потоках, у яких виконуються об’єкт спостереження та його спостерігач;
  • рядок 62: назва потоку;
  • рядок 63: поточний час, виражений у секундах і мілісекундах. Це дозволить нам простежити у часі передачу елементів об’єктом спостереження та їх обробку спостерігачем;
  • цей код має ту саму функціональність, що й попередній. Ми просто рефакторували останній;

Отримані результати такі:

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]
  • 1-й рядок результатів: до 56-го рядка коду ще нічого не відбулося. Обсервабел було просто налаштовано;
  • 2-й рядок результатів: 56-й рядок коду викликає метод [call] з 15-го рядка. У 3-му рядку до спостерігача надсилається число 80,39;
  • рядок 4: спостерігач отримує передане число;
  • рядки 5–8: попередній процес повторюється 2 рази;
  • рядок 9: об’єкт, що спостерігається, надсилає повідомлення про завершення передачі;
  • рядок 10: спостерігач отримує його;
  • рядок 11: відображається рядком 57 коду;

Отже, бачимо, що лише один рядок 56 підписки спричинив відображення рядків 2–10 результатів. Коли починаєш працювати з бібліотекою RxJava, виникає питання, як все взаємопов’язано, зокрема які зв’язки існують між спостерігачем та об’єктом спостереження. Тут видно, що рядок 56, підписка на спостережуване,

  • спричинила виведення всіх елементів об’єкта спостереження;
  • що спостережуване та спостерігач виконуються в одному потоці;
  • що через це спостерігається така послідовність: виведення елемента i, спостереження за елементом i, виведення елемента (i+1), спостереження за елементом (i+1), ...

Нагадаємо, що емітер чекав перед тим, як емітувати свої елементи:


                    // очікування
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // помилка
                        subscriber.onError(e);
}

де i у рядку 3 позначає номер передачі (0<=i<3). Якщо розглянути час передачі елементів об’єкта спостереження:

  • рядки 2, 3: елемент 0 був переданий приблизно через 500 мс після початку підписки;
  • рядки 3, 5: елемент 1 був відправлений приблизно через 400 мс після елемента 0;
  • рядки 5, 7: елемент 2 був відправлений приблизно через 300 мс після елемента 1;

7.2. Потік виконання, потік спостереження

7.2.1. Приклад-06: об’єкт спостереження та спостерігач у потоці, відмінному від [main]

  

Ми рефакторуємо попередній приклад [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) {

        // блокування
        CountDownLatch latch = new CountDownLatch(1);

        // конфігурація спостережуваних величин типу «реальні числа»
        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++) {
                    // очікування
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // помилка
                        subscriber.onError(e);
                    }
                    // дія
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // завершено
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // підписник
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // опускаємо бар'єр
                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));
            }
        };

        // продовження налаштування спостережуваного об’єкта
        obs1 = obs1.subscribeOn(Schedulers.computation());
        // підписка
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // очікування перед шлагбаумом
        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()));
    }
}
  • рядок 16: створюємо бар'єр (семафор) з об'єктом типу [CountDownLatch]. Цей об'єкт слугує для синхронізації потоків між собою. Тут він ініціалізується значенням 1, яке ми назвемо значенням бар'єру (або семафора). Потік переходить у стан очікування бар'єру за допомогою операції:

latch.await();

Поттік блокується, якщо значення бар'єра >0. Потік може збільшувати/зменшувати внутрішнє значення бар'єра. У рядку 48 значення бар'єра зменшується на 1.

  • рядок 63: спостережуваний об’єкт налаштований так, щоб виконуватися у потоці, наданому планувальником [Schedulers.computation()]. Цей планувальник може надавати стільки потоків, скільки є ядер на машині виконання. У розділі про приклад програми показано використання інших планувальників (див. розділ 2.8);

Принцип роботи коду такий:

  • метод [main] виконується в головному (main) потоці;
  • рядок 66: запускає виведення елементів спостережуваного об’єкта. Вони будуть виведені в потоці, відмінному від головного;
  • рядок 70: головний потік заблоковано, оскільки значення бар’єра дорівнює 1 (див. рядок 16). Він зможе продовжити роботу лише тоді, коли це значення зміниться на 0. Це відбувається в рядку 48. Саме спостерігач опускає бар’єр, коли отримує повідомлення про те, що об’єкт спостереження завершив передачу даних;

Виконання дає такі результати:

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]
  • рядок 1: відбудеться підписка;
  • рядок 2: це запускає виконання методу [call] у потоці [RxComputationThreadPool-1]. Тепер ми маємо паралельне виконання з двома потоками;
  • рядок 3: з невідомих причин потік [RxComputationThreadPool-1] передав контроль. Потім потік [main] переймає управління і блокується бар’єром (рядок 70 коду). З цього моменту працювати може лише потік [RxComputationThreadPool-1];
  • рядки 4–11: спостерігається поведінка, що була раніше між об’єктом спостереження та його спостерігачем, але тепер усе відбувається у потоці [RxComputationThreadPool-1];
  • рядки 12–13: спостерігач опустив шлагбаум (рядок 48 коду), і потік [RxComputationThreadPool-1] завершився. Потік [main] бере на себе управління та виводить два повідомлення;

7.2.2. Приклад-07: об’єкт спостереження та спостерігач у двох різних потоках

  

Ми змінюємо попередній приклад наступним чином:


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

        // охоронець шлагбаума
        CountDownLatch latch = new CountDownLatch(1);

        // конфігурація спостережуваних величин з числа дійсних
        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++) {
                    // очікування
                    try {
                        Thread.sleep(500 - i * 100);
                    } catch (InterruptedException e) {
                        // помилка
                        subscriber.onError(e);
                    }
                    // дія
                    double value = new Random().nextInt(100) * 1.2;
                    showInfos(String.format("Observable.call onNext(%s)", value));
                    subscriber.onNext(value);
                }
                // завершено
                showInfos(String.format("Observable.call onCompleted"));
                subscriber.onCompleted();
            }
        });

        // підписник
        Subscriber<Double> subscriber = new Subscriber<Double>() {
            @Override
            public void onCompleted() {
                showInfos("Subscriber.onCompleted");
                // опускаємо бар'єр
                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));
            }
        };

        // продовження конфігурації спостережуваного об’єкта
        obs1 = obs1.subscribeOn(Schedulers.computation()).observeOn(Schedulers.computation());
        // підписка
        showInfos("avant souscription");
        obs1.subscribe(subscriber);
        // очікування на перевищення бар'єру
        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()));
    }
}

Код ідентичний коду з попереднього прикладу, за винятком рядка 63:


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

яка налаштовує об’єкт спостереження (subscribeOn) та спостерігач (observeOn) для виконання в одному з потоків, що надаються планувальником [Schedulers.computation()].

Отримано такі результати:

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]

Можна відзначити наступне:

  • спостережуваний об’єкт виконується у потоці [RxComputationThreadPool-4] (рядки 3–4, 6, 8–9);
  • спостерігач виконується у потоці [RxComputationThreadPool-3] (рядки 5, 7, 10–11);
  • що вони виконуються автономно. Так, у рядках 8–9 об’єкт спостереження надсилає 2 сповіщення (onNext, onCompleted) до того, як об’єкт спостереження отримає сповіщення [onNext] (рядок 10);

Бібліотека RxJava відповідає за передачу даних (випромінювання) з потоку об’єкта спостереження до потоку спостерігача. Розробнику не потрібно про це турбуватися.

Ми розглянули, як створювати об’єкти спостереження (Observable.from, Observable.create). Тепер розглянемо попередньо визначені об’єкти спостереження з бібліотеки RxJava.

7.3. Заздалегідь визначені об’єкти, що спостерігаються

7.3.1. Приклад-08: метод [Observable.range]

 

Відтепер ми будемо використовувати спеціальні класи для спостережуваних процесів та їхніх спостерігачів. Ідея полягає в тому, щоб мати змогу реєструвати їхні імена, потоки виконання та час виконання, щоб можна було відстежувати їх у часі.

Клас [Process] буде просто Observable, якому можна присвоїти ім’я. Він реалізовуватиме такий інтерфейс [IProcess]:


package dvp.rxjava.observables.utils;

import rx.Observable;

public interface IProcess<T> {

    // назва спостережуваного параметра
    public String getName();

    // спостережувана величина
    public Observable<T> getObservable();

}

Цей інтерфейс може бути реалізований наступним класом [Process<T>]:


package dvp.rxjava.observables.utils;

import rx.Observable;
import rx.Scheduler;

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

    // назва спостережуваного параметра
    protected String name;
    // спостережуваний процес
    protected Observable<T> observable;

    // конструктори
    public Process(String name, Observable<T> observable) {
        // локальні ініціалізації
        this.name = name;
        this.observable = observable;
    }

    // методи getter та setter
    public String getName() {
        return name;
    }

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

}
  • рядок 9: назва процесу;
  • рядок 11: спостережувана величина;
  • рядки 14–18: конструктор;

Спостерігач буде описаний наступним класом [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> {

...
}
  • у рядку 11 клас Observateur<T> успадковує клас Subscriber<T>, який ми коротко розглянули в розділі 7.1.3. Ми будемо використовувати його як аргумент методу [Observable.subscribe]:

// спостережуване виконання (спостереження)
obs1.subscribe(observateur);

Метод [Observable.subscribe], використаний у рядку 2 вище, має таке визначення:

 

Роль [Subscriber] полягає головним чином у керуванні елементами, що надсилаються спостережуваним об’єктом, на який він підписався, за допомогою методів інтерфейсу [Observer]: onNext, onError, onCompleted. Клас [Subscriber] має такі методи:

 

У коді класу [Observateur] ми будемо використовувати методи [1] та isUnsubscribed, щоб дізнатися, чи була скасована підписка абонента. Повний код класу [Observateur<T>] виглядає так:


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

    // семафор
    private CountDownLatch latch;
    // метод відображення
    private Consumer<String> showInfos;
    // ім'я спостерігача
    private String observerName;
    // ім'я спостережуваного процесу
    private String processName;

    // конструктори
    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;
    }

    // --------------------------- реалізація інтерфейсу Observer<T>
    @Override
    public void onCompleted() {
        // кінець трансляцій
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber [%s,%s].onCompleted", observerName, processName));
        }
        // кінець блокування головного потоку
        latch.countDown();
    }

    @Override
    public void onError(Throwable e) {
        // помилка передачі
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber [%s, %s].onError (%s)", observerName, processName, e));
        }
    }

    @Override
    public void onNext(T value) {
        // додаткова передача
        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));
            }
        }
    }
}
  • крім характеристик Subscriber, спостерігач Observateur буде містити таку інформацію:
    • рядок 14: бар’єр або семафор, який слугуватиме для блокування головного потоку, доки спостерігач не отримає всі елементи, передані об’єктом спостереження. Це відбуватиметься у рядку 36 коду, коли спостерігач отримає від об’єкта спостереження повідомлення про завершення передачі;
    • рядок 16: екземпляр Consumer<String>, який слугуватиме для виведення повідомлення на консоль;
    • рядок 18: ім’я спостерігача, щоб розрізняти їх між собою, коли їх декілька;
    • рядок 20: ім’я спостережуваного процесу;
  • рядки 36, 46, 54: методи [onCompleted, onError, onNext] інтерфейсу [Observer<T>], реалізованого абстрактним класом [Subscriber<T>]. Цей клас їх не реалізує. Тому це потрібно зробити в його дочірніх класах. Перш ніж виконувати будь-які дії в цих методах, перевіряємо, чи спостерігач не відписався від об’єкта спостереження, за яким він спостерігає;
  • рядок 59: метод [onNext] спостерігача записує рядок jSON отриманого елемента. Це дозволить нам відображати різні типи елементів;

Отже, розглянемо новий метод класу Observable — метод [range]:

 

Спостережувана величина Observable.range(n,m) видає (m) цілих чисел у діапазоні від n до n+m-1. Ми розглянемо її за допомогою такого коду [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 {

        // кількість спостерігачів
        final int nbObservateurs = 2;

        // семафор
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // спостережувана конфігурація
        Observable<Integer> obs1 = Observable.range(15, 3).subscribeOn(Schedulers.computation());
        // спостережуване виконання (спостереження)
        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"));
        }
        // очікування
        showInfos.accept("main : attente fin observation");
        latch.await();
        // завершення
        showInfos.accept("main : fin observation");
    }

    // відображення
    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()));
}
  • рядок 16: ми будемо використовувати двох спостерігачів;
  • рядок 19: семафор ініціалізується значенням 2, оскільки кожен спостерігач буде розміщений у окремому потоці. Отже, головний потік повинен чекати на завершення роботи обох спостережних потоків;
  • рядок 22: ми налаштовуємо об’єкт спостереження таким чином, щоб він виконувався у потоці планувальника [Schedulers.computation()]. Спостерігач буде у тому самому потоці, що й об’єкт спостереження;
  • рядки 25–27: підписуємо двох спостерігачів на об’єкт спостереження. Це спричинить повне виконання об’єкта спостереження для кожного з спостерігачів: будуть виведені цілі числа 15, 16 та 17;
  • рядок 30: головний потік очікує завершення роботи спостерігачів;

Отримано такі результати:

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]
  • рядок 2: головний потік заблоковано в очікуванні завершення роботи 2 спостерігачів;
  • рядки 3–4: бачимо, що спостерігач 0 знаходиться у потоці [RxComputationThreadPool-1], а спостерігач 1 — у потоці [RxComputationThreadPool-2];
  • рядки 3–10: бачимо, що обидва спостерігачі отримують абсолютно однакові елементи;

Ми використаємо клас Observateur, визначений таким чином, щоб проілюструвати поведінку інших типів спостережуваних об’єктів.

7.3.2. Приклад-09: методи Observable.[interval, take, doNext]

  
 

Цей приклад ілюструє використання об’єкта спостереження Observable.interval (довгий інтервал, одиниця TimeUnit), який генерує довгі цілі числа через регулярні проміжки часу. Слід звернути увагу на пункт [1]: за замовчуванням спостережувана величина [Observable.interval] виконується в одному з потоків планувальника [Schedulers.computation].

Код матиме такий вигляд:


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 {

        // кількість спостерігачів
        final int nbObservateurs = 2;

        // семафор
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // спостережувана конфігурація
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // спостережуване виконання (спостереження)
        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"));
        }
        // очікування
        showInfos.accept("main : attente fin observation");
        latch.await();
        // завершення
        showInfos.accept("main : fin observation");
    }

    // відображення
    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()));
}
  • рядок 22: обсервабел видає довгі цілі числа кожні 500 мілісекунд. Послідовність починається з числа 0;
  • рядок 22: цей об’єкт, що спостерігається, генерує нескінченну кількість значень. Метод [Observable.take(n)] створює новий об’єкт, що спостерігається, який зберігає лише перші n згенерованих елементів;
 

Повернемося до коду обсервабеля:


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

У рядку 2 метод [Observable.doOnNext] виконується щоразу, коли об’єкт, що спостерігається, випускає новий елемент. Це часто використовується для реєстрації інформації в журналі. У даному випадку ми хочемо зареєструвати дату випуску елементів, щоб перевірити, чи дотримується інтервал у 500 мілісекунд. Метод [Observable.doOnNext] не змінює об’єкт спостереження, до якого він застосовується. Його визначення таке:

 

Виконання дає такі результати:

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]
  • рядки 3, 7 та 11: бачимо, що інтервал передачі приблизно дорівнює 500 мс;
  • Обидва спостерігачі, звісно, знаходяться у двох різних потоках, хоча об’єкт спостереження не було налаштовано на виконання з використанням конкретного планувальника. Тут ми бачимо стандартну роботу об’єкта спостереження [Observable.interval];

7.3.3. Приклади-10/12: методи Observable.[error, empty, never]

 

Відтепер ми будемо більш лаконічними у наших ілюстраціях методів класу [Observable]. Попередній код був таким:


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 {

        // кількість спостерігачів
        final int nbObservateurs = 2;

        // семафор
        CountDownLatch latch = new CountDownLatch(nbObservateurs);

        // спостережувана конфігурація
        Observable<Long> obs1 = Observable.interval(500L, TimeUnit.MILLISECONDS).take(3)
                .doOnNext(l -> showInfos.accept(l.toString()));
        // спостережуване виконання (спостереження)
        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"));
        }
        // очікування
        showInfos.accept("main : attente fin observation");
        latch.await();
        // завершення
        showInfos.accept("main : fin observation");
    }

    // відображення
    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()));
}

Цей код уже використовувався в попередньому прикладі. Змінювалися лише рядки 21–22. Тому ми винесемо більшу частину цього коду в наступний клас [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 {

        // семафор
        CountDownLatch latch = new CountDownLatch(nbObservateurs * processes.length);

        // спостережуване виконання (спостереження)
        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()));
            }
        }
        // очікування
        showInfos.accept("main : attente fin observation");
        latch.await();
        // кінець
        showInfos.accept("main : fin observation");
    }

    // відображення
    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()));
}
  • рядок 13: метод приймає два параметри:
    • nbObservateurs: кількість спостерігачів процесів, переданих як другий параметр;
    • processes: процеси (іменовані спостережувані), які потрібно спостерігати. Завдяки нотації [IProcess<?>] процеси зможуть генерувати елементи різних типів;
  • рядок 16: семафор повинен перейти у зелений стан, коли всі спостерігачі завершать усі свої спостереження. Початкове значення семафора дорівнює кількості спостерігачів, помноженій на кількість спостережень;
  • рядки 20–25: кожного спостерігача підписують на всі процеси, які потрібно спостерігати;
  • рядок 23: отримуємо спостережуване значення з процесу (див. параграф 7.3.1);
  • рядок 23: до нього підписується спостерігач. Йому передаються 4 відомості:
    • його ім’я;
    • семафор, який він повинен зменшувати, коли отримує повідомлення про завершення передачі спостережуваної величини, яку він спостерігає;
    • метод, який слід використовувати, коли він хоче записати інформацію в журнал консолі;
    • ім’я процесу, який він буде спостерігати;

Оскільки ці класи визначені, приклад 10 матиме такий вигляд:


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 {
        // спостережувана конфігурація
        Observable<?> obs = Observable.error(new RuntimeException("Erreur !!!")).subscribeOn(Schedulers.computation());
        // спостережуване виконання (спостереження)
        ProcessUtils.subscribe(2,new Process<>("process1", obs));
    }
}

У рядку 11 статичний метод [Observable.error] визначено таким чином:

 

Отже, у рядку 8 налаштовується об’єкт спостереження, який просто генерує виняток, що надсилається до методу [onError] його підписників. Виконання дає такі результати:


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]

У рядках 3 і 4 метод [onError] обох підписників отримав виняток, згенерований об’єктом спостереження.

Це виконання має одну особливість: методи [onCompleted] обох спостерігачів не були викликані. Як наслідок, бар’єр не було опущено, і головний потік залишається заблокованим у статичному методі [ProcessUtils.subscribe] у наступному рядку 3:


// очікування
showInfos.accept("main : attente fin observation");
latch.await();
// завершення
showInfos.accept("main : fin observation");

Тут ми бачимо, що в разі помилки спостережуваного методу [onCompleted] підписників не викликається. Тому ми змінюємо метод [Observateur.onError] наступним чином:


    @Override
    public void onError(Throwable e) {
        // помилка передачі
        if (!isUnsubscribed()) {
            showInfos.accept(String.format("Subscriber[%s, %s].onError (%s)", observerName, processName, e));
        }
        // завершення блокування головного потоку
        latch.countDown();
}

Ми додаємо рядки 7–8, щоб усунути перешкоду у разі помилки спостережуваної величини. З цим новим кодом виконання дає такі результати:


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]

Ми отримуємо рядок 5, якого раніше не було.

Приклад 11 буде таким:


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 {
        // спостережувана конфігурація
        Observable<?> obs1 = Observable.empty();
        // спостережуване виконання (спостереження)
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

У рядку 10 статичний метод [Observable.empty] створює спостережуваний об’єкт, який не випускає жодного елемента. Він випускає лише повідомлення про закінчення випуску;

 

Виконання коду з наведеного вище прикладу дає такі результати:

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]
  • рядки 2 і 3: бачимо, що обидва спостерігачі отримують повідомлення про завершення передачі, не отримавши раніше жодних елементів.

Можна запитати, для чого взагалі потрібен цей метод. Його можна використовувати аналогічно до колекції, яка спочатку порожня, а потім у неї додаються елементи:

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

У рядку 3 початкову спостережувану величину obs (рядок 1) об’єднують з іншими спостережуваними величинами.

Приклад 12 ілюструє статичний метод [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 {
        // спостережувана конфігурація
        Observable<?> obs1 = Observable.never();
        // спостережуване виконання (спостереження)
        ProcessUtils.subscribe(2,new Process<>("process1",obs1));
    }
}

Статичний метод [Observable.never] створює спостережуваний об’єкт, який ніколи не генерує сигналів:

 

Виконання прикладу дає такі результати:

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

У рядку 2 головний потік очікує нескінченно довго. Дійсно, жоден об’єкт спостереження не надсилає сповіщення [onCompleted], яке дозволяє перевести семафор (шлагбаум) у зелений стан (опустити шлагбаум).

7.4. Multi-threading

7.4.1. Приклад-13: потік дії, потік спостереження

У розділі 7.1.3 ми створили об’єкт спостереження за допомогою статичного методу [Observable.create]:

 
  • метод [create] повертає тип Observable<T>;
  • параметр методу [create] — це функція типу [Observable.OnSubscribe<T>], визначена таким чином:
 

Тип [Observable.OnSubscribe<T>] є функціональним інтерфейсом, який, у свою чергу, розширює функціональний інтерфейс [Action1<Subscriber<? super T>>]. Метод [call] цього інтерфейсу очікує тип [Subscriber] (абонент, підписник, спостерігач). Далі в цьому документі ми іноді будемо називати тип [Observable.OnSubscribe<T>] «дією». Ми створимо власні дії, які матимуть імена. Це будуть екземпляри наступного інтерфейсу [IProcessAction]:

  

package dvp.rxjava.observables.utils;

import rx.Observable;

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

    // дія має назву
    public String getName();
}
  • рядок 5: інтерфейс [IProcessAction<T>] має всі характеристики інтерфейсу [Observable.OnSubscribe<T>];
  • рядок 8: крім того, він має метод [getName], який повертає ім’я екземпляра, що реалізує інтерфейс;

Ми будемо використовувати таку дію з назвою [ProcessAction01]:


package dvp.rxjava.observables.utils;

import java.util.Random;

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

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

    // дані
    private String name;
    private int nbValues;
    private Func1<Integer, T> func1;

    // конструктори
    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++) {
            // очікування
            try {
                Thread.sleep(new Random().nextInt(500));
            } catch (InterruptedException e) {
                // помилка
                ProcessUtils.showInfos.accept(String.format("Observable (%s) onError", getName()));
                subscriber.onError(e);
            }
            // видача елемента
            T value = func1.call(i);
            ProcessUtils.showInfos.accept(String.format("Observable (%s,%s) onNext (%s)", getName(), i, value));
            subscriber.onNext(value);
        }
        // завершено
        ProcessUtils.showInfos.accept(String.format("Observable (%s) onCompleted", getName()));
        subscriber.onCompleted();
    }

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

}
  • рядок 8: клас [ProcessAction01<T>] реалізує інтерфейс [IProcessAction<T>], а отже, і інтерфейс [Observable.OnSubscribe<T>];
  • рядок 11: назва дії;
  • рядок 12: кількість значень, що підлягають виведенню;
  • рядок 13: екземпляр типу [Func1<Integer, T>], який на основі цілого числа створює тип T, що буде виведений спостережуваним об’єктом (рядки 35 та 37);
  • рядки 16–20: у конструктор передаються назва дії, кількість значень, що мають бути виведені, та функція виведення;
  • рядки 23–42: код процесу;
  • рядок 23: метод [call] отримує як параметр підписника об’єкта спостереження, пов’язаного з процесом;
  • рядок 28: процес випускає свої елементи після очікування тривалістю, що визначається випадковим чином;
  • рядок 32: генерація помилки;
  • рядок 37: звичайна передача;
  • рядок 41: надсилання повідомлення про завершення надсилання;
  • рядки 25–38: дія надсилає nbValues реальних значень після випадкового часу очікування (рядок 30);
  • рядок 35: значення, яке потрібно вивести, надається функцією [func1], що передається як параметр конструктору (рядок 16);

Ми рефакторуємо клас [Process] (див. розділ 7.3.1), щоб його можна було створити також із іменованою дією. Додаємо до нього такий конструктор:


public Process(IProcessAction<T> na, Scheduler schedulerObserved, Scheduler schedulerObserver) {
        // назва процесу=назва дії
        name = na.getName();
        // дія --> спостережувана величина
        observable = Observable.create(na);
        // потік виконання спостережуваного процесу
        if (schedulerObserved != null) {
            observable = observable.subscribeOn(schedulerObserved);
        }
        // потік спостереження спостерігача
        if (schedulerObserver != null) {
            observable = observable.observeOn(schedulerObserver);
        }
    }
  • у рядку 1 конструктор приймає 3 параметри:
    1. іменовану дію, яка буде використовуватися для створення об’єкта спостереження (рядок 5);
    2. планувальник спостережуваного процесу (може бути null);
    3. планувальник спостерігача (може бути null);
  • рядок 5: спостережувана величина створюється на основі дії, переданої як параметр;

Наступний код [Exemple13] спостерігає за різними об’єктами спостереження:


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 {
        // процес 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 1, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        // процес 3
        Process<Integer> process3 = new Process<>(new ProcessAction01<Integer>("process3", 3, i -> i * 2), null,
                Schedulers.computation());
        // процес 4
        Process<Boolean> process4 = new Process<>(new ProcessAction01<Boolean>("process4", 4, i -> i % 2 == 0), null, null);
        // підписки
        ProcessUtils.subscribe(1, process1);
        ProcessUtils.subscribe(1, process2);
        ProcessUtils.subscribe(1, process3);
        ProcessUtils.subscribe(1, process4);
    }
}
  • рядки 13–15: процес process1 генерує 1 дійсне число в обчислювальному потоці, яке буде спостерігатися в іншому обчислювальному потоці;
  • рядки 17–18: процес process2 генерує 2 символьні рядки у обчислювальному потоці, при цьому не вказано, у якому потоці відбуватиметься спостереження. Результати показують, що за замовчуванням спостереження відбувається в тому самому потоці, що й виконання процесу;
  • рядки 20–21: процес process3 генерує 3 цілі числа у невизначеному потоці, які будуть спостерігатися у обчислювальному потоці. Результати показують, що виконання процесу за замовчуванням відбувається у головному потоці;
  • рядок 23: процес process4 генерує 4 логічні значення у невизначеному потоці, які будуть спостерігатися у невизначеному потоці. Результати показують, що виконання процесу та його спостереження за замовчуванням відбуваються у головному потоці;

Результат виконання цього коду такий:

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]
  • процес process1 генерує 1 дійсне число (рядок 4) у обчислювальному потоці [RxComputationThreadPool-4], яке спостерігається в обчислювальному потоці [RxComputationThreadPool-3] (рядок 6);
  • процес process2 генерує 2 символьні рядки (рядки 12, 14) у обчислювальному потоці [RxComputationThreadPool-5], які спостерігаються в цьому ж потоці (рядки 13, 15);
  • процес process3 генерує 3 цілих числа (рядки 21, 23, 25) у головному потоці, які спостерігаються у обчислювальному потоці [RxComputationThreadPool-6] (рядки 22, 24, 28);
  • процес process4 генерує 4 логічні значення (рядки 34, 36, 38, 40) у головному потоці, які спостерігаються в цьому ж головному потоці (рядки 33, 35, 37, 39);

Пропонуємо читачеві простежити вищезазначене:

  • життєвий цикл спостережуваного процесу та його потоку;
  • життєвий цикл його спостерігача та його потоку;

Значна частина переваг бібліотек Rx полягає саме в цій багатопотоковості, якою розробнику не доводиться займатися самостійно.

7.5. Комбінації декількох спостережуваних величин

7.5.1. Приклад-14: об’єднання двох об’єктів спостереження за допомогою [Observable.merge]

Тепер ми розглянемо статичні методи класу [Observable], що дозволяють об’єднати кілька об’єктів спостереження в один об’єкт спостереження-результат.

Перший приклад такого типу буде таким:


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 {
        // процес 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), Schedulers.computation(), null);
        // об'єднання
        Process<?> process12 = new Process<>("process12",
                Observable.merge(process1.getObservable(), process2.getObservable()));
        // підписки
        ProcessUtils.subscribe(1, process12);
    }
}
  • рядки 15–17: процес із назвою [process1] видасть 3 дійсні числа на обчислювальному потоці. Його також спостерігатимуть на обчислювальному потоці;
  • рядки 19–20: процес із назвою [process2] виведе 2 символьні рядки у обчислювальному потоці. Потік спостереження не задається. Раніше ми бачили, що в цьому випадку потік спостереження є обчислювальним потоком;
  • рядок 23: обидва процеси об’єднуються, тобто створюється об’єкт спостереження, елементи якого надходять одночасно з обох процесів. Для цього використовується статичний метод [Observable.merge]:
 

На відміну від того, що може здатися на наведеному вище схематичному зображенні, під час об’єднання елементи потоку 1 можуть вставлятися між елементами потоку 2. Про це свідчать результати виконання:

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]
  • рядок 3: процес [process1] виконується на обчислювальному потоці [RxComputationThreadPool-4];
  • рядок 4: процес [process2] виконується на обчислювальному потоці [RxComputationThreadPool-5];
  • рядок 9: процес [process12] спостерігається на обчислювальному потоці [RxComputationThreadPool-3]. Мені невідомо, за яким правилом було зроблено цей вибір;
  • рядки 9–11: бачимо, що спостерігач відстежує елементи обох процесів [process1] (рядок 5) та [process2] (рядки 6, 7), хоча жоден із них не завершився (має місце змішування);
  • процес [process12] завершується (рядок 17), коли обидва процеси process1 та process2 завершені;

7.5.2. Приклад-15: об’єднання двох спостережуваних величин за допомогою [Observable.concat]

Тепер розглянемо такий код:


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 {
        // процес 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, Schedulers.computation());
        // об'єднання
        Process<?> process12 = new Process<>("process12",
                Observable.concat(process1.getObservable(), process2.getObservable()));
        // підписки
        ProcessUtils.subscribe(1, process12);
    }
}
  • рядки 15–17: процес із назвою [process1] виведе 3 дійсні числа у обчислювальному потоці. Він також буде спостерігатися в обчислювальному потоці;
  • рядки 19–20: процес із назвою [process2] виведе 2 символьні рядки у невизначений потік, у даному випадку — у головний потік за замовчуванням. Він буде спостерігатися у обчислювальному потоці;
  • рядок 23: обидва процеси об’єднуються, тобто створюється об’єкт спостереження, елементи якого походять з обох процесів. Змішування виведених значень не відбувається. Процес [process12] спочатку виведе всі значення процесу [process1], а потім — значення процесу [process2]. Для цього використовується статичний метод [Observable.concat]:
 

Результати виконання такі:

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]
  • рядки 3–10: виконується процес [process1], а процес [process12] виводить значення, отримані від [process1];
  • рядок 9: процес [process1] завершено;
  • рядки 11–17: виконується процес [process2], а процес [process12] виводить значення, отримані від [process2];

У процесі process2 спостерігається дивна особливість: для нього не було вказано потоку виконання. Отже, можна було б очікувати, що за замовчуванням ним стане головний потік. Однак це не так. Потоком виконання став обчислювальний потік [RxComputationThreadPool-3] (рядок 11). Отже, коли не задано потоку виконання або спостереження, неможливо зробити припущення щодо того, який потік буде обрано.

7.5.3. Приклад-16: об’єднання двох спостережуваних величин за допомогою [Observable.zip]

Тепер розглянемо такий код:


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 {
        // процес 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<String> process2 = new Process<>(
                new ProcessAction01<String>("process2", 2, i -> String.format("valeur-%s", i)), null, null);
        // функція об'єднання 2 процесів
        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");
                }
            }
        };
        // архівування двох процесів
        Process<String> process12 = new Process<>("process12",
                Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));
        // підписки
        ProcessUtils.subscribe(1, process12);
    }
}
  • рядки 16–18: процес із назвою [process1] виведе 3 дійсні числа у обчислювальному потоці. Він також спостерігатиметься в обчислювальному потоці;
  • рядки 20–21: процес із назвою [process2] виведе 2 символьні рядки у невизначеному потоці. Потік спостереження також не визначений;
  • рядки 23–32: інстанціювання типу [FuncN<String>] за допомогою анонімного класу. FuncN є функціональним інтерфейсом:
 

Метод [FuncN.call] очікує масив об’єктів і повертає тип R. Функція [funcn] буде використана для об’єднання процесів process1 та process2 у цьому порядку. У методі [FuncN.call]:

  • args[0] буде Double;
  • args[1] буде String;

У цьому випадку результатом [funcn.call] буде рядок символів із рядка 27. Для побудови цього результату не потрібно знати типи аргументів методу call.

Обидва процеси поєднуються таким чином:


// архівування двох процесів
Process<String> process12 = new Process<>("process12",
Observable.zip(Arrays.asList(process1.getObservable(), process2.getObservable()), funcn));

Метод [Observable.zip] працює наступним чином:

 

Бачимо, що:

  • перший аргумент zip — це Iterable<Observable>. У нашому прикладі ми маємо фактичний параметр типу List<Observable>, сформований з наших двох спостережуваних величин;
  • другий аргумент zip має тип FuncN. У нашому прикладі фактичним параметром є [funcn];

Виконання дає такі результати:

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]
  • рядки 7, 11: процес process12 генерує два елементи;
  • рядок 8: додатковий елемент, виведений процесом process1, який не має партнера в процесі process2, не виводиться процесом-результатом process12;

Бачимо, що процес process2, якому не було призначено ані потоку виконання, ані потоку спостереження, використовував головний потік для обох цілей.

7.5.4. Приклад-17: об’єднання двох спостережуваних величин за допомогою [Observable.combineLatest]

Тепер розглянемо такий код:


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 {
        // процес 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null,
                Schedulers.computation());
        // об'єднання двох процесів
        Process<Double> process12 = new Process<>("process12",
                Observable.combineLatest(process1.getObservable(), process2.getObservable(), (d1, d2) -> d1 + d2));
        // підписки
        ProcessUtils.subscribe(1, process12);
    }
}
  • рядки 14–16: процес із назвою [process1] виведе 3 дійсні числа на обчислювальному потоці. Він також спостерігатиметься на обчислювальному потоці;
  • рядки 18–20: процес із назвою [process2] видасть 2 дійсні числа на невизначеному потоці. Вони будуть спостережуватися на обчислювальному потоці;
  • рядок 23: обидва спостережувані величини об’єднуються за допомогою наступного статичного методу [Observable.combineLatest]:
 

Спостережувана величина [combineLatest] працює наступним чином: коли одна з двох спостережуваних величин генерує елемент E1, цей елемент комбінується за допомогою [combineFunction] з останнім елементом, згенерованим іншою спостережуваною величиною.

Виконання цього коду дає такий результат:

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]
  • рядок 5: передача process2 (56) об'єднується з останнім елементом, переданим process1 (54, рядок 4), і дає результат, наведений у рядку 7;
  • рядок 6: передача process1 (51,6) поєднується з останнім елементом, переданим process2 (56, рядок 5), і дає результат, наведений у рядку 8;
  • рядок 9: передача process2 (261,8) поєднується з останнім елементом, переданим process1 (51,6, рядок 6), і дає результат, наведений у рядку 12;
  • рядок 13: передача від process1 (80,39) поєднується з останнім елементом, переданим process2 (261,8, рядок 9), і дає результат, наведений у рядку 15;

Тут ми маємо справу з варіантом спостережуваного об’єкта [zip], де цього разу об’єднані елементи не обов’язково є елементами, що займають однакову позицію у потоках. Тут слід зауважити, що процес process2, для якого не було задано потоку виконання, у цьому випадку було виконано в головному потоці (рядок 2).

7.5.5. Приклад-18: об’єднання двох спостережуваних величин за допомогою [Observable.amb]

Тепер розглянемо такий код:


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 {
        // процес 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<Double> process2 = new Process<>(
                new ProcessAction01<Double>("process2", 2, i -> new Random().nextInt(200) * 1.4), null, null);
        // об'єднання двох процесів
        Process<Double> process12 = new Process<>("process12",
                Observable.amb(process1.getObservable(), process2.getObservable()));
        // підписки
        ProcessUtils.subscribe(1, process12);
    }
}
  • рядки 14–16: процес із назвою [process1] видасть 3 дійсні числа у обчислювальному потоці. Він також спостерігатиметься в обчислювальному потоці;
  • рядки 18–20: процес із назвою [process2] видасть 2 дійсні числа на невизначеному потоці. Вони будуть спостережуватися на невизначеному потоці;
  • рядок 22: обидва спостережувані величини об’єднуються за допомогою наступного статичного методу [Observable.amb]:
 

Як показано на схемі вище, спостережувана величина [Observable.amb(Observable o1, Observable o2)] генерує елементи спостережуваної величини, яка генерує першою. Це підтверджують результати наведеного прикладу:

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]
  • у рядку 4 першим передає дані процес process2;
  • рядки 8, 12: процес process12 випускає всі елементи, випущені процесом process2 (рядки 4, 11);

7.6. Ланцюг обробки спостережуваного об’єкта

7.6.1. Приклад-19: перетворення спостережуваного за допомогою [Observable.map]

У попередніх прикладах ми розглянули різні комбінації двох спостережуваних величин для отримання третьої спостережуваної величини. Тепер ми розглянемо статичні методи класу [Observable], які дозволяють виконувати операції перетворення, фільтрування та агрегації над спостережуваним об’єктом. Тут ми зустрінемо методи, аналогічні тим, що належать до класу [Stream], які ми вивчали в розділі 5.

Нашим першим прикладом буде наступний:


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 {
        // процес 1
        Process<Double> process1 = new Process<>(
                new ProcessAction01<Double>("process1", 3, i -> new Random().nextInt(100) * 1.2), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<String> process2 = new Process<>("process2",
                process1.getObservable().map(d -> String.format("valeur-%s", d)));
        // підписки
        ProcessUtils.subscribe(1, process2);
    }
}
  • рядки 14–16: процес із назвою process1 видаватиме 3 дійсні числа у обчислювальному потоці. Його також спостерігатимуть у обчислювальному потоці;
  • рядки 17–18: числа, виведені процесом process1, будуть перетворені на символьні рядки у процесі process2;
  • рядок 20: спостерігається process2;

Метод [Observable.map] у рядку 18 аналогічний методу [Stream.map], розглянутому в розділі 5.5:

 

Результати прикладу такі:

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]
  • рядки 4, 5 та 8: випромінювання process1. Це дійсні числа;
  • рядки 6, 7, 10: випромінювання process2, що спостерігаються. Це символьні рядки;

7.6.2. Приклад-20: фільтрування спостережуваної величини за допомогою [Observable.filter]

Приклад буде таким:


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 {
        // процес 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().filter(i -> i % 2 == 0));
        // підписки
        ProcessUtils.subscribe(1, process2);
    }
}
  • рядки 11–12: процес із назвою process1 буде генерувати цілі числа від 0 до 2 у обчислювальному потоці. Його також спостерігатимуть у обчислювальному потоці;
  • рядок 14: числа, виведені процесом process1, будуть відфільтровані, щоб у процесі process2 залишилися лише парні числа;
  • рядок 20: спостерігаємо за process2;

Метод [Observable.filter] у рядку 18 аналогічний методу [Stream.filter], розглянутому в розділі 5.4:

 

Результати прикладу такі:

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]
  • рядки 4, 5 та 7: передачі process1;
  • рядки 6, 9: спостережувані випромінювання process2. Це парні елементи process1;

7.6.3. Приклад-21: перетворення спостережуваної величини за допомогою [Observable.flatMap]

Приклад буде таким:


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 {
        // процес 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<Integer> process2 = new Process<>("process2", process1.getObservable().flatMap(i -> {
            int value = i * 10;
            return Observable.just(value, value + 1, value + 2);
        }));
        // підписки
        ProcessUtils.subscribe(1, process2);
    }
}
  • рядки 12–13: процес із назвою process1 буде генерувати цілі числа від 0 до 2 у обчислювальному потоці. Він також буде спостерігатися в обчислювальному потоці;
  • рядки 15–18: кожне число n, виведене process1, перетворюється на об’єкт спостереження, що виводить 3 числа (10*n, 10*n+1, 10*n+2). Якби в рядку 15 використовувався метод [map], то process2 виводив би тип Observable<Integer>, а не тип Integer. Використаний метод [flatMap] дозволяє звести (flatten) цю послідовність елементів типу Observable<Integer> до послідовності елементів типу Integer, що складається з кожного елемента кожного з Observable<Integer>;
  • рядок 20: спостерігається process2;

Метод [Observable.flatMap] у рядку 15 аналогічний методу [Stream.flatMap], розглянутому в розділі 5.6.12:

 

Результати прикладу такі:

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]
  • рядки 5–7: три передачі process2 після передачі рядка 4 process1;
  • рядки 9–11: три передачі process2 після передачі рядка 8 з process1;
  • рядки 14–16: три передачі process2 після передачі рядка 12 process1;

Наступний код демонструє, як створити тип Observable<Integer[]> на основі process1 та [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 {
        // процес 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<Integer[]> process2 = new Process<>("process2", process1.getObservable().map(i -> {
            int value = i * 10;
            return new Integer[] { value, value + 1, value + 2 };
        }));
        // підписки
        ProcessUtils.subscribe(1, process2);
    }
}
  • рядок 14: використовується метод [Observable.map];
  • рядок 16: який повертає тип Integer[];

Результати такі:

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]
  • рядки 6, 7, 10: бачимо результати map;

Усі ці перетворення спостережуваної величини можна об'єднувати в ланцюжок, оскільки кожне перетворення створює нову спостережувану величину. Це ілюструє наступний приклад [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 {
        // процес 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // процес 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));
        // підписки
        ProcessUtils.subscribe(1, process2);
    }
}
  • рядки 15–18: за flatMap йде filter;

Результати виконання такі:

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]
  • рядки 8–13: process2 вивів лише парні елементи з flatMap;

Метод, схожий на [flatMap], — це метод [flatMapIterable], проілюстрований наступним прикладом [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 {
        // процес 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // процес 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));
        // підписки
        ProcessUtils.subscribe(1, process2);
    }
}

У рядку 16 замість методу [flatMap] використовується метод [flatMapIterable]. У цьому випадку функція перетворення повинна генерувати тип Iterable<T> (рядок 18) замість типу Observable<T>.

Отримуємо ті самі результати, що й раніше.

Повернемося до визначення методу [flatMap]:

 

Як бачимо вище, між двома зеленими елементами [1-2] вставлено синій елемент [3]. Це означає, що під час операції згладжування Observable<T> метод [flatMap] дотримується порядку випромінювання цих різних внутрішніх спостережуваних величин. Це ілюструється наступним прикладом [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 {
        // процес 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // процес 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().flatMap(i -> process2.getObservable()));
        // підписки
        ProcessUtils.subscribe(1, process3);
    }
}
  • рядки 11–12: процес process1 видає цілі числа [0,1];
  • рядки 14–15: процес process2 генерує цілі числа [10,11,12];
  • рядки 17–18: кожному елементу, що генерується процесом process1, прив’язується спостережувана величина процесу process2. Це означає, що:
    • елементу [0] процесу process1 буде прив’язана спостережувана величина, що генерує [10,11,12];
    • те саме стосується елемента 1;

У підсумку буде згенеровано 6 чисел [10, 11, 12, 10, 11, 12]. Ми хочемо побачити, в якому порядку.

Результати виконання такі:

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]

Бачимо, що порядок виведення процесу process3 був таким: [10, 10, 11, 12, 11, 12] (рядки 11, 12, 14, 17, 19, 22). Отже, дійсно відбулося змішування елементів, виведених процесом process2. Цього можна уникнути, використовуючи метод [concatMap] замість методу [flatMap]. Це ілюструє наступний код [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 {
        // процес 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // процес 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().concatMap(i -> process2.getObservable()));
        // підписки
        ProcessUtils.subscribe(1, process3);
    }
}

У рядку 18 [flatMap] замінено на [concatMap]. Результати виконання такі:

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]

Видно, що порядок виведення процесу process3 був таким: [10, 11, 12, 10, 11, 12] (рядки 12–14, 17, 19, 22). Елементи, виведені процесом process2, не були змішані.

Іншим варіантом методу [map] є метод [switchMap]:

 

Вищезазначений об’єкт спостереження [1] породжує ще 3 об’єкти спостереження [2], що складаються з 2 елементів, які потім згладжуються, як у [flatMap] та [3]. Можна помітити, що результат має 5 елементів, а не 6. Це пов’язано з тим, що до того, як друга спостережувана величина випустила свій елемент № 2 [6], третя спостережувана величина випустила свій перший елемент [5], внаслідок чого друга спостережувана величина відкидається. Отже, елемент [6] не зустрічається в спостережуваному результаті [3].

Для ілюстрації [switchMap] ми використаємо такий приклад [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 {
        // процес 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 2, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<Integer> process2 = new Process<>(new ProcessAction01<>("process2", 3, i -> i + 10),
                Schedulers.computation(), Schedulers.computation());
        // процес 3
        Process<Integer> process3 = new Process<>("process3",
                process1.getObservable().switchMap(i -> process2.getObservable()));
        // підписки
        ProcessUtils.subscribe(1, process3);
    }
}

Виконання прикладу дає такі результати:

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 генерує 2 елементи, які утворюють 2 спостережувані об’єкти process2, що складаються з 3 елементів;
  • рядок 14: спостерігач отримує елемент № 0, випущений першим спостережуваним об’єктом process2 у рядку 6;
  • рядок 15: спостерігач отримує елемент № 0, переданий другим спостережуваним об’єктом process2 у рядку 13. Невідомо, чому він раніше не отримав елементи 1 і 2, надіслані першим об’єктом спостереження process2 у рядках 7 і 8. У будь-якому разі перший об’єкт спостереження process2 відкидається;
  • у підсумку спостерігач бачить лише 4 елементи (рядки 14, 15, 17, 20) замість 6, які були надіслані;

7.6.4. Приклади-22: інші методи класу [Observable]

Клас [Observable] переймає багато методів класу [Stream], які працюють аналогічно. Ось декілька з них. Ми обмежимося наведенням коду та його результатів.

[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 {
        // процес
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).take(3));
        // підписки
        ProcessUtils.subscribe(1, process);
    }
}

результати

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 {
        // процес
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).takeLast(2));
        // підписки
        ProcessUtils.subscribe(1, process);
    }
}

результати

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 {
        // процеси
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).skip(5).take(2));
        // підписки
        ProcessUtils.subscribe(1, process);
    }
}

результати

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 {
        // процеси
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).reduce(0, (i, a) -> i + a));
        // підписки
        ProcessUtils.subscribe(1, process);
    }
}
  • рядок 10: обчислює суму елементів спостережуваної величини. Результатом є спостережувана величина, яка випромінює цю суму;

результати

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 {
        // процеси
        Process<Boolean> process = new Process<>("process", Observable.range(1, 10).all(i -> i > 10));
        // підписки
        ProcessUtils.subscribe(1, process);
    }
}
  • рядок 10: повертає Observable<Boolean>, який випромінює елемент true, якщо предикат методу [all] є істинним для всіх елементів, інакше — false;

результати

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 {
        // процеси
        Process<Integer> process = new Process<>("process", Observable.range(1, 10).count());
        // підписки
        ProcessUtils.subscribe(1, process);
    }
}
  • рядок 10: [Observable.count] створює спостережуваний об’єкт з 1 елементом, який є сумою спостережуваних елементів;

результати

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 {
        // процеси
        Process<Integer> process = new Process<>("process", Observable.just(1, 2, 1, 3).distinct());
        // підписки
        ProcessUtils.subscribe(1, process);
    }
}

результати

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 {
        // процеси
        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()));
        // підписки
        ProcessUtils.subscribe(1, process);
    }
}
  • рядок 11: метод [groupBy] об'єднує 10 виведених елементів у 2 групи: парні та непарні числа. Результатом є тип Observable<GroupedObservable<Boolean, Integer>>, тобто спостережуваний об’єкт, елементи якого мають тип GroupedObservable<Boolean, Integer>, де Boolean — це тип ключа групи (тут — false, true), який також є типом результату лямбда-виразу, переданого як параметр методу [groupBy], а Integer — типом елементів групи;
  • рядок 12: тип GroupedObservable має метод [asObservable], який дозволяє створити спостережуваний об’єкт на основі цього типу. Отже, ми отримаємо 2 типи Observable<Integer>: один для парних чисел, інший — для непарних. З цих двох об’єктів спостереження метод [concatMap] створить один;

результати

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 {
        // процес 1
        Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
                Schedulers.computation());
        // процес 2
        Process<Timestamped<Integer>> process2 = new Process<>("process2", process1.getObservable().timestamp());
        // підписки
        ProcessUtils.subscribe(1, process2);
    }
}
  • у рядку 15 метод [timestamp] прив’язує час до кожного елемента оброблюваної величини;

результати

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]

У цьому прикладі важко сказати, що саме означає інформація timestamp:

  • рядки 4–5: бачимо, що елемент 1 з process1 був переданий через 139 мс після елемента 0;
  • рядки 6 і 7: бачимо, що елемент 1 з process2 було зафіксовано через 234 мс після елемента 0;
  • рядки 5, 8: бачимо, що елемент 2 з process1 був переданий через 33 мс після елемента 1;
  • рядки 7 і 10: бачимо, що елемент 2 з process2 було спостережено через 37 мс після елемента 1;

Ці зсуви зумовлені тим, що потоки спостереження та виконання спостережуваних величин не збігаються. Якщо замінити рядки 12–13 на такі (Приклад 22j):


// процес 1
Process<Integer> process1 = new Process<>(new ProcessAction01<>("process1", 3, i -> i), Schedulers.computation(),
null);
  • рядки 2–3: не задається потік спостереження. Відомо, що в цьому випадку спостережувана величина спостерігається там, де вона виконується;

Це дає такі результати:

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]
  • рядки 4 та 6: процес process1 надсилає свій елемент № 1 через 587 мс після елемента № 0;
  • рядки 5 і 7: спостерігач спостерігає ці 2 елементи з інтервалом у 586 мс;
  • рядки 6 і 8: процес process1 надсилає свій елемент № 2 через 396 мс після елемента № 1;
  • рядки 7 і 9: спостерігач фіксує ці 2 елементи з інтервалом у 396 мс;

У цьому випадку значення timestamp є узгодженими: вони дійсно відображають дату надсилання елемента.

7.7. Планувальники

7.7.1. Приклад-23: планувальник [Schedulers.computation]

Тепер розглянемо планувальники виконання. Спостереження проводитиметься на потоці виконання.

Тема планувальників є дещо складною. Різні планувальники представлені в цьому питанні на сайті StackOverflow [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases]:

 

Ми спробуємо проілюструвати використання цих різних планувальників на прикладах. Перший приклад ілюструє планувальник [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 {
        // процеси
        @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);
        }
        // підписки
        ProcessUtils.subscribe(1, processes);
    }
}
  • рядки 14–19: створюється масив із 10 процесів, що виконуються в одному обчислювальному потоці;
  • рядок 17: кожен процес генерує випадкове дійсне число;
  • рядок 21: підписуємося на всі ці процеси;

Результати такі:

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]
  • рядки 2–10: перші 8 процесів запускаються на 8 різних потоках (використовувана машина має 8 ядер). Можна помітити, що всі вони запускаються приблизно одночасно;
  • рядки 17–19: 3 процеси завершуються, звільняючи таким чином 3 потоки;
  • рядки 23–24: два останні процеси можуть запуститися, зайнявши 2 з цих звільнених потоків;

Отже, слід запам’ятати, що планувальник [Schedulers.computation] надає пул із n потоків, де n — кількість ядер машини. Потоки виконуються паралельно на цих ядрах.

7.7.2. Приклад-24: планувальник [Schedulers.io]

Ми запускаємо попередній код із планувальником [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 {
        // процеси
        @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);
        }
        // підписки
        ProcessUtils.subscribe(1, processes);
    }
}
  • рядок 18: процеси виконуються з потоками планувальника [Schedulers.io];

Це дає такі результати:

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]
  • рядки 2–10: кожний із 10 процесів запускається в окремому потоці. На відміну від попереднього випадку, всі процеси вдалося запустити. Можна помітити, що ці запуски тривають 6 мс, тоді як раніше це займало 1 мс;
  • рядки 13–18: спостережувані величини генерують дані по черзі, а не майже паралельно, як це було раніше;

У чому полягає різниця між планувальниками [Schedulers.io] та [Schedulers.computation]? Відповідь можна знайти в URL та [http://stackoverflow.com/questions/31276164/rxjava-schedulers-use-cases]:

 

7.7.3. Приклад-25: планувальник [Schedulers.newThread]

Ми запускаємо попередній код із планувальником [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 {
        // процеси
        @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);
        }
        // підписки
        ProcessUtils.subscribe(1, processes);
    }
}

Отримані результати збігаються з результатами для планувальника [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]

У розділі URL [http://stackoverflow.com/questions/33415881/retrofit-with-rxjava-schedulers-newthread-vs-schedulers-io] пояснюється, що планувальник [Schedulers.io] надає пул потоків, чого не робить планувальник [Schedulers.newThread]. Пул потоків автоматично створює n потоків. Він розподіляє їх між процесами, які їх потребують. Коли ці процеси завершуються, їхні потоки не видаляються, а повертаються до пулу і можуть бути повторно використані іншим процесом. Це економічніше, ніж постійно створювати та видаляти потоки. Отже, можна вважати, що краще використовувати планувальник [Schedulers.io].

7.7.4. Приклад-26: планувальники [Schedulers.immediate, Schedulers.trampoline]

Повернемося до пояснення, наданого для цих двох планувальників:

 

Пояснення досить просте для розуміння, але коли хочеш його проілюструвати, виявляється, що ти його не зрозумів. Саме книга [Learning Reactive Programming With Java 8] дозволила мені створити приклад, який повторює приклад, знайдений у цій книзі, але спрощує його. Ось він:


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 {

        // планувальник
        Scheduler scheduler = Schedulers.immediate();
        // робочий процес цього планувальника
        Worker worker = scheduler.createWorker();
        // тип Action0, що має виконуватися на робочому процесі
        Action0 action02 = new Action0() {
            @Override
            public void call() {
                // журнал Action02
                ProcessUtils.showInfos.accept("action02");
            }
        };

        // тип Action0, що має виконуватися на робочому процесі
        Action0 action01 = new Action0() {
            @Override
            public void call() {
                // заплановано нову дію на тому самому робочому процесі
                worker.schedule(action02);
                // журнал дії action01
                ProcessUtils.showInfos.accept("action01");
            }
        };
        // action01 заплановано на робочому процесі
        worker.schedule(action01);
    }

    // перегляди
    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()));

}
  • рядок 17: планувальник. Це буде або [Schedulers.immediate], як тут, або [Schedulers.trampoline] пізніше;
  • рядок 19: можна запускати дії типу Action0 (рядки 21, 20) на робочих процесах планувальника. Метод [Scheduler.createWorker] дозволяє створити робочий процес. Метод [Worker.schedule(Action0)] дозволяє запустити дію типу Action0 на робочому процесі;
  • рядки 21–27: перша дія з назвою [action02], яка буде виконана (рядок 40) робочим процесом із рядка 19;
  • рядки 30–38: друга дія з назвою [action01]. Її особливістю є те, що вона запускає дію action02 на тому самому робочому процесі, що й вона сама (рядок 34). Саме в цьому полягає різниця між [Schedulers.immediate] та [Schedulers.trampoline]:
    • якщо планувальником є [Schedulers.immediate], то в рядку 34 дія action02 буде виконана негайно (звідси й назва планувальника), а дія action01, що виконується в цей момент, буде перервана. Після цього з’явиться повідомлення з рядка 25. Після завершення дії action02 дія action01 відновиться, і з’явиться повідомлення з рядка 36;
    • якщо планувальником є [Schedulers.trampoline], то, як показано в рядку 34, дія action02 ставиться в режим очікування. Вона буде виконана лише після завершення поточного завдання action01. Тоді з’явиться повідомлення з рядка 36. Після завершення дії action01 буде виконано дію action02, і ми побачимо повідомлення з рядка 25;

Виконання наведеного вище коду дає такі результати:

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

Якщо в рядку 17 використовувати планувальник [Schedulers.trampoline], отримаємо протилежні результати:

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

Однак важко встановити зв’язок із спостережуваними величинами. Я не знайшов переконливого прикладу, який би продемонстрував переваги виконання спостережуваної величини в одному з цих двох потоків. Проте ось один приклад, який, на мою думку, зовсім не виглядає природним:


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 = Schedulers.immediate().createWorker();
        // робочий процес worker = Schedulers.trampoline().createWorker();
        // спостережуваний об’єкт 1 на робочому процесі
        worker.schedule(() -> Observable.range(1, 2).subscribe(new Action1<Integer>() {

            @Override
            public void call(Integer i) {
                ProcessUtils.showInfos.accept(String.valueOf(i));
                // спостережуваний об’єкт 2 на тому самому робочому процесі
                worker.schedule(() -> Observable.range(100, 2).subscribe(new Action1<Integer>() {
                    @Override
                    public void call(Integer i) {
                        ProcessUtils.showInfos.accept(String.valueOf(i));
                    }
                }));
            }
        }));
    }
}
  • рядки 13–14: створюється робочий процес на основі одного з двох планувальників [Schedulers.immediate] та [Schedulers.trampoline];
  • рядок 16: на цьому робочому процесі заплановано перший об’єкт спостереження obs1 для генерації чисел [1,2]
  • рядок 22: щоразу, коли спостерігається елемент цього спостережуваного об’єкта obs1, на тому самому робочому процесі запускається спостереження за другим спостережуваним об’єктом obs2 для генерації чисел [100,101];

За допомогою планувальника [Schedulers.immediate] отримуємо такі результати:

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]

Тоді як із планувальником [Schedulers.trampoline] отримуємо такі результати:

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

Попереду ще багато роботи. Щоб глибше ознайомитися з бібліотекою RxJava, читачеві пропонується продовжити навчання за посиланнями, наведеними на початку цього документа. Проте ми маємо основи для використання RxJava у середовищах Swing та Android. Саме це ми й продемонструємо зараз.