Skip to content

8. RxJava w środowisku Swing

8.1. Introduction

W tym miejscu powrócimy do aplikacji Swing przedstawionej w akapicie 2.

  

Aby pracować z RxJava w środowisku Swing, wykorzystamy bibliotekę RxSwing, która rozszerza RxJava o klasy i interfejsy przydatne w środowisku Swing. W tym celu plik Gradle dla przykładu Swing wygląda następująco:

  

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'
}
  • wiersz 15: zależność od RxSwing;

Będziemy korzystać wyłącznie z jednego obiektu właściwego dla RxSwing: harmonogramu [SwingScheduler.getInstance()], który uruchamia i obserwuje obiekty obserwowalne w wątku pętli zdarzeń Swing. Będziemy go używać wyłącznie do obserwowania obiektów obserwowalnych działających w wątkach innych niż wątek pętli zdarzeń. Przypomnijmy architekturę przykładowej aplikacji:

Image

  • warstwa usług asynchronicznych zawiera metody zwracające obiekty obserwowalne. Obiekty te uruchamiamy w wątkach innych niż wątek pętli zdarzeń. Dzięki temu interfejs graficzny nie jest zablokowany i może reagować na działania użytkownika. Najbardziej oczywistym przykładem jest umożliwienie użytkownikowi kliknięcia przycisku [Annuler] w celu przerwania zbyt długiej operacji asynchronicznej. Aby było to możliwe, wystarczy, że interfejs graficzny jest zablokowany (frozen);
  • warstwa Swing chce wykorzystać wyniki zwracane przez operacje asynchroniczne i na ich podstawie zaktualizować interfejs graficzny. Można to jednak zrobić wyłącznie w wątku pętli zdarzeń. W tym celu wyniki te są monitorowane w harmonogramie [SwingScheduler.getInstance()];

W ten sposób w kodzie obsługi zdarzeń interfejsu graficznego interakcja z warstwą asynchroniczną [rxService] przebiega w następujący sposób:


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

gdzie harmonogram [Schedulers.computation()] może zostać zastąpiony innym harmonogramem w zależności od konkretnego zastosowania.

Zachęcamy czytelnika do ponownego zapoznania się z akapitem 2. Dysponuje on teraz wiedzą niezbędną do pełnego zrozumienia tego zagadnienia.

8.2. Struktura kodu

Kod realizuje następującą architekturę:

Image

Projekt IntelliJ IDEA realizujący tę architekturę wygląda następująco:

  
  • pakiet [rxswing.service] implementuje warstwy usług synchronicznych (IService, Service) oraz asynchronicznych (IRxService, RxService);
  • pakiet [rxswing.ui] implementuje interfejs Swing;

8.3. Uruchamianie projektu

Aby uruchomić projekt w środowisku IntelliJ IDEA, należy wykonać następujące czynności:

 

8.4. Usługa synchroniczna

Image

  

Warstwa usługi synchronicznej posiada następujący interfejs [IService]:


package dvp.rxswing.service;

public interface IService {
  // liczby losowe w przedziale [a,b]
  // generowanych jest n liczb, gdzie n samo w sobie jest liczbą losową z przedziału [minCount, maxCount]
  // liczby są generowane po upływie opóźnienia wynoszącego delay milisekund,
  // gdzie [delay] jest liczbą losową z przedziału [minDelay, maxDelay]
  public ServiceResponse getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay);
}

Typ odpowiedzi usługi [ServiceResponse] jest następujący:


package dvp.rxswing.service;

import java.util.List;

public class ServiceResponse {

  // czas oczekiwania na obsługę
  private int delay;
  // liczby losowe
  private List<Integer> aleas;
  // wątek wykonawczy
  private String executedOn;

  // konstruktorzy

  public ServiceResponse() {
      // wątek wykonawczy
    executedOn = Thread.currentThread().getName();
  }

  public ServiceResponse(int delay, List<Integer> aleas) {
      // konstruktor lokalny
    this();
    // inne inicjalizacje
    this.delay = delay;
    this.aleas = aleas;
  }

  // metody pobierające i ustawiające
...
}

Interfejs [IService] jest zaimplementowany przez następującą klasę [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) {
    // liczby losowe w przedziale [a,b]
    // generowanych jest n liczb, gdzie n samo w sobie jest liczbą losową z przedziału [minCount, maxCount]
    // liczby są generowane po upływie opóźnienia wynoszącego delay milisekund,
    // gdzie [delay] jest liczbą losową z przedziału [minDelay, maxDelay]

    // kilka sprawdzeń
    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;
    }
    // błędy?
    if (!messages.isEmpty()) {
      throw new AleasException(String.join(" [---] ", messages), erreur);
    }
    // generator liczb losowych
    Random random = new Random();
    // oczekiwanie?
    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);
      }
    }
    // generowanie wyniku
    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));
    }
    // zwrot wyniku
    return new ServiceResponse(delay,nombres);
  }

}

Klasa wyjątku [AleasException] używana przez usługę jest następująca:


package dvp.rxswing.service;

public class AleasException extends RuntimeException {

    private static final long serialVersionUID = 1L;
    // kod błędu
  private int code;

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

  // metody pobierające i ustawiające
...
}
  • wiersz 3: rozszerza klasę [RuntimeException]. Jest to zatem wyjątek niekontrolowany;
  • wiersz 7: wzbogaca swoją klasę nadrzędną o kod błędu (0 = brak błędu);

8.5. Usługa asynchroniczna

Image

  

Warstwa usług asynchronicznych posiada następujący interfejs [IRxService]:


package dvp.rxswing.service;

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

public interface IRxService {
  // liczby losowe w przedziale [a,b]
  // generowanych jest n liczb, gdzie n samo w sobie jest liczbą losową z przedziału [minCount, maxCount]
  // liczby są generowane po upływie opóźnienia wynoszącego delay milisekund,
  // gdzie [delay] jest liczbą losową z przedziału [minDelay, maxDelay]
  public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse);
}
  • wiersz 11: metoda [getAleas] usługi zwraca teraz obserwowalną wartość;

Metoda [getAleas] zwraca odpowiedź typu [UiResponse] przeznaczoną dla warstwy [Ui]. Typ ten jest następujący:


package dvp.rxswing.ui;

import dvp.rxswing.service.ServiceResponse;

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

public class UiResponse {

  // identyfikator klienta
  private int idClient;
  // odpowiedź serwisu
  private ServiceResponse serviceResponse;
  // nazwa wątku obserwacyjnego
  private String observedOn;
  // czas wysłania zapytania
  private String requestAt;
  // czas odpowiedzi
  private String responseAt;

  // konstruktorzy

  public UiResponse() {
      // wątek obserwacyjny
    observedOn = Thread.currentThread().getName();
    // czas żądania
    requestAt = getTimeStamp();
  }

  // metody prywatne

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

  // metody pobierające i ustawiające
...
}
  • liczby losowe znajdują się w polu w wierszu 13;
  • pozostałe pola służą do określenia wątków wykonawczych i obserwacyjnych dla obserwowalnej wartości usługi asynchronicznej, a także godzin wysłania żądania do usługi i otrzymania odpowiedzi;

Interfejs asynchroniczny jest zaimplementowany przez następującą klasę [RxService]:


package dvp.rxswing.service;

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

public class RxService implements IRxService {

  // usługa synchroniczna
  private IService service;

  // konstruktor
  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) {
      // tworzymy obserwowalny obiekt, który emituje wartość zwracaną przez usługę synchroniczną
    return Observable.create(subscriber -> {
      try {
          // wywołanie synchroniczne
        uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
        // przekazujemy wynik do obserwatora
        subscriber.onNext(uiResponse);
      } catch (Exception e) {
          // przekazujemy błąd do obserwatora
        subscriber.onError(e);
      } finally {
          // powiadamia się obserwatora, że emisje zostały zakończone
        subscriber.onCompleted();
      }
    });
  }
}
  • wiersze 12–14: klasa [RxService] usługi asynchronicznej jest tworzona na podstawie instancji interfejsu synchronicznego [IService];
  • wiersze 20–33: utworzenie obserwowalnej, wynik metody [getAleas];
  • wiersz 22: wywoływana jest metoda synchroniczna [service.getAleas]. Jej wynik typu [ServiceResponse] jest umieszczany w obiekcie typu [UiResponse], który ma zostać przekazany do warstwy [swing]. Obiekt ten został pierwotnie przekazany w parametrach wywołania metody (ostatni parametr, wiersz 17);
  • wiersz 24: odpowiedź [UiResponse] jest wysyłana do obserwatora (warstwa [swing]). Obiekt [UiResponse] zawiera nie tylko informacje utworzone przez usługę synchroniczną w wierszu 22. Zawiera on również inne informacje utworzone przez metodę wywołującą metodę [getAleas] z wiersza 17. Z tego powodu ta metoda wywołująca przekazała obiekt [UiResponse] jako parametr do metody [getAleas] (ostatni parametr, wiersz 17);
  • wiersz 30: nie zapominamy o zasygnalizowaniu zakończenia transmisji. Mamy tu obserwowalną, która emituje tylko jedną wartość: tę zwracaną przez usługę synchroniczną;
  • wiersz 27: informujemy obserwatora o ewentualnym błędzie;

8.6. Interfejs graficzny

Image

  
  • interfejs graficzny został stworzony za pomocą narzędzia IDE [Netbeans], które posiada dobry edytor graficzny. Edytor ten wygenerował plik [AbstractJFrameAleas.form], który może być wykorzystywany wyłącznie przez to narzędzie IDE;
  • klasa [AbstractJFrameAleas] została również wygenerowana przez edytor graficzny NetBeans. Następnie została ona zrefaktoryzowana w następujący sposób: zdarzenia interfejsu graficznego, które chcieliśmy obsłużyć, są przetwarzane w klasie [AbstractJFrameAleas] za pomocą metod abstrakcyjnych zaimplementowanych w klasie potomnej [JFrameAleasEvents]. Ostatecznie,
    • klasa abstrakcyjna [AbstractJFrameAleas] odpowiada za tworzenie i wyświetlanie interfejsu graficznego;
    • klasa potomna [JFrameAleasEvents] odpowiada za obsługę zdarzeń tego interfejsu;

Komponenty graficznego interfejsu użytkownika zakładki [Request] są następujące:

 
nr
typ
nazwa
rola
1
JTabbedPane
jTabbedPane1
kontener zakładek. Zawiera dwie zakładki (JPanel) [jPanelRequest] dla zapytania, [jPanelresponse] dla odpowiedzi;
2
JTextField
jTextFieldNbValeurs
liczba żądań, które należy wysłać do serwisu generującego liczby losowe. W przypadku serwisu asynchronicznego uruchomionego na harmonogramie [Schedulers.io] żądania te będą współdzielić jeden procesor;
3
JTextField
jTextFieldA
punkt a przedziału [a,b]
4
JTextField
jTextFieldB
zacisk b przedziału [a,b]
5
JTextField
jTextFieldMinCount
zacisk minCount z przedziału [minCount, maxCount]
6
JTextField
jTextFieldMaxCount
zacisk maxCount z przedziału [minCount, maxCount]
7
JTextField
jTextFieldMinDelay
zacisk minDelay z przedziału [minDelay, maxDelay]
8
JTextField
jTextFieldMaxDelay
zacisk maxDelay z przedziału [minDelay, maxDelay]
9
JCheckBox
jCheckBoxRxSwing
jeśli pole jest zaznaczone, zapytania są kierowane do interfejsu asynchronicznego. W przeciwnym razie są kierowane do interfejsu synchronicznego
10
JComboBox
jComboBoxSchedulers
w przypadku żądań asynchronicznych będą one wykonywane zgodnie z harmonogramem wybranym w tym miejscu
11
JButton
jButtonGenerate
uruchamia wykonywanie zapytań w trybie synchronicznym lub asynchronicznym

Elementy interfejsu graficznego zakładki [Response] są następujące:

 
nr
typ
nazwa
rola
1
JLabel
jLabelDuree
całkowity czas wykonania zapytań w milisekundach
2
JLabel
jLabelNbReponses
całkowita liczba zaobserwowanych odpowiedzi (może różnić się od liczby zapytań, ponieważ każde zapytanie może dostarczyć kilka wartości do obserwacji)
3
JList
jListNumbers
wyświetlanie wartości zaobserwowanych (otrzymanych)
4
JButton
jButtonAnnuler
anuluje żądania w trakcie wykonywania

8.7. Uruchomienie interfejsu graficznego

  

Klasa [JFrameAleasEvents] obsługuje zdarzenia interfejsu graficznego, w szczególności kliknięcie przycisku [Générer]. Jest to klasa wykonywalna, która uruchamia się w następującym kontekście:


public class JFrameAleasEvents extends AbstractJFrameAleas {

    private static final long serialVersionUID = 1L;
    // usługa generowania synchronicznego
    private IService service;
    // usługa generacji asynchronicznej
    private IRxService rxService;

    // wpisy
    private int nbRequests;
    private int a;
    private int b;
    private int minDelay;
    private int maxDelay;
    private int minCount;
    private int maxCount;

    // komunikaty o błędach
    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 ";

    // subskrypcje obserwowalnych
    protected List<Subscription> subscriptions = new ArrayList<Subscription>();
    // początek i koniec wykonania
    private long debut;
    // mapper jSON
    private ObjectMapper jsonMapper;
    // model odpowiedzi
    private DefaultListModel<String> model;

    // konstruktor
    public JFrameAleasEvents() {
        // element nadrzędny
        super();
        // lokalny
        initJFrame();
        // usługi
        service = new Service();
        rxService = new RxService(service);
        // mapujący jSON
        jsonMapper = new ObjectMapper();
    }

    private void initJFrame() {
        // ukrywanie komunikatów o błędach
        jLabelCountError.setText("");
        jLabelDelayError.setText("");
        jLabelIntervalError.setText("");
        jLabelNbValuesError.setText("");
        // ukrywanie domyślnych tekstów
        jTextFieldA.setText("100");
        jTextFieldB.setText("200");
        jTextFieldMinCount.setText("5");
        jTextFieldMaxCount.setText("10");
        jTextFieldMinDelay.setText("100");
        jTextFieldMaxDelay.setText("500");
        jTextFieldNbValeurs.setText("10");
        jLabelDuree.setText("");
        // szablon odpowiedzi
        model = new DefaultListModel<>();
        jListNumbers.setModel(model);
        // liczba rdzeni
        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);
        }

        /* Utwórz i wyświetl formularz */
        java.awt.EventQueue.invokeLater(() -> {
            new JFrameAleasEvents().setVisible(true);
        });
    }
  • wiersz 1: klasa [JFrameAleasEvents] dziedziczy po klasie [AbstractJFrameAleas], która z kolei dziedziczy po klasie Swing [JFrame]. Klasa [JFrameAleasEvents] jest zatem oknem Swing;
  • wiersze 68–75: metoda [main], która zostanie wykonana;
  • wiersz 70: ustawia wygląd i styl interfejsu graficznego;
  • wiersz 79: wywoływany jest konstruktor klasy [JFrameAleasEvents]: interfejs graficzny zostanie utworzony i zainicjowany. Po zakończeniu tej operacji interfejs staje się widoczny;
  • wiersze 34–44: konstruktor;
  • wiersz 36: wywołanie konstruktora klasy nadrzędnej spowoduje zainicjowanie interfejsu graficznego. W tym momencie wygląda on tak, jak zaprojektował go programista. Nie jest jeszcze widoczny;
  • wiersz 38: niektóre komponenty interfejsu graficznego są inicjowane;
  • wiersz 40: instancjonowanie usługi synchronicznej;
  • wiersz 41: utworzenie instancji usługi asynchronicznej;

8.8. Wykonanie zapytań synchronicznych

Kliknięcie przycisku [Générer] powoduje wykonanie następującej metody [doGenerate]:


    @Override
    protected void doGenerate() {
        // czy dane są prawidłowe?
        if (!isPageValid()) {
            return;
        }
        // rx czy nie?
        if (jCheckBoxRxSwing.isSelected()) {
            // żądania asynchroniczne
            doGenerateWithRxService();
        } else {
            // żądania synchroniczne
            doGenerateWithService();
        }
}
  • wiersze 4–6: sprawdzane jest, czy dane wprowadzone przez użytkownika są prawidłowe. Nie będziemy omawiać metody [isPageValid]. Jest ona prosta;
  • wiersz 8: sprawdzany jest stan pola wyboru RxSwing;
  • wiersz 13: żądania są wykonywane synchronicznie;

Metoda [doGenerateWithService] wygląda następująco:


    // generowanie synchroniczne
    private void doGenerateWithService() {
        // początek oczekiwania
        beginWaiting();
        try {
            for (int i = 0; i < nbRequests; i++) {
                // przygotowanie odpowiedzi
                UiResponse uiResponse = new UiResponse();
                // nr klienta
                uiResponse.setIdClient(i);
                // wywołanie synchroniczne
                uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
                // czas odpowiedzi
                uiResponse.setResponseAt();
                // aktualizacja szablonu JList o otrzymane odpowiedzi
                model.add(0, jsonMapper.writeValueAsString(uiResponse));
                // aktualizacja liczby odpowiedzi
                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);
        }
        // koniec oczekiwania
        endWaiting();
}
  • wiersz 12: synchroniczne wywołanie usługi generującej liczby losowe;
  • wykonanie metody [doGenerateWithService] odbywa się w całości w wątku pętli zdarzeń Swing. Dopóki metoda nie zostanie zakończona, interfejs graficzny nie przetwarza żadnych nowych zdarzeń. Jest on zawieszony (frozen). Na przykład aktualizacje interfejsu graficznego z wierszy 16 i 18 nigdy nie zostaną wyświetlone. Będą one widoczne dopiero z ich ostatecznymi wartościami, i to dopiero po zakończeniu wykonywania wszystkich żądań;

Metoda [beginWaiting] (wiersz 4) wygląda następująco:


    private void beginWaiting() {
        // przyciski
        jButtonGenerate.setVisible(false);
        jButtonCancel.setVisible(true);
        // wskaźnik oczekiwania
        jTabbedPane1.setCursor(Cursor.getPredefinedCursor(Cursor.WAIT_CURSOR));
        jButtonCancel.setCursor(Cursor.getDefaultCursor());
        // 0 odpowiedzi
        model.clear();
        // subskrypcje Rx
        subscriptions.clear();
        // wyświetla się widok odpowiedzi
        jTabbedPane1.setSelectedIndex(1);
        jLabelNbReponses.setText("0");
        jLabelDuree.setText("");
        // rozpoczęcie wykonywania
        debut = new Date().getTime();
}
  • wiersz 3: przycisk [Générer] jest ukryty. Powoduje to wygenerowanie zdarzenia, które również może zostać wykonane dopiero po zakończeniu wykonywania wszystkich zapytań. Dlatego nigdy nie widać go jako ukrytego, ponieważ metoda [endWaiting] w wierszu 25 metody [doGenerateWithService] ponownie go wyświetla;
  • wiersz 13: wybieramy kartę [Response], aby zobaczyć napływające odpowiedzi. Ponownie, to zdarzenie zostanie wykonane dopiero po zakończeniu realizacji wszystkich zapytań, kiedy to zobaczymy wszystkie odpowiedzi naraz, podczas gdy chcieliśmy je widzieć po kolei;

Interfejs synchroniczny ma wyraźne wady. Można je przezwyciężyć dzięki interfejsowi asynchronicznemu.

8.9. Wykonanie zapytań asynchronicznych

Kod służący do wykonywania zapytań asynchronicznych wygląda następująco:


private void doGenerateWithRxService() {
        // początek oczekiwania
        beginWaiting();
        // uzyskamy liczby losowe w postaci obserwowalnej
        Observable<UiResponse> observable = Observable.empty();
        // Harmonogram wykonywania poszczególnych obserwowalnych
        Scheduler[] schedulers = { Schedulers.io(), Schedulers.computation(), Schedulers.newThread(),
                Schedulers.trampoline(), Schedulers.immediate() };
        Scheduler scheduler = schedulers[jComboBoxSchedulers.getSelectedIndex()];
        // konfiguracja obserwowalnych
        for (int i = 0; i < nbRequests; i++) {
            // przygotowanie odpowiedzi
            UiResponse uiResponse = new UiResponse();
            uiResponse.setIdClient(i);
            // obserwowalna jest skonfigurowana do działania zgodnie z harmonogramem wybranym przez użytkownika
            // następnie zsumowanie uzyskanej obserwowalnej z obserwowalną całości
            observable = observable.mergeWith(
                    rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(scheduler));
        }
        // obserwator
        observable = observable.observeOn(SwingScheduler.getInstance());
        // na razie przeprowadzono jedynie konfigurację
        // nie wysłano jeszcze żadnego żądania do synchronicznej usługi generowania liczb losowych
        // subskrybujemy obserwowalną – to właśnie spowoduje wywołanie synchronicznej usługi generowania liczb losowych
        try {
            // mamy tu tylko subskrypcję – wynikiem jest subskrypcja
            subscriptions.add(observable.subscribe(
                    // powiadomienie o emisji
                    uiResponse -> {
                        // aktualizujemy interfejs użytkownika o odpowiedź
                        // jest to możliwe, ponieważ obserwacja odbywa się w wątku interfejsu użytkownika
                        updateUi(uiResponse);
                    } ,
                    // powiadomienie o błędzie
                    th -> {
                        // w przypadku błędu – wyświetla się komunikat
                        String message = getInfoForThrowable("L'erreur suivante s'est produite", th);
                        JOptionPane.showMessageDialog(this, message, "Informations", JOptionPane.PLAIN_MESSAGE);
                        // anulowanie żądań
                        doCancel();
                    } ,
                    // powiadomienie [onCompleted]
                    // koniec oczekiwania
                    this::endWaiting));
        } catch (Throwable th) {
            // wyjątek + ogólny – wyświetl
            String message = getInfoForThrowable("L'erreur suivante s'est produite", th);
            JOptionPane.showMessageDialog(this, message, "Informations", JOptionPane.PLAIN_MESSAGE);
            // anulowanie żądań
            doCancel();
        }
    }
  • wiersz 3: interfejs graficzny zostaje zmodyfikowany, aby pokazać, że trwa operacja, która może potrwać długo;
  • wiersz 5: tworzony jest pusty obiekt obserwowalny. To właśnie ten obiekt będzie obserwowany przez warstwę [swing];
  • wiersz 7: tablica możliwych harmonogramów;
  • wiersz 9: umożliwiliśmy użytkownikowi wybór harmonogramu, na którym mają być wykonywane zapytania. Pobieramy wybrany przez niego harmonogram;
  • wiersze 11–19: każde z zapytań zwraca obserwowalną, której elementy są kumulowane (mergeWith) (wiersz 17) w obserwowalnej z wiersza 5;
  • wiersze 13–14: tworzony jest obiekt [UiResponse]. Przypominamy, że obiekt ten jest zarówno parametrem wejściowym metody [RxService.getAleas], jak i jej wynikiem (wiersze 17–18);
  • wiersz 14: każde żądanie jest identyfikowane za pomocą numeru, zwanego tutaj [idClient]. Jest to konieczne, ponieważ w środowisku asynchronicznym kolejność otrzymywania odpowiedzi może różnić się od kolejności wysyłania żądań. [idClient] pozwala ustalić, do którego żądania należy dana odpowiedź;
  • wiersze 17–18: wysyłane jest asynchroniczne żądanie o nazwie [rxService.getAleas]. Jest ono wykonywane na harmonogramie wybranym przez użytkownika. Jego wynik typu Observable<UiResponse> jest sumowany z obserwowalną z wiersza 5. Należy pamiętać, że metoda [rxService.getAleas] jest tutaj wykonywana i zwraca obserwowalną. Nie oznacza to jednak, że uzyskano liczby losowe. Obserwowalna jest bowiem wykonywana dopiero po zasubskrybowaniu jej. Na razie tak się nie stało;
  • wiersz 21: to jest kluczowa instrukcja: żądamy, aby obserwacja elementów wysyłanych przez obserwowalną z wiersza 5 odbywała się w wątku interfejsu użytkownika. Wykorzystujemy tutaj harmonogram właściwy dla biblioteki RxSwing;
  • wiersze 25–51: subskrybujemy obserwowalną z wiersza 5. Dopiero teraz liczby losowe zostaną zażądane od synchronicznej usługi generującej te liczby. Najważniejsze elementy znajdują się w instrukcjach z wierszy 29–33. Pozostała część kodu zajmuje się głównie obsługą błędów oraz powiadomieniem o obserwowalnej [onCompleted];
  • wiersze 28–44: należy pamiętać, że poproszono o obserwację procesu z wiersza 5 w wątku interfejsu użytkownika. Zatem kod z wierszy 28–44 jest wykonywany w wątku interfejsu użytkownika;
  • wiersze 29–33: przetwarzamy powiadomienie [onNext] z obserwowalnego obiektu. Otrzymujemy typ [UiResponse] wysłany przez obserwowany proces. Jest to wynik jednego z asynchronicznych żądań. Aktualizujemy interfejs graficzny tą odpowiedzią;
  • wiersze 34–41: przetwarzamy powiadomienie [onError] z obserwowalnego obiektu. Wyświetlamy okno dialogowe z komunikatem o błędzie (wiersze 37–38), a następnie anulujemy żądania (wiersz 40);
  • wiersze 42–44: przetwarzamy powiadomienie [onCompleted] z obserwowalnego. Aktualizujemy interfejs graficzny, aby pokazać, że żądana usługa została zakończona. Wiersz 44 można by również zapisać w następujący sposób
 ()->{endWaiting();}

W tym przypadku zdecydowano się jednak na użycie odwołania do metody;

  • wiersze 45–51: niektóre wyjątki nie przechodzą przez wiersze 34–41. Dzieje się tak w przypadku wysłania zbyt dużej liczby żądań. Po przekroczeniu pewnego limitu, który zależy od środowiska pracy w momencie wykonywania, pojawia się komunikat [StackOverflowError], który jest przechwytywany przez wiersze 45–51;
  • wiersz 27: subskrypcja generuje typ [Subscription], który jest dodawany do listy subskrypcji. Lista ta będzie zawierała w tym przypadku tylko jeden element;

W wierszu 32 aktualizujemy interfejs graficzny za pomocą następującej metody [updateUi]:


    private void updateUi(UiResponse uiResponse) {
        // czas odpowiedzi
        uiResponse.setResponseAt();
        // wątek obserwacyjny
        uiResponse.setObservedOn();
        // liczba odpowiedzi
        jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
        // czas wykonania
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
        // dodanie ciągu jSON z odpowiedzi do szablonu JList odpowiedzi
        try {
            model.add(0, jsonMapper.writeValueAsString(uiResponse));
        } catch (JsonProcessingException e) {
            e.printStackTrace();
        }
}

Widać tutaj, że komponenty interfejsu graficznego są aktualizowane (wiersze 7, 9, 12). Aby było to możliwe, konieczne jest przebywanie w wątku interfejsu użytkownika (pętla zdarzeń).

Metoda [endWaiting] wygląda następująco:


    private void endWaiting() {
        // przycisk [Générer] widoczny
        jButtonGenerate.setVisible(true);
        // przycisk [Annuler] jest ukryty
        jButtonCancel.setVisible(false);
        // ukryty kursor oczekiwania
        jTabbedPane1.setCursor(Cursor.getDefaultCursor());
        // zakładka odpowiedzi zaznaczona
        jTabbedPane1.setSelectedIndex(1);
        // czas ostatniej aktualizacji
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
}

Metoda [doCancel] jest wywoływana w przypadku wystąpienia błędu podczas wykonywania żądań asynchronicznych lub gdy użytkownik kliknie przycisk [Annuler]. Jej kod wygląda następująco:


// subskrypcje obserwowalnych
    private List<Subscription> subscriptions = new ArrayList<Subscription>();
....

    @Override
    protected void doCancel() {
        // koniec oczekiwania
        endWaiting();
        // w przypadku subskrypcji
        if (jCheckBoxRxSwing.isSelected() && subscriptions != null) {
            subscriptions.forEach(Subscription::unsubscribe);
            //subscriptions.forEach(s -> s.unsubscribe());
        }
    }

  • wiersz 2: [subscriptions] to lista subskrypcji;
  • wiersz 11: wszystkie subskrypcje są anulowane;
  • wiersz 12: inne zapisy z wiersza 11. Metoda [forEach] oczekuje tutaj instancji typu Consumer<Subscription> (patrz punkt 4.4);

Wróćmy do kodu metody [doGenerateWithService]: można go podzielić na dwa etapy:

  1. etap konfiguracji obserwowalnych. Odbywa się to w wątku wywołującego metodę [doGenerateWithService], czyli w wątku interfejsu użytkownika;
  2. subskrypcja, która spowoduje wykonanie obserwowalnych;

Jeśli dla obserwowalnych obiektów jako harmonogram wybrano jeden z harmonogramów [Schedulers.computation(), Scheduler.io(), Schedulers.newThread()], to będą one wykonywane poza wątkiem interfejsu użytkownika. Te różne wątki będą konkurować o rdzeń lub rdzenie komputera. Ponieważ zapytania są operacjami długotrwałymi (trwającymi kilkaset milisekund), metoda [doGenerateWithService] wykonywana w wątku interfejsu użytkownika zakończy się, zanim zapytania zwrócą swoje odpowiedzi. Metoda ta została jednak uruchomiona w odpowiedzi na zdarzenie kliknięcia przycisku [Générer]. Po przetworzeniu tego zdarzenia wątek interfejsu użytkownika będzie mógł przejść do przetwarzania kolejnych zdarzeń. Jest ich kilka. Metoda [beginWaiting] ustawiła bowiem kilka z nich:


    private void beginWaiting() {
        // przyciski
        jButtonGenerate.setVisible(false);
        jButtonCancel.setVisible(true);
        // kursor oczekiwania
        jTabbedPane1.setCursor(Cursor.getPredefinedCursor(Cursor.WAIT_CURSOR));
        jButtonCancel.setCursor(Cursor.getDefaultCursor());
        // wyzeruj odpowiedzi
        model.clear();
        // subskrypcje Rx
        subscriptions.clear();
        // wyświetla widok odpowiedzi
        jTabbedPane1.setSelectedIndex(1);
        jLabelNbReponses.setText("0");
        jLabelDuree.setText("");
        // rozpoczęcie wykonywania
        debut = new Date().getTime();
}

Praktycznie wszystkie linijki tego kodu mają wpływ na interfejs graficzny. Aktualizacja ta nie następuje natychmiast: zdarzenia są umieszczane w kolejce pętli zdarzeń. Po przetworzeniu zdarzenia kliknięcia przycisku [Générer] zdarzenia te są wykonywane po kolei, a użytkownik może zaobserwować zmiany w interfejsie graficznym:

  • wyświetlana jest zakładka [Response] (wiersz 13) i przypisywany jest do niej kursor oczekiwania (wiersz 6)
  • wyświetlany jest przycisk [Annuler] (wiersz 4), a użytkownik będzie mógł go kliknąć;
  • pole JList z odpowiedziami zostaje wyczyszczone (wiersz 9);
  • JLabel liczby odpowiedzi wyświetla 0;
  • JLabel czasu wykonania wyświetla pusty ciąg znaków;

Przez cały czas wykonywania zapytań wątek UI ma regularny dostęp do procesora. Może wówczas przetwarzać oczekujące zdarzenia. Wśród nich znajdują się te ustawione przez metodę [updateUi]:


    private void updateUi(UiResponse uiResponse) {
        // czas odpowiedzi
        uiResponse.setResponseAt();
        // wątek obserwacyjny
        uiResponse.setObservedOn();
        // liczba odpowiedzi
        jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
        // czas trwania wykonania
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
        // dodanie ciągu jSON z odpowiedzi do szablonu JList odpowiedzi
        try {
            model.add(0, jsonMapper.writeValueAsString(uiResponse));
        } catch (JsonProcessingException e) {
            e.printStackTrace();
        }
}

Gdy wątek interfejsu użytkownika ma kontrolę:

  • wartość JLabel określająca liczbę odpowiedzi jest aktualizowana (wiersz 7);
  • aktualizowany jest JLabel dotyczący czasu wykonania (wiersz 9);
  • JList dotyczący odpowiedzi jest aktualizowany na podstawie swojego szablonu (wiersz 12);

W ten sposób użytkownik widzi postęp realizacji zapytań. Ponadto może je anulować za pomocą przycisku [Annuler]. Na tym właśnie polega zaleta posiadania usług asynchronicznych przed warstwą [swing], a RxJava jest technologią z wyboru do ich implementacji.

Na koniec należy zauważyć, że jeśli użytkownik wybierze jeden z harmonogramów [Schedulers.immediate(), Schedulers.trampoline()], obserwowalne są wówczas wykonywane w tym samym wątku co wywołujący, tj. w wątku interfejsu użytkownika. W ten sposób powracamy do działania synchronicznego.

Wyniki uzyskane przy użyciu różnych harmonogramów przedstawiono w punktach 2.8.1, 2.8.2, 2.8.3 i 2.8.4.