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:

- 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ę:

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

![]() |
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

![]() |
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

![]() |
- 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
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:
- etap konfiguracji obserwowalnych. Odbywa się to w wątku wywołującego metodę [doGenerateWithService], czyli w wątku interfejsu użytkownika;
- 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.








