8. RxJava in der Swing-Umgebung
8.1. Introduction
Wir werden hier auf die in Abschnitt 2 vorgestellte Swing-Anwendung zurückkommen.
![]() |
Um mit RxJava in einer Swing-Umgebung zu arbeiten, verwenden wir die Bibliothek RxSwing, die RxJava um Klassen und Schnittstellen erweitert, die in einer Swing-Umgebung nützlich sind. Dazu sieht die Gradle-Datei des Swing-Beispiels wie folgt aus:
![]() |
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'
}
- Zeile 15: die Abhängigkeit von RxSwing;
Wir werden nur ein einziges Objekt verwenden, das spezifisch für RxSwing ist: den Scheduler [SwingScheduler.getInstance()], der die Ausführung bzw. Beobachtung der Observables im Thread der Swing-Event-Loop übernimmt. Wir werden ihn ausschließlich dazu verwenden, Observables zu beobachten, die auf anderen Threads als dem der Event-Loop ausgeführt werden. Zur Erinnerung: Die Architektur der Beispielanwendung sieht wie folgt aus:

- Die asynchrone Service-Schicht verfügt über Methoden, die Observables zurückgeben. Wir führen diese Observables in anderen Threads als dem der Event-Loop aus. So bleibt die grafische Benutzeroberfläche nicht eingefroren. Sie kann auf Benutzeraktionen reagieren. Am offensichtlichsten ist es, dem Benutzer zu ermöglichen, auf eine Schaltfläche ([Annuler]) zu klicken, um einen zu lang andauernden asynchronen Vorgang abzubrechen. Damit dies möglich ist, darf die grafische Benutzeroberfläche lediglich kurzzeitig eingefroren sein (frozen);
- die Swing-Schicht möchte die von den asynchronen Vorgängen zurückgegebenen Ergebnisse auswerten und anhand dieser die grafische Benutzeroberfläche aktualisieren. Dies kann jedoch nur im Thread der Event-Loop erfolgen. Zu diesem Zweck werden diese Ergebnisse im Scheduler [SwingScheduler.getInstance()] überwacht;
Im Code zur Ereignisverarbeitung der grafischen Benutzeroberfläche erfolgt die Interaktion mit der asynchronen Schicht [rxService] somit in folgender Form:
Observable obs=rxService.doSomething(...).subscribeOn(Schedulers.computation()).observeOn(SwingScheduler.getInstance()) ;
wobei der Scheduler [Schedulers.computation()] je nach Anwendungsfall durch einen anderen Scheduler ersetzt werden kann.
Der Leser wird gebeten, Absatz 2 noch einmal durchzulesen. Er verfügt nun über das nötige Wissen, um ihn vollständig zu verstehen.
8.2. Die Struktur des Codes
Der Code implementiert die folgende Architektur:

Das IntelliJ IDEA-Projekt, das diese Architektur implementiert, lautet wie folgt:
![]() |
- Das Paket [rxswing.service] implementiert die synchronen (IService, Service) und asynchronen (IRxService, RxService) Service-Schichten;
- Das Paket [rxswing.ui] implementiert die Swing-Schnittstelle;
8.3. Ausführung des Projekts
Um das Projekt in IntelliJ IDEA auszuführen, gehen Sie wie folgt vor:
![]() |
8.4. Der synchrone Dienst

![]() |
Die synchrone Service-Schicht verfügt über die folgende Schnittstelle: [IService]:
package dvp.rxswing.service;
public interface IService {
// Zufallszahlen im Intervall [a,b]
// n Zahlen werden generiert, wobei n selbst eine Zufallszahl im Intervall [minCount, maxCount] ist
// Die Zahlen werden nach einer Wartezeit von delay Millisekunden generiert,
// wobei [delay] selbst eine Zufallszahl im Intervall [minDelay, maxDelay] ist
public ServiceResponse getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay);
}
Der Typ [ServiceResponse] der Dienstantwort lautet wie folgt:
package dvp.rxswing.service;
import java.util.List;
public class ServiceResponse {
// Wartezeit des Dienstes
private int delay;
// Zufallszahlen
private List<Integer> aleas;
// Ausführungsthread
private String executedOn;
// Konstruktoren
public ServiceResponse() {
// Ausführungsthread
executedOn = Thread.currentThread().getName();
}
public ServiceResponse(int delay, List<Integer> aleas) {
// Lokaler Konstruktor
this();
// sonstige Initialisierungen
this.delay = delay;
this.aleas = aleas;
}
// Getter und Setter
...
}
Die Schnittstelle [IService] wird durch die folgende Klasse [Service] implementiert:
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) {
// Zufallszahlen im Intervall [a,b]
// n Zahlen werden generiert, wobei n selbst eine Zufallszahl im Intervall [minCount, maxCount] ist
// Die Zahlen werden nach einer Wartezeit von delay Millisekunden generiert,
// wobei [delay] selbst eine Zufallszahl im Intervall [minDelay, maxDelay] ist
// einige Überprüfungen
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;
}
// Fehler?
if (!messages.isEmpty()) {
throw new AleasException(String.join(" [---] ", messages), erreur);
}
// Zufallszahlengenerator
Random random = new Random();
// Warten?
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);
}
}
// Ergebnisgenerierung
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));
}
// Ergebnisrückgabe
return new ServiceResponse(delay,nombres);
}
}
Die vom Dienst verwendete Ausnahmeklasse [AleasException] lautet wie folgt:
package dvp.rxswing.service;
public class AleasException extends RuntimeException {
private static final long serialVersionUID = 1L;
// Fehlercode
private int code;
// Konstruktoren
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;
}
// Getter und Setter
...
}
- Zeile 3: Sie erweitert die Klasse [RuntimeException]. Es handelt sich also um eine unkontrollierte Ausnahme;
- Zeile 7: Sie erweitert ihre übergeordnete Klasse um einen Fehlercode (0 = kein Fehler);
8.5. Der asynchrone Dienst

![]() |
Die asynchrone Service-Schicht verfügt über die folgende Schnittstelle [IRxService]:
package dvp.rxswing.service;
import dvp.rxswing.ui.UiResponse;
import rx.Observable;
public interface IRxService {
// Zufallszahlen im Intervall [a,b]
// n Zahlen werden generiert, wobei n selbst eine Zufallszahl im Intervall ist [minCount, maxCount]
// Die Zahlen werden nach einer Wartezeit von delay Millisekunden generiert,
// wobei [delay] selbst eine Zufallszahl im Intervall [minDelay, maxDelay] ist
public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse);
}
- Zeile 11: Die Methode [getAleas] des Dienstes gibt nun ein Observable zurück;
Die Methode [getAleas] gibt eine Antwort vom Typ [UiResponse] zurück, die für die Schicht [Ui] bestimmt ist. Dieser Typ ist wie folgt definiert:
package dvp.rxswing.ui;
import dvp.rxswing.service.ServiceResponse;
import java.text.SimpleDateFormat;
import java.util.Calendar;
public class UiResponse {
// Kunden-ID
private int idClient;
// Antwort des Dienstes
private ServiceResponse serviceResponse;
// Name des Beobachtungs-Threads
private String observedOn;
// Zeitpunkt der Anfrage
private String requestAt;
// Zeitpunkt der Antwort
private String responseAt;
// Konstruktoren
public UiResponse() {
// Beobachtungs-Thread
observedOn = Thread.currentThread().getName();
// Zeitpunkt der Anfrage
requestAt = getTimeStamp();
}
// private Methoden
private String getTimeStamp() {
return new SimpleDateFormat("hh:mm:ss:SSS").format(Calendar.getInstance().getTime());
}
// Getter und Setter
...
}
- Die Zufallszahlen befinden sich im Feld in Zeile 13;
- die übrigen Felder dienen zur Angabe der Ausführungs- und Beobachtungs-Threads der Beobachtungsgröße des asynchronen Dienstes sowie der Zeitpunkte der an den Dienst gerichteten Anfrage und der erhaltenen Antwort;
Die asynchrone Schnittstelle wird durch die folgende Klasse „[RxService]“ implementiert:
package dvp.rxswing.service;
import dvp.rxswing.ui.UiResponse;
import rx.Observable;
public class RxService implements IRxService {
// synchroner Dienst
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) {
// Es wird ein Observable erstellt, das den vom synchronen Dienst zurückgegebenen Wert ausgibt
return Observable.create(subscriber -> {
try {
// synchroner Aufruf
uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
// Das Ergebnis wird an den Beobachter übergeben
subscriber.onNext(uiResponse);
} catch (Exception e) {
// Der Fehler wird an den Beobachter übergeben
subscriber.onError(e);
} finally {
// dem Beobachter wird gemeldet, dass die Emissionen beendet sind
subscriber.onCompleted();
}
});
}
}
- Zeilen 12–14: Die Klasse [RxService] des asynchronen Dienstes wird aus einer Instanz der synchronen Schnittstelle [IService] erstellt;
- Zeilen 20–33: Erstellung des Observables, Ergebnis der Methode [getAleas];
- Zeile 22: Die synchrone Methode [service.getAleas] wird aufgerufen. Ihr Ergebnis vom Typ [ServiceResponse] wird in das Objekt vom Typ [UiResponse] aufgenommen, das an die Schicht [swing] übergeben werden soll. Dieses Objekt wurde ursprünglich in den Aufrufparametern der Methode übergeben (letzter Parameter, Zeile 17);
- Zeile 24: Die Antwort [UiResponse] wird an den Beobachter (die Schicht [swing]) gesendet. Das Objekt [UiResponse] enthält nicht nur die vom synchronen Dienst in Zeile 22 erstellten Informationen. Es enthält auch weitere Informationen, die von der aufrufenden Methode der Methode [getAleas] in Zeile 17 erstellt wurden. Aus diesem Grund hat diese aufrufende Methode das Objekt [UiResponse] als Parameter an die Methode [getAleas] übergeben (letzter Parameter, Zeile 17);
- Zeile 30: Man vergisst nicht, das Ende der Übertragungen zu melden. Hier gibt es eine Beobachtungsgröße, die nur einen Wert ausgibt: den vom synchronen Dienst zurückgegebenen Wert;
- Zeile 27: Dem Beobachter wird ein eventueller Fehler gemeldet;
8.6. Die grafische Benutzeroberfläche

![]() |
- Die grafische Benutzeroberfläche wurde mit dem IDE [Netbeans] erstellt, der über einen guten grafischen Editor verfügt. Dieser Editor hat die Datei [AbstractJFrameAleas.form] generiert, die nur von diesem IDE verarbeitet werden kann;
- Die Klasse [AbstractJFrameAleas] wurde ebenfalls vom grafischen Editor von NetBeans generiert. Sie wurde anschließend wie folgt umgestaltet: Die Ereignisse der grafischen Benutzeroberfläche, die wir verarbeiten wollten, werden in der Klasse [AbstractJFrameAleas] durch abstrakte Methoden behandelt, die in der Unterklasse [JFrameAleasEvents] implementiert sind. Letztendlich,
- Die abstrakte Klasse [AbstractJFrameAleas] ist für den Aufbau und die Anzeige der grafischen Benutzeroberfläche zuständig;
- die untergeordnete Klasse [JFrameAleasEvents] ist für die Verwaltung der Ereignisse dieser Benutzeroberfläche zuständig;
Die Komponenten der grafischen Benutzeroberfläche der Registerkarte [Request] sind folgende:
![]() |
Nr. | Typ | Name | Rolle |
1 | JTabbedPane | jTabbedPane1 | ein Registerkarten-Container. Enthält zwei Registerkarten (JPanel) [jPanelRequest] für die Anfrage, [jPanelresponse] für die Antwort; |
2 | JTextField | jTextFieldNbValeurs | die Anzahl der Anfragen an den Zufallszahlendienst. Im Fall des asynchronen Dienstes, der auf dem Scheduler [Schedulers.io] ausgeführt wird, teilen sich diese Anfragen einen Prozessor; |
3 | JTextField | jTextFieldA | Endpunkt a des Intervalls [a,b] |
4 | JTextField | jTextFieldB | Klemme b des Intervalls [a,b] |
5 | JTextField | jTextFieldMinCount | Klemme minCount des Intervalls [minCount, maxCount] |
6 | JTextField | jTextFieldMaxCount | Klemme maxCount aus dem Bereich [minCount, maxCount] |
7 | JTextField | jTextFieldMinDelay | Klemme minDelay aus dem Bereich [minDelay, maxDelay] |
8 | JTextField | jTextFieldMaxDelay | Klemme maxDelay aus dem Bereich [minDelay, maxDelay] |
9 | JCheckBox | jCheckBoxRxSwing | Wenn das Kontrollkästchen aktiviert ist, werden die Anfragen über die asynchrone Schnittstelle gestellt. Andernfalls erfolgen sie über die synchrone Schnittstelle |
10 | JComboBox | jComboBoxSchedulers | Bei asynchronen Anfragen werden diese mit dem hier ausgewählten Scheduler ausgeführt |
11 | JButton | jButtonGenerate | Löst die Ausführung der Abfragen im synchronen oder asynchronen Modus aus |
Die Komponenten der grafischen Benutzeroberfläche der Registerkarte „[Response]“ sind folgende:
![]() |
Nr. | Typ | Name | Rolle |
1 | JLabel | jLabelDuree | die Gesamtlaufzeit der Abfragen in Millisekunden |
2 | JLabel | jLabelNbReponses | die Gesamtzahl der beobachteten Antworten (kann von der Anzahl der Abfragen abweichen, da jede Abfrage mehrere zu beobachtende Werte liefern kann) |
3 | JList | jListNumbers | Anzeige der beobachteten (empfangenen) Werte |
4 | JButton | jButtonAnnuler | Bricht die aktuell ausgeführten Abfragen ab |
8.7. Instanziierung der grafischen Benutzeroberfläche
![]() |
Die Klasse [JFrameAleasEvents] verwaltet die Ereignisse der grafischen Benutzeroberfläche, insbesondere den Klick auf die Schaltfläche [Générer]. Es handelt sich um eine ausführbare Klasse, die im folgenden Kontext gestartet wird:
public class JFrameAleasEvents extends AbstractJFrameAleas {
private static final long serialVersionUID = 1L;
// synchroner Generierungsdienst
private IService service;
// asynchroner Erzeugungsdienst
private IRxService rxService;
// die Eingaben
private int nbRequests;
private int a;
private int b;
private int minDelay;
private int maxDelay;
private int minCount;
private int maxCount;
// Fehlermeldungen
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 ";
// Abonnements für Observables
protected List<Subscription> subscriptions = new ArrayList<Subscription>();
// Start und Ende der Ausführung
private long debut;
// Mapper jSON
private ObjectMapper jsonMapper;
// Antwortmodell
private DefaultListModel<String> model;
// Konstruktor
public JFrameAleasEvents() {
// übergeordnetes Element
super();
// lokal
initJFrame();
// Dienstleistungen
service = new Service();
rxService = new RxService(service);
// Mapper jSON
jsonMapper = new ObjectMapper();
}
private void initJFrame() {
// Fehlermeldungen werden ausgeblendet
jLabelCountError.setText("");
jLabelDelayError.setText("");
jLabelIntervalError.setText("");
jLabelNbValuesError.setText("");
// Standardtexte werden ausgeblendet
jTextFieldA.setText("100");
jTextFieldB.setText("200");
jTextFieldMinCount.setText("5");
jTextFieldMaxCount.setText("10");
jTextFieldMinDelay.setText("100");
jTextFieldMaxDelay.setText("500");
jTextFieldNbValeurs.setText("10");
jLabelDuree.setText("");
// Antwortvorlage
model = new DefaultListModel<>();
jListNumbers.setModel(model);
// Anzahl der Kerne
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);
}
/* Formular erstellen und anzeigen */
java.awt.EventQueue.invokeLater(() -> {
new JFrameAleasEvents().setVisible(true);
});
}
- Zeile 1: Die Klasse [JFrameAleasEvents] erweitert die Klasse [AbstractJFrameAleas], die ihrerseits die Swing-Klasse [JFrame] erweitert. Die Klasse [JFrameAleasEvents] ist somit ein Swing-Fenster;
- Zeilen 68–75: Die Methode [main], die ausgeführt wird;
- Zeile 70: Legt das „Look and Feel“ der grafischen Benutzeroberfläche fest;
- Zeile 79: Der Konstruktor der Klasse [JFrameAleasEvents] wird aufgerufen: Die grafische Benutzeroberfläche wird erstellt und initialisiert. Anschließend wird sie sichtbar gemacht;
- Zeilen 34–44: der Konstruktor;
- Zeile 36: Der Aufruf des übergeordneten Konstruktors initialisiert die grafische Benutzeroberfläche. Zu diesem Zeitpunkt entspricht sie genau dem Entwurf des Entwicklers. Sie ist noch nicht sichtbar;
- Zeile 38: Bestimmte Komponenten der grafischen Benutzeroberfläche werden initialisiert;
- Zeile 40: Instanziierung des synchronen Dienstes;
- Zeile 41: Instanziierung des asynchronen Dienstes;
8.8. Ausführung synchroner Abfragen
Ein Klick auf die Schaltfläche [Générer] löst die Ausführung der folgenden Methode [doGenerate] aus:
@Override
protected void doGenerate() {
// Sind die Eingaben gültig?
if (!isPageValid()) {
return;
}
// rx oder nicht?
if (jCheckBoxRxSwing.isSelected()) {
// Asynchrone Anfragen
doGenerateWithRxService();
} else {
// Synchrone Anfragen
doGenerateWithService();
}
}
- Zeilen 4–6: Es wird überprüft, ob die Eingaben des Benutzers gültig sind. Auf die Methode [isPageValid] gehen wir hier nicht näher ein. Sie ist recht einfach;
- Zeile 8: Der Status des Kontrollkästchens RxSwing wird geprüft;
- Zeile 13: Die Abfragen werden synchron ausgeführt;
Die Methode [doGenerateWithService] lautet wie folgt:
// Synchrone Generierung
private void doGenerateWithService() {
// Wartezeit beginnt
beginWaiting();
try {
for (int i = 0; i < nbRequests; i++) {
// Antwortvorbereitung
UiResponse uiResponse = new UiResponse();
// Kundennummer
uiResponse.setIdClient(i);
// synchroner Aufruf
uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
// Antwortzeit
uiResponse.setResponseAt();
// Aktualisierung der Vorlage „JList“ mit den empfangenen Antworten
model.add(0, jsonMapper.writeValueAsString(uiResponse));
// Aktualisierung der Anzahl der Antworten
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);
}
// Wartezeit beendet
endWaiting();
}
- Zeile 12: synchroner Aufruf des Dienstes zur Erzeugung von Zufallszahlen;
- Die Ausführung der Methode [doGenerateWithService] erfolgt vollständig im Thread der Swing-Ereignisschleife. Solange die Methode nicht beendet ist, verarbeitet die grafische Benutzeroberfläche keine neuen Ereignisse. Sie ist eingefroren (frozen). So werden beispielsweise die Aktualisierungen der grafischen Benutzeroberfläche in den Zeilen 16 und 18 niemals angezeigt. Sie sind erst mit ihren Endwerten sichtbar, und zwar erst nach Abschluss der Ausführung aller Abfragen;
Die Methode [beginWaiting] (Zeile 4) lautet wie folgt:
private void beginWaiting() {
// Schaltflächen
jButtonGenerate.setVisible(false);
jButtonCancel.setVisible(true);
// Warte-Cursor
jTabbedPane1.setCursor(Cursor.getPredefinedCursor(Cursor.WAIT_CURSOR));
jButtonCancel.setCursor(Cursor.getDefaultCursor());
// Antworten zurücksetzen
model.clear();
// Rx-Abonnements
subscriptions.clear();
// Ansicht der Antworten wird angezeigt
jTabbedPane1.setSelectedIndex(1);
jLabelNbReponses.setText("0");
jLabelDuree.setText("");
// Ausführung starten
debut = new Date().getTime();
}
- Zeile 3: Die Schaltfläche [Générer] ist ausgeblendet. Dadurch wird ein Ereignis ausgelöst, das ebenfalls erst nach Abschluss der Ausführung aller Abfragen ausgeführt werden kann. Man sieht sie daher nie ausgeblendet, da die Methode [endWaiting] in Zeile 25 der Methode [doGenerateWithService] sie wieder anzeigt;
- Zeile 13: Man wählt die Registerkarte „[Response]“ aus, um die eintreffenden Antworten zu sehen. Auch dieses Ereignis wird erst nach Abschluss aller Abfragen ausgeführt, sodass man dann alle Antworten auf einmal sieht, obwohl man eigentlich wollte, dass sie nacheinander eintreffen;
Die synchrone Schnittstelle weist eindeutig Mängel auf. Diese werden durch die asynchrone Schnittstelle überwunden.
8.9. Ausführung asynchroner Abfragen
Der Code zur Ausführung der asynchronen Abfragen lautet wie folgt:
private void doGenerateWithRxService() {
// Wartephase beginnt
beginWaiting();
// Die Zufallszahlen werden in Form eines Observables abgerufen
Observable<UiResponse> observable = Observable.empty();
// Scheduler für die Ausführung der verschiedenen Observables
Scheduler[] schedulers = { Schedulers.io(), Schedulers.computation(), Schedulers.newThread(),
Schedulers.trampoline(), Schedulers.immediate() };
Scheduler scheduler = schedulers[jComboBoxSchedulers.getSelectedIndex()];
// Konfiguration der Observables
for (int i = 0; i < nbRequests; i++) {
// Antwortvorbereitung
UiResponse uiResponse = new UiResponse();
uiResponse.setIdClient(i);
// Das Observable ist so konfiguriert, dass es auf dem vom Benutzer gewählten Scheduler ausgeführt wird
// Anschließend wird die erhaltene Beobachtungsgröße zur Gesamtbeobachtungsgröße addiert
observable = observable.mergeWith(
rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(scheduler));
}
// Beobachter
observable = observable.observeOn(SwingScheduler.getInstance());
// Bislang haben wir lediglich die Konfiguration vorgenommen
// Es wurde noch keine Anfrage an den synchronen Zufallszahlengenerator gestellt
// Wir abonnieren das Observable – dies löst den Aufruf des synchronen Dienstes zur Zufallszahlengenerierung aus
try {
// Hier gibt es nur ein Abonnement – das Ergebnis ist ein Abonnement
subscriptions.add(observable.subscribe(
// Benachrichtigung über die Ausgabe
uiResponse -> {
// Die Benutzeroberfläche wird mit der Antwort aktualisiert
// Dies ist möglich, da die Beobachtung im UI-Thread stattfindet
updateUi(uiResponse);
} ,
// Fehlermeldung
th -> {
// Fehlerfall – wird angezeigt
String message = getInfoForThrowable("L'erreur suivante s'est produite", th);
JOptionPane.showMessageDialog(this, message, "Informations", JOptionPane.PLAIN_MESSAGE);
// Abfragen abbrechen
doCancel();
} ,
// Benachrichtigung [onCompleted]
// Wartezeit abgelaufen
this::endWaiting));
} catch (Throwable th) {
// Ausnahmefall + allgemein – wird angezeigt
String message = getInfoForThrowable("L'erreur suivante s'est produite", th);
JOptionPane.showMessageDialog(this, message, "Informations", JOptionPane.PLAIN_MESSAGE);
// Anfragen werden abgebrochen
doCancel();
}
}
- Zeile 3: Die grafische Benutzeroberfläche wird angepasst, um anzuzeigen, dass ein möglicherweise langwieriger Vorgang läuft;
- Zeile 5: Es wird ein leeres Observable erstellt. Dieses Observable wird von der Schicht [swing] überwacht;
- Zeile 7: Das Array der möglichen Scheduler;
- Zeile 9: Wir haben dem Benutzer die Möglichkeit gegeben, den Scheduler auszuwählen, auf dem die Abfragen ausgeführt werden sollen. Wir rufen den von ihm gewählten Scheduler ab;
- Zeilen 11–19: Jede der Abfragen gibt ein Observable zurück, dessen Elemente (mergeWith) (Zeile 17) im Observable aus Zeile 5 kumuliert werden;
- Zeilen 13–14: Das Objekt [UiResponse] wird erstellt. Zur Erinnerung: Dieses Objekt ist sowohl Eingabeparameter der Methode [RxService.getAleas] als auch deren Ergebnis (Zeilen 17–18);
- Zeile 14: Jede Anfrage wird durch ihre Nummer gekennzeichnet, die hier [idClient] lautet. Dies ist notwendig, da in einer asynchronen Umgebung die Reihenfolge des Empfangs der Antworten von der Reihenfolge des Versendens der Anfragen abweichen kann. Anhand von [idClient] lässt sich feststellen, zu welcher Anfrage die Antwort gehört;
- Zeilen 17–18: Die asynchrone Anfrage [rxService.getAleas] wird gestellt. Sie wird auf dem vom Benutzer gewählten Scheduler ausgeführt. Ihr Ergebnis vom Typ `Observable<UiResponse>` wird mit dem Observable aus Zeile 5 kumuliert. Man muss sich bewusst machen, dass die Methode [rxService.getAleas] hier ausgeführt wird und ein Observable zurückgibt. Das bedeutet jedoch nicht, dass Zufallszahlen generiert wurden. Tatsächlich wird ein Observable erst ausgeführt, wenn man es abonniert. Das ist noch nicht der Fall;
- Zeile 21: Dies ist die wichtige Anweisung: Es wird festgelegt, dass die Beobachtung der vom Observable in Zeile 5 ausgegebenen Elemente auf dem UI-Thread erfolgen soll. Hier wird ein eigener Scheduler der Bibliothek RxSwing verwendet;
- Zeilen 25–51: Man abonniert das Observable aus Zeile 5. Erst jetzt werden die Zufallszahlen vom synchronen Dienst zur Erzeugung dieser Zahlen angefordert. Das Wesentliche findet sich in den Anweisungen der Zeilen 29–33. Der Rest dient im Wesentlichen der Fehlerbehandlung und der Benachrichtigung des Observables [onCompleted];
- Zeilen 28–44: Man muss bedenken, dass wir die Beobachtung des Prozesses aus Zeile 5 im UI-Thread angefordert haben. Daher wird der Code in den Zeilen 28–44 im UI-Thread ausgeführt;
- Zeilen 29–33: Die Benachrichtigung [onNext] des Observables wird verarbeitet. Wir erhalten einen vom beobachteten Prozess ausgegebenen Typ [UiResponse]. Dies ist das Ergebnis einer der asynchronen Abfragen. Die Benutzeroberfläche wird mit dieser Antwort aktualisiert;
- Zeilen 34–41: Die Benachrichtigung [onError] des Observables wird verarbeitet. Es wird ein Dialogfeld mit dem Fehler angezeigt (Zeilen 37–38), anschließend werden die Anfragen abgebrochen (Zeile 40);
- Zeilen 42–44: Die Benachrichtigung [onCompleted] des Observables wird verarbeitet. Die Benutzeroberfläche wird aktualisiert, um anzuzeigen, dass der angeforderte Dienst abgeschlossen ist. Zeile 44 hätte auch wie folgt geschrieben werden können
Hier wurde jedoch der Verwendung einer Methodenreferenz der Vorzug gegeben;
- Zeilen 45–51: Bestimmte Ausnahmen durchlaufen nicht die Zeilen 34–41. Dies ist der Fall, wenn zu viele Anfragen gestellt werden. Nach Überschreiten einer bestimmten Grenze, die von der Laufzeitumgebung zum Zeitpunkt der Ausführung abhängt, entsteht ein [StackOverflowError], der von den Zeilen 45–51 abgefangen wird;
- Zeile 27: Das Abonnement erzeugt einen Typ [Subscription], der einer Abonnementliste hinzugefügt wird. Diese enthält hier nur ein Element;
In Zeile 32 wird die grafische Benutzeroberfläche mit der folgenden Methode [updateUi] aktualisiert:
private void updateUi(UiResponse uiResponse) {
// Antwortzeit
uiResponse.setResponseAt();
// Beobachtungs-Thread
uiResponse.setObservedOn();
// Anzahl der Antworten
jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
// Ausführungsdauer
jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
// Hinzufügen der Zeichenfolge jSON aus der Antwort zum Antwortmuster JList
try {
model.add(0, jsonMapper.writeValueAsString(uiResponse));
} catch (JsonProcessingException e) {
e.printStackTrace();
}
}
Hier ist zu sehen, dass Komponenten der grafischen Benutzeroberfläche aktualisiert werden (Zeilen 7, 9, 12). Damit dies möglich ist, muss man sich zwingend im UI-Thread (Event-Loop) befinden.
Die Methode [endWaiting] lautet wie folgt:
private void endWaiting() {
// Schaltfläche [Générer] sichtbar
jButtonGenerate.setVisible(true);
// Schaltfläche „[Annuler]“ ausgeblendet
jButtonCancel.setVisible(false);
// Warte-Cursor ausgeblendet
jTabbedPane1.setCursor(Cursor.getDefaultCursor());
// Registerkarte „Antworten“ ausgewählt
jTabbedPane1.setSelectedIndex(1);
// Zeitpunkt der letzten Aktualisierung
jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
}
Die Methode [doCancel] wird aufgerufen, wenn bei der Ausführung asynchroner Abfragen ein Fehler auftritt oder wenn der Benutzer auf die Schaltfläche [Annuler] klickt. Ihr Code lautet wie folgt:
// Abonnements für Beobachtungsgrößen
private List<Subscription> subscriptions = new ArrayList<Subscription>();
....
@Override
protected void doCancel() {
// Wartezeit beendet
endWaiting();
// im Falle von Abonnements
if (jCheckBoxRxSwing.isSelected() && subscriptions != null) {
subscriptions.forEach(Subscription::unsubscribe);
//subscriptions.forEach(s -> s.unsubscribe());
}
}
- Zeile 2: [subscriptions] ist eine Liste eines Abonnements;
- Zeile 11: Alle Abonnements werden gekündigt;
- Zeile 12: eine weitere Darstellung von Zeile 11. Die Methode [forEach] erwartet hier eine Instanz vom Typ Consumer<Subscription> (siehe Abschnitt 4.4);
Kehren wir zum Code der Methode [doGenerateWithService] zurück: Er lässt sich in zwei Schritte unterteilen:
- Schritt zur Konfiguration der Observables. Dies erfolgt im Thread des Aufrufers der Methode [doGenerateWithService], d. h. im Thread der Benutzeroberfläche;
- das Abonnement, das die Ausführung der Observables auslöst;
Wenn die Observables einen der Scheduler „[Schedulers.computation(), Scheduler.io(), Schedulers.newThread()]“ als Scheduler haben, werden sie außerhalb des UI-Threads ausgeführt. Diese verschiedenen Threads konkurrieren dann um den oder die Prozessorkerne des Rechners. Da es sich bei den Abfragen um lang andauernde Vorgänge handelt (mehrere hundert Millisekunden), wird die im UI-Thread ausgeführte Methode [doGenerateWithService] beendet sein, bevor die Abfragen ihre Antworten zurückgegeben haben. Diese Methode war jedoch beim Klick auf die Schaltfläche [Générer] ausgeführt worden. Nachdem dieses Ereignis verarbeitet wurde, kann der UI-Thread mit der Verarbeitung der folgenden Ereignisse fortfahren. Davon gibt es mehrere. So hatte die Methode [beginWaiting] mehrere davon gesetzt:
private void beginWaiting() {
// Schaltflächen
jButtonGenerate.setVisible(false);
jButtonCancel.setVisible(true);
// Warte-Cursor
jTabbedPane1.setCursor(Cursor.getPredefinedCursor(Cursor.WAIT_CURSOR));
jButtonCancel.setCursor(Cursor.getDefaultCursor());
// Antworten löschen
model.clear();
// Rx-Abonnements
subscriptions.clear();
// Antwortansicht wird angezeigt
jTabbedPane1.setSelectedIndex(1);
jLabelNbReponses.setText("0");
jLabelDuree.setText("");
// Ausführung starten
debut = new Date().getTime();
}
Praktisch jede Zeile dieses Codes hat Auswirkungen auf die grafische Benutzeroberfläche. Diese Aktualisierung erfolgt nicht sofort: Ereignisse werden in die Warteschlange der Ereignisschleife gestellt. Sobald das Klick-Ereignis auf die Schaltfläche [Générer] verarbeitet ist, werden diese Ereignisse der Reihe nach ausgeführt, und der Benutzer kann beobachten, wie sich die grafische Benutzeroberfläche ändert:
- Die Registerkarte „[Response]“ wird angezeigt (Zeile 13) und mit einem Wartekursor versehen (Zeile 6)
- die Schaltfläche „[Annuler]“ wird angezeigt (Zeile 4) und der Benutzer kann darauf klicken;
- das Feld „JList“ für die Antworten wird geleert (Zeile 9);
- das Feld JLabel für die Anzahl der Antworten zeigt 0 an;
- das JLabel für die Ausführungsdauer zeigt eine leere Zeichenkette an;
Während der gesamten Ausführungszeit der Abfragen hat der Thread des UI regelmäßig Zugriff auf den Prozessor. Er kann dann die ausstehenden Ereignisse verarbeiten. Dazu gehören auch diejenigen, die von der Methode [updateUi] gesetzt wurden:
private void updateUi(UiResponse uiResponse) {
// Antwortzeit
uiResponse.setResponseAt();
// Beobachtungs-Thread
uiResponse.setObservedOn();
// Anzahl der Antworten
jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
// Ausführungsdauer
jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
// Hinzufügen der Zeichenfolge „jSON“ aus der Antwort zum Muster „JList“ der Antworten
try {
model.add(0, jsonMapper.writeValueAsString(uiResponse));
} catch (JsonProcessingException e) {
e.printStackTrace();
}
}
Wenn der UI-Thread die Kontrolle hat:
- wird der Wert von JLabel für die Anzahl der Antworten aktualisiert (Zeile 7);
- wird das JLabel für die Ausführungsdauer aktualisiert (Zeile 9);
- wird der JList für die Antworten über sein Modell aktualisiert (Zeile 12);
So kann der Benutzer den Fortschritt der Abfrageausführung verfolgen. Außerdem kann er sie über die Schaltfläche [Annuler] abbrechen. Genau darin liegt der Vorteil asynchroner Dienste vor der Schicht [swing], und RxJava ist die Technologie der Wahl für deren Implementierung.
Abschließend sei angemerkt, dass, wenn der Benutzer einen der Scheduler [Schedulers.immediate(), Schedulers.trampoline()] wählt, die Observables auf demselben Thread wie der Aufrufer ausgeführt werden, d. h. dem Thread der Benutzeroberfläche. Man kehrt dann zu einem synchronen Betrieb zurück.
Die mit den verschiedenen Schedulern erzielten Ergebnisse wurden in den Abschnitten 2.8.1, 2.8.2, 2.8.3 und 2.8.4 dargestellt.








