Skip to content

8. RxJava у середовищі Swing

8.1. Introduction

Тут ми повернемося до програми Swing, представленої в розділі 2.

  

Для роботи з RxJava у середовищі Swing ми використовуватимемо бібліотеку RxSwing, яка доповнює RxJava класами та інтерфейсами, корисними у середовищі Swing. Для цього файл Gradle прикладу Swing має такий вигляд:

  

buildscript {
    repositories {
        mavenCentral()
    }
}
apply plugin: 'java'
jar {
    baseName = 'exemples-01'
    version = '0.0.1-SNAPSHOT'
}
repositories {
    mavenCentral()
}
dependencies {
    compile('io.reactivex:rxswing:0.25.0')
    compile('io.reactivex:rxjava:1.1.3')
    compile('com.fasterxml.jackson.core:jackson-databind:2.7.3')
}
task wrapper(type: Wrapper) {
    gradleVersion = '2.9'
}
  • рядок 15: залежність від RxSwing;

Ми будемо використовувати лише один об’єкт, властивий RxSwing: планувальник [SwingScheduler.getInstance()], який забезпечує виконання/спостереження за спостережуваними величинами у потоці циклу подій Swing. Ми будемо використовувати його виключно для спостереження за об’єктами спостереження, що виконуються в інших потоках, відмінних від потоку циклу подій. Нагадаємо архітектуру прикладу додатка:

Image

  • асинхронний сервісний рівень містить методи, які повертають спостережувані об’єкти. Ми виконуємо ці спостережувані об’єкти в потоках, відмінних від потоку циклу подій. Таким чином, графічний інтерфейс не застигає. Він може реагувати на дії користувача. Найбільш очевидним є надання користувачеві можливості натиснути кнопку [Annuler], щоб перервати занадто тривалу асинхронну операцію. Щоб це стало можливим, графічний інтерфейс повинен бути зафіксованим (frozen);
  • шар Swing має обробляти результати, повернуті асинхронними операціями, і на їх основі оновлювати графічний інтерфейс. Однак це можна зробити лише у потоці циклу подій. Для цього ці результати відстежуються у планувальнику [SwingScheduler.getInstance()];

Отже, у коді обробки подій графічного інтерфейсу взаємодія з асинхронним рівнем [rxService] відбувається у такій формі:


Observable obs=rxService.doSomething(...).subscribeOn(Schedulers.computation()).observeOn(SwingScheduler.getInstance()) ;

де планувальник [Schedulers.computation()] може бути замінений іншим планувальником залежно від конкретного випадку використання.

Читачеві пропонується ще раз перечитати параграф 2. Тепер він має знання, необхідні для його повного розуміння.

8.2. Структура коду

Код реалізує таку архітектуру:

Image

Проєкт IntelliJ IDEA, що реалізує цю архітектуру, має такий вигляд:

  
  • пакет [rxswing.service] реалізує синхронні (IService, Service) та асинхронні (IRxService, RxService) рівні сервісів;
  • пакет [rxswing.ui] реалізує інтерфейс Swing;

8.3. Виконання проєкту

Щоб запустити проект у IntelliJ IDEA, виконайте такі дії:

 

8.4. Синхронний сервіс

Image

  

Рівень синхронних сервісів має такий інтерфейс [IService]:


package dvp.rxswing.service;

public interface IService {
  // випадкові числа в інтервалі [a,b]
  // n чисел генерується, де n — саме випадкове число в інтервалі [minCount, maxCount]
  // числа генеруються після затримки в delay мілісекунд,
  // де [delay] — це випадкове число в інтервалі [minDelay, maxDelay]
  public ServiceResponse getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay);
}

Тип відповіді служби [ServiceResponse] є таким:


package dvp.rxswing.service;

import java.util.List;

public class ServiceResponse {

  // час очікування служби
  private int delay;
  // випадкові числа
  private List<Integer> aleas;
  // потік виконання
  private String executedOn;

  // конструктори

  public ServiceResponse() {
      // потік виконання
    executedOn = Thread.currentThread().getName();
  }

  public ServiceResponse(int delay, List<Integer> aleas) {
      // локальний конструктор
    this();
    // інші ініціалізації
    this.delay = delay;
    this.aleas = aleas;
  }

  // методи getter та setter
...
}

Інтерфейс [IService] реалізовано наступним класом [Service]:


package dvp.rxswing.service;

import java.util.*;

public class Service implements IService {

  @Override
  public ServiceResponse getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay) {
    // випадкові числа в інтервалі [a,b]
    // n чисел генерується, де n — саме випадкове число в інтервалі [minCount, maxCount]
    // числа генеруються після затримки в delay мілісекунд,
    // де [delay] — це випадкове число в інтервалі [minDelay, maxDelay]

    // деякі перевірки
    List<String> messages = new ArrayList<>();
    int erreur = 0;
    if (a < 0) {
      messages.add("Le nombre a de l'intervalle [a,b] de génération doit être supérieur à 0");
      erreur |= 2;
    }
    if (a >= b) {
      messages.add("Dans l'intervalle [a,b] de génération, on doit avoir a< b");
      erreur |= 4;
    }
    if (minCount < 0) {
      messages.add("Le nombre min de l'intervalle [min,count] du nombre de valeurs générées doit être supérieur à 0");
      erreur |= 16;
    }
    if (minCount > maxCount) {
      messages.add("Dans l'intervalle [min,count] du nombre de valeurs générées, on doit avoir min<= max");
      erreur |= 32;
    }
    if (minDelay < 0) {
      messages.add("Le nombre min de l'intervalle [min,count] du délai d'attente doit être supérieur à 0");
      erreur |= 64;
    }
    if (minCount > maxCount) {
      messages.add("Dans l'intervalle [min,count] du délai d'attente, on doit avoir min<= max");
      erreur |= 128;
    }
    if (maxDelay > 5000) {
      messages.add("L'attente en millisecondes avant la génération des nombres doit être dans l'intervalle [0,5000]");
      erreur |= 256;
    }
    // помилки?
    if (!messages.isEmpty()) {
      throw new AleasException(String.join(" [---] ", messages), erreur);
    }
    // генератор випадкових чисел
    Random random = new Random();
    // очікування?
    int delay = minDelay + random.nextInt(maxDelay - minDelay + 1);
    if (delay > 0) {
      try {
        Thread.sleep(delay);
      } catch (InterruptedException e) {
        throw new AleasException(String.format("[%s : %s]", e.getClass().getName(), e.getMessage()), 1024);
      }
    }
    // генерація результату
    int count = minCount + random.nextInt(maxCount - minCount + 1);
    List<Integer> nombres = new ArrayList<>();
    for (int i = 0; i < count; i++) {
      nombres.add(a + random.nextInt(b - a + 1));
    }
    // повернення результату
    return new ServiceResponse(delay,nombres);
  }

}

Клас винятків [AleasException], який використовується службою, має такий вигляд:


package dvp.rxswing.service;

public class AleasException extends RuntimeException {

    private static final long serialVersionUID = 1L;
    // код помилки
  private int code;

  // конструктори
  public AleasException() {
  }

  public AleasException(String detailMessage, int code) {
    super(detailMessage);
    this.code = code;
  }

  public AleasException(Throwable throwable, int code) {
    super(throwable);
    this.code = code;
  }

  public AleasException(String detailMessage, Throwable throwable, int code) {
    super(detailMessage, throwable);
    this.code = code;
  }

  // гетери та сеттери
...
}
  • рядок 3: він розширює клас [RuntimeException]. Отже, це неконтрольоване виключення;
  • рядок 7: вона доповнює свій батьківський клас кодом помилки (0 = помилки немає);

8.5. Асинхронний сервіс

Image

  

Рівень асинхронних служб має такий інтерфейс [IRxService]:


package dvp.rxswing.service;

import dvp.rxswing.ui.UiResponse;
import rx.Observable;

public interface IRxService {
  // випадкові числа в інтервалі [a,b]
  // n чисел генерується, де n — саме випадкове число в інтервалі [minCount, maxCount]
  // числа генеруються після затримки в delay мілісекунд,
  // де [delay] — це випадкове число в інтервалі [minDelay, maxDelay]
  public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse);
}
  • рядок 11: тепер метод [getAleas] сервісу повертає об’єкт-спостерігач;

Метод [getAleas] повертає відповідь типу [UiResponse], призначену для шару [Ui]. Цей тип має такий вигляд:


package dvp.rxswing.ui;

import dvp.rxswing.service.ServiceResponse;

import java.text.SimpleDateFormat;
import java.util.Calendar;

public class UiResponse {

  // ідентифікатор клієнта
  private int idClient;
  // відповідь сервісу
  private ServiceResponse serviceResponse;
  // ім'я потоку спостереження
  private String observedOn;
  // час запиту
  private String requestAt;
  // час відповіді
  private String responseAt;

  // конструктори

  public UiResponse() {
      // потік спостереження
    observedOn = Thread.currentThread().getName();
    // час запиту
    requestAt = getTimeStamp();
  }

  // приватні методи

  private String getTimeStamp() {
    return new SimpleDateFormat("hh:mm:ss:SSS").format(Calendar.getInstance().getTime());
  }

  // гетери та сеттери
...
}
  • випадкові числа містяться у полі рядка 13;
  • інші поля призначені для вказівки потоків виконання та спостереження за спостережуваною величиною асинхронної служби, а також часу надсилання запиту до служби та отримання відповіді;

Асинхронний інтерфейс реалізовано за допомогою наступного класу [RxService]:


package dvp.rxswing.service;

import dvp.rxswing.ui.UiResponse;
import rx.Observable;

public class RxService implements IRxService {

  // синхронний сервіс
  private IService service;

  // конструктор
  public RxService(IService service) {
    this.service = service;
  }

  @Override
  public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse) {
      // створюється об’єкт-спостерігач, що передає значення, повернене синхронним сервісом
    return Observable.create(subscriber -> {
      try {
          // синхронний виклик
        uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
        // результат передається спостерігачеві
        subscriber.onNext(uiResponse);
      } catch (Exception e) {
          // подаємо помилку спостерігачеві
        subscriber.onError(e);
      } finally {
          // сповіщаємо спостерігача про завершення передачі даних
        subscriber.onCompleted();
      }
    });
  }
}
  • рядки 12–14: клас [RxService] асинхронної служби створюється на основі екземпляра синхронного інтерфейсу [IService];
  • рядки 20–33: створення об’єкта спостереження, що є результатом методу [getAleas];
  • рядок 22: викликається синхронний метод [service.getAleas]. Його результат типу [ServiceResponse] включається в об’єкт типу [UiResponse], який має бути переданий на рівень [swing]. Цей об’єкт спочатку був переданий у параметрах виклику методу (останній параметр, рядок 17);
  • рядок 24: відповідь [UiResponse] надсилається спостерігачеві (рівень [swing]). Об’єкт [UiResponse] містить не лише інформацію, сформовану синхронною службою у рядку 22. Він також містить іншу інформацію, сформовану методом, що викликає метод [getAleas] у рядку 17. Саме з цієї причини цей викликаючий метод передав об’єкт [UiResponse] як параметр методу [getAleas] (останній параметр, рядок 17);
  • рядок 30: не забуваємо повідомити про закінчення передач. Тут маємо спостережуваний об’єкт, який видає лише одне значення: те, що повертає синхронна служба;
  • рядок 27: сповіщаємо спостерігача про можливу помилку;

8.6. Графічний інтерфейс

Image

  
  • графічний інтерфейс було створено за допомогою IDE [Netbeans], який має хороший графічний редактор. Цей редактор згенерував файл [AbstractJFrameAleas.form], який може використовуватися лише цим IDE;
  • клас [AbstractJFrameAleas] також було згенеровано графічним редактором NetBeans. Потім його було рефакторовано наступним чином: події графічного інтерфейсу, які ми хотіли обробляти, обробляються у класі [AbstractJFrameAleas] за допомогою абстрактних методів, реалізованих у дочірньому класі [JFrameAleasEvents]. У підсумку,
    • абстрактний клас [AbstractJFrameAleas] відповідає за побудову та відображення графічного інтерфейсу;
    • дочірній клас [JFrameAleasEvents] відповідає за обробку подій цього інтерфейсу;

Компоненти графічного інтерфейсу вкладки [Request] такі:

 
тип
назва
роль
1
JTabbedPane
jTabbedPane1
контейнер з вкладками. Містить дві вкладки (JPanel) [jPanelRequest] для запиту, [jPanelresponse] для відповіді;
2
JTextField
jTextFieldNbValeurs
кількість запитів, які потрібно надіслати до служби генерації випадкових чисел. У разі асинхронної служби, що виконується на планувальнику [Schedulers.io], ці запити будуть спільно використовувати один процесор;
3
JTextField
jTextFieldA
межа a інтервалу [a,b]
4
JTextField
jTextFieldB
точка b на відрізку [a,b]
5
JTextField
jTextFieldMinCount
клем minCount інтервалу [minCount, maxCount]
6
JTextField
jTextFieldMaxCount
клем maxCount з діапазону [minCount, maxCount]
7
JTextField
jTextFieldMinDelay
клем minDelay з діапазону [minDelay, maxDelay]
8
JTextField
jTextFieldMaxDelay
клем maxDelay з діапазону [minDelay, maxDelay]
9
JCheckBox
jCheckBoxRxSwing
якщо поле відмічено, запити надсилаються до асинхронного інтерфейсу. В іншому випадку вони надсилаються до синхронного інтерфейсу
10
JComboBox
jComboBoxSchedulers
у разі асинхронних запитів вони будуть виконуватися за допомогою обраного тут планувальника
11
JButton
jButtonGenerate
запускає виконання запитів у синхронному або асинхронному режимі

Компоненти графічного інтерфейсу вкладки [Response] такі:

 
тип
назва
роль
1
JLabel
jLabelDuree
загальний час виконання запитів у мілісекундах
2
JLabel
jLabelNbReponses
загальна кількість спостережуваних відповідей (може відрізнятися від кількості запитів, оскільки кожен запит може надавати кілька значень для спостереження)
3
JList
jListNumbers
відображення спостережуваних (отриманих) значень
4
JButton
jButtonAnnuler
скасовує запити, що виконуються

8.7. Ініціалізація графічного інтерфейсу

  

Клас [JFrameAleasEvents] обробляє події графічного інтерфейсу, зокрема натискання кнопки [Générer]. Це виконуваний клас, який запускається в такому контексті:


public class JFrameAleasEvents extends AbstractJFrameAleas {

    private static final long serialVersionUID = 1L;
    // синхронна служба генерації
    private IService service;
    // служба асинхронного генерації
    private IRxService rxService;

    // записи
    private int nbRequests;
    private int a;
    private int b;
    private int minDelay;
    private int maxDelay;
    private int minCount;
    private int maxCount;

    // повідомлення про помилки
    private final String jLabelNbValuesErrorText = "Tapez un nombre entier >=1";
    private final String jLabelCountErrorText = "minCount doit être >=0 et maxCount>=minCount ";
    private final String jLabelDelayErrorText = "minDelay doit être >=0 et maxDelay>=minDelay et  maxDelay<=5000";
    private final String jLabelIntervalErrorText = "a doit être >=0 et b>=a ";

    // підписки на спостережувані величини
    protected List<Subscription> subscriptions = new ArrayList<Subscription>();
    // початок-кінець виконання
    private long debut;
    // маппер jSON
    private ObjectMapper jsonMapper;
    // модель відповідей
    private DefaultListModel<String> model;

    // конструктор
    public JFrameAleasEvents() {
        // батьківський елемент
        super();
        // локальний
        initJFrame();
        // послуги
        service = new Service();
        rxService = new RxService(service);
        // маппер jSON
        jsonMapper = new ObjectMapper();
    }

    private void initJFrame() {
        // приховувати повідомлення про помилки
        jLabelCountError.setText("");
        jLabelDelayError.setText("");
        jLabelIntervalError.setText("");
        jLabelNbValuesError.setText("");
        // за замовчуванням приховувати тексти
        jTextFieldA.setText("100");
        jTextFieldB.setText("200");
        jTextFieldMinCount.setText("5");
        jTextFieldMaxCount.setText("10");
        jTextFieldMinDelay.setText("100");
        jTextFieldMaxDelay.setText("500");
        jTextFieldNbValeurs.setText("10");
        jLabelDuree.setText("");
        // шаблон відповідей
        model = new DefaultListModel<>();
        jListNumbers.setModel(model);
        // кількість ядер
        System.out.printf("La JVM a détecté [%s] coeurs sur votre machine%n", Runtime.getRuntime().availableProcessors());
    }

    public static void main(String args[]) {
        try {
            UIManager.setLookAndFeel(UIManager.getSystemLookAndFeelClassName());
        } catch (UnsupportedLookAndFeelException | ClassNotFoundException | InstantiationException
                | IllegalAccessException e) {
            System.out.println(e);
            System.exit(0);
        }

        /* Створити та відобразити форму */
        java.awt.EventQueue.invokeLater(() -> {
            new JFrameAleasEvents().setVisible(true);
        });
    }
  • рядок 1: клас [JFrameAleasEvents] успадковує клас [AbstractJFrameAleas], який, у свою чергу, успадковує клас Swing [JFrame]. Отже, клас [JFrameAleasEvents] є вікном Swing;
  • рядки 68–75: метод [main], який буде виконано;
  • рядок 70: встановлює зовнішній вигляд графічного інтерфейсу;
  • рядок 79: викликається конструктор класу [JFrameAleasEvents]: графічний інтерфейс буде побудовано та ініціалізовано. Після цього він стає видимим;
  • рядки 34–44: конструктор;
  • рядок 36: виклик конструктора батьківського класу ініціалізує графічний інтерфейс. На цей момент він відповідає тому, як його намалював розробник. Він ще не видимий;
  • рядок 38: ініціалізуються деякі компоненти графічного інтерфейсу;
  • рядок 40: створення екземпляра синхронного сервісу;
  • рядок 41: створення екземпляра асинхронного сервісу;

8.8. Виконання синхронних запитів

Натискання кнопки [Générer] призводить до виконання наступного методу [doGenerate]:


    @Override
    protected void doGenerate() {
        // дані введені правильно?
        if (!isPageValid()) {
            return;
        }
        // rx чи ні?
        if (jCheckBoxRxSwing.isSelected()) {
            // асинхронні запити
            doGenerateWithRxService();
        } else {
            // синхронні запити
            doGenerateWithService();
        }
}
  • рядки 4–6: перевіряється правильність введених користувачем даних. Метод [isPageValid] ми не коментуватимемо. Він є базовим;
  • рядок 8: перевіряється стан прапорця RxSwing;
  • рядок 13: запити виконуються синхронно;

Метод [doGenerateWithService] виглядає наступним чином:


    // синхронне генерування
    private void doGenerateWithService() {
        // початок очікування
        beginWaiting();
        try {
            for (int i = 0; i < nbRequests; i++) {
                // підготовка відповіді
                UiResponse uiResponse = new UiResponse();
                // номер клієнта
                uiResponse.setIdClient(i);
                // синхронний виклик
                uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
                // час відповіді
                uiResponse.setResponseAt();
                // оновлення шаблону JList з урахуванням отриманих відповідей
                model.add(0, jsonMapper.writeValueAsString(uiResponse));
                // оновлення кількості відповідей
                jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
            }
        } catch (JsonProcessingException | RuntimeException e) {
            JOptionPane.showMessageDialog(this, getInfoForThrowable("L'erreur suivante s'est produite", e), "Informations",
                    JOptionPane.PLAIN_MESSAGE);
        }
        // завершення очікування
        endWaiting();
}
  • рядок 12: синхронний виклик служби генерації випадкових чисел;
  • виконання методу [doGenerateWithService] відбувається повністю у потоці циклу подій Swing. Доки метод не завершиться, графічний інтерфейс не обробляє жодних нових подій. Він залишається застиглим (frozen). Так, наприклад, оновлення графічного інтерфейсу в рядках 16 і 18 ніколи не будуть відображені. Вони стануть видимими лише з їхніми кінцевими значеннями, і це відбудеться після завершення виконання всіх запитів;

Метод [beginWaiting] (рядок 4) має такий вигляд:


    private void beginWaiting() {
        // кнопки
        jButtonGenerate.setVisible(false);
        jButtonCancel.setVisible(true);
        // курсор очікування
        jTabbedPane1.setCursor(Cursor.getPredefinedCursor(Cursor.WAIT_CURSOR));
        jButtonCancel.setCursor(Cursor.getDefaultCursor());
        // очистити відповіді
        model.clear();
        // підписки Rx
        subscriptions.clear();
        // відображення перегляду відповідей
        jTabbedPane1.setSelectedIndex(1);
        jLabelNbReponses.setText("0");
        jLabelDuree.setText("");
        // початок виконання
        debut = new Date().getTime();
}
  • рядок 3: кнопка [Générer] прихована. Це створює подію, яка також зможе бути виконана лише після завершення виконання всіх запитів. Тому ми ніколи не бачимо його прихованим, оскільки метод [endWaiting] у рядку 25 методу [doGenerateWithService] знову його відображає;
  • рядок 13: вибираємо вкладку [Response], щоб побачити, як надходять відповіді. Знову ж таки, ця подія буде виконана лише після завершення виконання всіх запитів, і тоді ми побачимо всі відповіді разом, хоча хотіли бачити їх по черзі;

Синхронний інтерфейс має очевидні недоліки. Їх можна подолати за допомогою асинхронного інтерфейсу.

8.9. Виконання асинхронних запитів

Код виконання асинхронних запитів виглядає так:


private void doGenerateWithRxService() {
        // початок очікування
        beginWaiting();
        // отримаємо випадкові числа у вигляді спостережуваної величини
        Observable<UiResponse> observable = Observable.empty();
        // Планувальник виконання різних спостережуваних величин
        Scheduler[] schedulers = { Schedulers.io(), Schedulers.computation(), Schedulers.newThread(),
                Schedulers.trampoline(), Schedulers.immediate() };
        Scheduler scheduler = schedulers[jComboBoxSchedulers.getSelectedIndex()];
        // конфігурація спостережуваних величин
        for (int i = 0; i < nbRequests; i++) {
            // підготовка відповіді
            UiResponse uiResponse = new UiResponse();
            uiResponse.setIdClient(i);
            // спостережувана величина налаштована для виконання за розкладом, обраним користувачем
            // потім об’єднання отриманого спостережуваного з загальним спостережуваним
            observable = observable.mergeWith(
                    rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(scheduler));
        }
        // спостерігач
        observable = observable.observeOn(SwingScheduler.getInstance());
        // наразі ми лише виконали налаштування
        // до синхронної служби генерації випадкових чисел ще не надходило жодного запиту
        // підписуємося на спостережуваний об’єкт — саме це спричинить виклик синхронної служби генерації випадкових чисел
        try {
            // тут є лише одна підписка — результатом є підписка
            subscriptions.add(observable.subscribe(
                    // повідомлення про відправлення
                    uiResponse -> {
                        // оновлюємо інтерфейс користувача відповіддю
                        // це можливо, оскільки спостереження відбувається у потоці інтерфейсу користувача
                        updateUi(uiResponse);
                    } ,
                    // повідомлення про помилку
                    th -> {
                        // випадок помилки — її відображають
                        String message = getInfoForThrowable("L'erreur suivante s'est produite", th);
                        JOptionPane.showMessageDialog(this, message, "Informations", JOptionPane.PLAIN_MESSAGE);
                        // скасування запитів
                        doCancel();
                    } ,
                    // повідомлення [onCompleted]
                    // завершення очікування
                    this::endWaiting));
        } catch (Throwable th) {
            // виняток + загальний — відображається
            String message = getInfoForThrowable("L'erreur suivante s'est produite", th);
            JOptionPane.showMessageDialog(this, message, "Informations", JOptionPane.PLAIN_MESSAGE);
            // запити скасовано
            doCancel();
        }
    }
  • рядок 3: графічний інтерфейс модифікується, щоб показати, що виконується операція, яка може тривати довго;
  • рядок 5: створюється порожній об’єкт спостереження. Саме цей об’єкт спостереження буде відстежуватися рівнем [swing];
  • рядок 7: масив можливих планувальників;
  • рядок 9: ми надали користувачеві можливість вибрати планувальник, на якому будуть виконуватися запити. Ми отримуємо обраний ним планувальник;
  • рядки 11–19: кожен із запитів повертає об’єкт, елементи якого накопичуються (mergeWith) (рядок 17) в об’єкті з рядка 5;
  • рядки 13–14: створюється об’єкт [UiResponse]. Нагадаємо, що цей об’єкт є одночасно вхідним параметром методу [RxService.getAleas] і його результатом (рядки 17–18);
  • рядок 14: кожен запит ідентифікується за номером, який тут позначено як [idClient]. Це необхідно, оскільки в асинхронному середовищі порядок отримання відповідей може відрізнятися від порядку надсилання запитів. [idClient] дозволяє визначити, до якого запиту належить відповідь;
  • рядки 17–18: асинхронний запит має вигляд [rxService.getAleas]. Він виконується на планувальнику, обраному користувачем. Його результат типу Observable<UiResponse> об’єднується з об’єктом Observable з рядка 5. Слід чітко усвідомлювати, що метод [rxService.getAleas] тут виконується і повертає об’єкт Observable. Однак це не означає, що випадкові числа вже отримано. Адже об’єкт Observable виконується лише тоді, коли на нього підписано. Наразі цього ще не сталося;
  • рядок 21: це важлива інструкція: ми вимагаємо, щоб спостереження за елементами, що генеруються об’єктом спостереження з рядка 5, здійснювалося у потоці інтерфейсу користувача. Тут використовується власний планувальник бібліотеки RxSwing;
  • рядки 25–51: підписуємося на об’єкт спостереження з рядка 5. Лише тепер випадкові числа будуть запитуватися у синхронної служби генерації цих чисел. Основне міститься в інструкціях рядків 29–33. Решта коду в основному обробляє випадки помилок та сповіщення [onCompleted] про об’єкт спостереження;
  • рядки 28–44: слід пам’ятати, що ми подали запит на спостереження за процесом із рядка 5 у потоці інтерфейсу користувача. Отже, код у рядках 28–44 виконується у цьому потоці;
  • рядки 29–33: обробляється сповіщення [onNext] від об’єкта спостереження. Отримується тип [UiResponse], виведений спостережуваним процесом. Це результат одного з асинхронних запитів. Графічний інтерфейс оновлюється за допомогою цієї відповіді;
  • рядки 34–41: обробляється сповіщення [onError] від об’єкта спостереження. Відображається діалогове вікно з повідомленням про помилку (рядки 37–38), після чого запити скасовуються (рядок 40);
  • рядки 42–44: обробляється сповіщення [onCompleted] від об’єкта спостереження. Оновлюється графічний інтерфейс, щоб показати, що запитувана послуга завершена. Рядок 44 також можна було б написати таким чином
 ()->{endWaiting();}

Тут ми вирішили використовувати посилання на метод;

  • рядки 45–51: деякі винятки не проходять через рядки 34–41. Це відбувається, коли надсилається занадто багато запитів. Після перевищення певного ліміту, який залежить від робочого середовища на момент виконання, з’являється [StackOverflowError], який перехоплюється рядками 45–51;
  • рядок 27: підписка генерує тип [Subscription], який додається до списку підписок. Цей список тут матиме лише один елемент;

У рядку 32 оновлюється графічний інтерфейс за допомогою наступного методу [updateUi]:


    private void updateUi(UiResponse uiResponse) {
        // час відповіді
        uiResponse.setResponseAt();
        // потоки спостереження
        uiResponse.setObservedOn();
        // кількість відповідей
        jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
        // тривалість виконання
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
        // додавання рядка jSON з відповіді до шаблону JList відповідей
        try {
            model.add(0, jsonMapper.writeValueAsString(uiResponse));
        } catch (JsonProcessingException e) {
            e.printStackTrace();
        }
}

Тут видно, що компоненти графічного інтерфейсу оновлюються (рядки 7, 9, 12). Щоб це було можливо, необхідно обов’язково перебувати у потоці Ui (event loop).

Метод [endWaiting] виглядає так:


    private void endWaiting() {
        // кнопка [Générer] видима
        jButtonGenerate.setVisible(true);
        // кнопка [Annuler] прихована
        jButtonCancel.setVisible(false);
        // курсор очікування прихований
        jTabbedPane1.setCursor(Cursor.getDefaultCursor());
        // вкладка відповідей вибрана
        jTabbedPane1.setSelectedIndex(1);
        // час останнього оновлення
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
}

Метод [doCancel] викликається, коли виникає помилка під час виконання асинхронних запитів або коли користувач натискає кнопку [Annuler]. Його код такий:


// підписки на спостережувані величини
    private List<Subscription> subscriptions = new ArrayList<Subscription>();
....

    @Override
    protected void doCancel() {
        // завершення очікування
        endWaiting();
        // у разі підписок
        if (jCheckBoxRxSwing.isSelected() && subscriptions != null) {
            subscriptions.forEach(Subscription::unsubscribe);
            //subscriptions.forEach(s -> s.unsubscribe());
        }
    }

  • рядок 2: [subscriptions] — це список підписок;
  • рядок 11: усі підписки скасовано;
  • рядок 12: інший варіант запису з рядка 11. Метод [forEach] тут очікує екземпляр типу Consumer<Subscription> (див. параграф 4.4);

Повернемося до коду методу [doGenerateWithService]: його можна розкласти на два етапи:

  1. етап налаштування спостережуваних величин. Це відбувається у потоці виклику методу [doGenerateWithService], тобто у потоці інтерфейсу користувача;
  2. підписка, яка спричинить виконання спостережуваних об’єктів;

Якщо для об’єктів спостереження призначено один із планувальників [Schedulers.computation(), Scheduler.io(), Schedulers.newThread()], то вони будуть виконуватися поза потоком інтерфейсу користувача. Ці різні потоки будуть конкурувати за ядро або ядра комп’ютера. Оскільки запити є тривалими операціями (кілька сотень мілісекунд), метод [doGenerateWithService], що виконується у потоці інтерфейсу користувача, завершиться до того, як запити повернуть свої відповіді. Однак цей метод було запущено у відповідь на подію натискання кнопки [Générer]. Після обробки цієї події потік інтерфейсу користувача зможе перейти до обробки наступних подій. Їх є кілька. Так, метод [beginWaiting] встановив кілька таких подій:


    private void beginWaiting() {
        // кнопки
        jButtonGenerate.setVisible(false);
        jButtonCancel.setVisible(true);
        // курсор очікування
        jTabbedPane1.setCursor(Cursor.getPredefinedCursor(Cursor.WAIT_CURSOR));
        jButtonCancel.setCursor(Cursor.getDefaultCursor());
        // очистити відповіді
        model.clear();
        // підписки Rx
        subscriptions.clear();
        // відображення перегляду відповідей
        jTabbedPane1.setSelectedIndex(1);
        jLabelNbReponses.setText("0");
        jLabelDuree.setText("");
        // початок виконання
        debut = new Date().getTime();
}

Практично кожен рядок цього коду впливає на графічний інтерфейс. Це оновлення відбувається не відразу: події ставляться в чергу циклу подій. Після обробки події кліка на кнопці [Générer] ці події виконуються по черзі, і користувач може побачити, як змінюється графічний інтерфейс:

  • відображається вкладка [Response] (рядок 13), і до неї додається курсор очікування (рядок 6)
  • відображається її кнопка [Annuler] (рядок 4), і користувач зможе натиснути на неї;
  • поле JList для відповідей очищується (рядок 9);
  • JLabel кількості відповідей відображає 0;
  • JLabel тривалості виконання відображає порожній рядок;

Протягом усього часу виконання запитів потік UI регулярно отримує доступ до процесора. Тоді він може обробляти події, що знаходяться в черзі. Серед них є ті, що встановлюються методом [updateUi]:


    private void updateUi(UiResponse uiResponse) {
        // час відповіді
        uiResponse.setResponseAt();
        // потік спостереження
        uiResponse.setObservedOn();
        // кількість відповідей
        jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
        // тривалість виконання
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
        // додавання рядка jSON з відповіді до шаблону JList відповідей
        try {
            model.add(0, jsonMapper.writeValueAsString(uiResponse));
        } catch (JsonProcessingException e) {
            e.printStackTrace();
        }
}

Коли потік Ui має контроль:

  • значення JLabel, що відображає кількість відповідей, оновлюється (рядок 7);
  • оновлюється JLabel, що відображає тривалість виконання (рядок 9);
  • JList відповідей оновлюється за допомогою його шаблону (рядок 12);

Таким чином, користувач бачить хід виконання запитів. Крім того, він може скасувати їх за допомогою кнопки [Annuler]. У цьому і полягає вся перевага використання асинхронних сервісів перед шаром [swing], а RxJava є оптимальною технологією для їх реалізації.

Наостанок зазначимо, що якщо користувач обирає один із планувальників [Schedulers.immediate(), Schedulers.trampoline()], то спостережувані об’єкти виконуються в тому самому потоці, що й виклик, тобто в потоці інтерфейсу користувача. У цьому випадку ми повертаємося до синхронного режиму роботи.

Результати, отримані з використанням різних планувальників, наведено в розділах 2.8.1, 2.8.2, 2.8.3, 2.8.4.