Skip to content

8. RxJava in de Swing-omgeving

8.1. Introduction

We komen hier terug op de Swing-toepassing die in paragraaf 2 is gepresenteerd.

  

Om met RxJava in een Swing-omgeving te werken, gebruiken we de bibliotheek RxSwing, die aan RxJava klassen en interfaces toevoegt die nuttig zijn in een Swing-omgeving. Hiervoor ziet het Gradle-bestand van het Swing-voorbeeld er als volgt uit:

  

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'
}
  • regel 15: de afhankelijkheid van RxSwing;

We zullen slechts één enkel object gebruiken dat specifiek is voor RxSwing: de scheduler [SwingScheduler.getInstance()], die de observables uitvoert en observeert op de thread van de Swing-eventloop. We gaan deze uitsluitend gebruiken om observables te observeren die op andere threads dan die van de event loop worden uitgevoerd. Laten we de architectuur van de voorbeeldtoepassing nog eens in herinnering brengen:

Image

  • De asynchrone servicelaag bevat methoden die observables retourneren. We voeren deze observables uit in threads die verschillen van die van de event loop. Zo blijft de grafische interface niet vastgelopen. Ze kan reageren op acties van de gebruiker. De meest voor de hand liggende is dat de gebruiker op een knop [Annuler] kan klikken om een te lang durende asynchrone bewerking te onderbreken. Om dit te kunnen doen, mag de grafische interface alleen even vastlopen (frozen);
  • de Swing-laag wil de resultaten van de asynchrone bewerkingen verwerken en op basis daarvan de grafische interface bijwerken. Dit kan echter alleen in de thread van de event loop. Hiervoor worden deze resultaten opgepikt in de scheduler [SwingScheduler.getInstance()];

In de code voor de gebeurtenisafhandeling van de grafische interface verloopt de interactie met de asynchrone laag [rxService] dus als volgt:


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

waarbij de scheduler [Schedulers.computation()] afhankelijk van het gebruiksscenario door een andere scheduler kan worden vervangen.

De lezer wordt verzocht paragraaf 2 nogmaals door te nemen. Hij beschikt nu over de nodige kennis om deze volledig te begrijpen.

8.2. De structuur van de code

De code implementeert de volgende architectuur:

Image

Het IntelliJ Idea-project dat deze architectuur implementeert, is het volgende:

  
  • het pakket [rxswing.service] implementeert de synchrone (IService, Service) en asynchrone (IRxService, RxService) servicelagen;
  • het pakket [rxswing.ui] implementeert de Swing-interface;

8.3. Het project uitvoeren

Ga als volgt te werk om het project in IntelliJ IDEA uit te voeren:

 

8.4. De synchrone service

Image

  

De synchrone servicelaag heeft de volgende interface: [IService]:


package dvp.rxswing.service;

public interface IService {
  // willekeurige getallen in het interval [a,b]
  // n getallen worden gegenereerd, waarbij n zelf een willekeurig getal is in het interval [minCount, maxCount]
  // de getallen worden gegenereerd na een wachttijd van delay milliseconden,
  // waarbij [delay] zelf een willekeurig getal is in het interval [minDelay, maxDelay]
  public ServiceResponse getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay);
}

Het type [ServiceResponse] van het antwoord van de service is als volgt:


package dvp.rxswing.service;

import java.util.List;

public class ServiceResponse {

  // wachttijd van de dienst
  private int delay;
  // willekeurige getallen
  private List<Integer> aleas;
  // uitvoeringsthread
  private String executedOn;

  // constructors

  public ServiceResponse() {
      // uitvoeringsthread
    executedOn = Thread.currentThread().getName();
  }

  public ServiceResponse(int delay, List<Integer> aleas) {
      // lokale constructor
    this();
    // overige initialisaties
    this.delay = delay;
    this.aleas = aleas;
  }

  // getters en setters
...
}

De interface [IService] wordt geïmplementeerd door de volgende klasse [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) {
    // willekeurige getallen in het interval [a,b]
    // n getallen worden gegenereerd, waarbij n zelf een willekeurig getal is in het interval [minCount, maxCount]
    // de getallen worden gegenereerd na een wachttijd van delay milliseconden,
    // waarbij [delay] zelf een willekeurig getal is in het interval [minDelay, maxDelay]

    // enkele controles
    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;
    }
    // fouten?
    if (!messages.isEmpty()) {
      throw new AleasException(String.join(" [---] ", messages), erreur);
    }
    // willekeurige-getallengenerator
    Random random = new Random();
    // in afwachting?
    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);
      }
    }
    // resultaat genereren
    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));
    }
    // resultaat teruggeven
    return new ServiceResponse(delay,nombres);
  }

}

De uitzonderingsklasse [AleasException] die door de service wordt gebruikt, is als volgt:


package dvp.rxswing.service;

public class AleasException extends RuntimeException {

    private static final long serialVersionUID = 1L;
    // foutcode
  private int code;

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

  // getters en setters
...
}
  • regel 3: deze breidt de klasse [RuntimeException] uit. Het gaat dus om een ongecontroleerde uitzondering;
  • regel 7: deze voegt een foutcode toe aan de bovenliggende klasse (0 = geen fout);

8.5. De asynchrone service

Image

  

De asynchrone servicelaag heeft de volgende interface [IRxService]:


package dvp.rxswing.service;

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

public interface IRxService {
  // willekeurige getallen in het interval [a,b]
  // n getallen worden gegenereerd, waarbij n zelf een willekeurig getal is in het interval [minCount, maxCount]
  // de getallen worden gegenereerd na een wachttijd van delay milliseconden,
  // waarbij [delay] zelf een willekeurig getal is in het interval [minDelay, maxDelay]
  public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse);
}
  • regel 11: de methode [getAleas] van de service retourneert nu een observable;

De methode [getAleas] retourneert een antwoord van het type [UiResponse], bestemd voor de laag [Ui]. Dit type is als volgt:


package dvp.rxswing.ui;

import dvp.rxswing.service.ServiceResponse;

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

public class UiResponse {

  // klant-id
  private int idClient;
  // antwoord van de dienst
  private ServiceResponse serviceResponse;
  // naam van de observatiethread
  private String observedOn;
  // tijdstip van het verzoek
  private String requestAt;
  // tijdstip van het antwoord
  private String responseAt;

  // constructoren

  public UiResponse() {
      // observatiethread
    observedOn = Thread.currentThread().getName();
    // tijdstip van de aanvraag
    requestAt = getTimeStamp();
  }

  // privé-methoden

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

  // getters en setters
...
}
  • de willekeurige getallen staan in het veld op regel 13;
  • de overige velden dienen om de uitvoerings- en observatiethreads van de observabele waarde van de asynchrone service te specificeren, evenals de tijdstippen van het verzoek aan de service en van het ontvangen antwoord;

De asynchrone interface wordt geïmplementeerd door de volgende klasse [RxService]:


package dvp.rxswing.service;

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

public class RxService implements IRxService {

  // synchrone service
  private IService service;

  // constructor
  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) {
      // we maken een observable aan die de waarde doorgeeft die door de synchrone service wordt geretourneerd
    return Observable.create(subscriber -> {
      try {
          // synchrone aanroep
        uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
        // het resultaat wordt doorgegeven aan de observer
        subscriber.onNext(uiResponse);
      } catch (Exception e) {
          // de fout wordt doorgegeven aan de observer
        subscriber.onError(e);
      } finally {
          // de observer wordt geïnformeerd dat de verzendingen zijn voltooid
        subscriber.onCompleted();
      }
    });
  }
}
  • regels 12-14: de klasse [RxService] van de asynchrone service wordt geconstrueerd op basis van een instantie van de synchrone interface [IService];
  • regels 20-33: aanmaak van de observable, het resultaat van de methode [getAleas];
  • regel 22: de synchrone methode [service.getAleas] wordt aangeroepen. Het resultaat ervan, van het type [ServiceResponse], wordt opgenomen in het object van het type [UiResponse] dat aan de laag [swing] moet worden geleverd. Dit object is aanvankelijk doorgegeven in de aanroepparameters van de methode (laatste parameter, regel 17);
  • regel 24: het antwoord [UiResponse] wordt naar de waarnemer (de laag [swing]) verzonden. Het object [UiResponse] bevat niet alleen de informatie die door de synchrone service in regel 22 is samengesteld. Het bevat ook andere informatie die is samengesteld door de aanroepende methode van de methode [getAleas] uit regel 17. Om deze reden heeft deze aanroepende methode het object [UiResponse] als parameter doorgegeven aan de methode [getAleas] (laatste parameter, regel 17);
  • regel 30: we vergeten niet het einde van de uitzendingen te melden. Hier hebben we een observabel die slechts één waarde uitzendt: de waarde die door de synchrone dienst wordt geretourneerd;
  • regel 27: een eventuele fout wordt aan de waarnemer gemeld;

8.6. De grafische interface

Image

  
  • de grafische interface is gebouwd met de IDE [Netbeans], die over een goede grafische editor beschikt. Deze editor heeft het bestand [AbstractJFrameAleas.form] gegenereerd, dat alleen door deze IDE kan worden gebruikt;
  • de klasse [AbstractJFrameAleas] is eveneens gegenereerd door de grafische editor van NetBeans. Vervolgens is deze als volgt geherstructureerd: de gebeurtenissen van de grafische interface die we wilden afhandelen, worden in de klasse [AbstractJFrameAleas] verwerkt door middel van abstracte methoden die zijn geïmplementeerd in de onderliggende klasse [JFrameAleasEvents]. Uiteindelijk,
    • de abstracte klasse [AbstractJFrameAleas] zorgt voor het opbouwen en weergeven van de grafische interface;
    • de onderliggende klasse [JFrameAleasEvents] is verantwoordelijk voor het beheer van de gebeurtenissen daarvan;

De componenten van de grafische interface van het tabblad [Request] zijn de volgende:

 
nr.
type
naam
rol
1
JTabbedPane
jTabbedPane1
een tabbladcontainer. Bevat twee tabbladen (JPanel) [jPanelRequest] voor de aanvraag, [jPanelresponse] voor het antwoord;
2
JTextField
jTextFieldNbValeurs
het aantal verzoeken dat aan de dienst voor willekeurige getallen moet worden gedaan. In het geval van de asynchrone dienst die wordt uitgevoerd op de scheduler [Schedulers.io], zullen deze verzoeken één processor delen;
3
JTextField
jTextFieldA
eindpunt a van het interval [a,b]
4
JTextField
jTextFieldB
aansluitpunt b van het interval [a,b]
5
JTextField
jTextFieldMinCount
aansluiting minCount van het interval [minCount, maxCount]
6
JTextField
jTextFieldMaxCount
aansluitpunt maxCount van het interval [minCount, maxCount]
7
JTextField
jTextFieldMinDelay
aansluitpunt minDelay van het interval [minDelay, maxDelay]
8
JTextField
jTextFieldMaxDelay
aansluiting maxDelay van het interval [minDelay, maxDelay]
9
JCheckBox
jCheckBoxRxSwing
als het vakje is aangevinkt, worden de verzoeken via de asynchrone interface verzonden. Anders worden ze via de synchrone interface verzonden
10
JComboBox
jComboBoxSchedulers
bij asynchrone verzoeken worden deze uitgevoerd met de hier gekozen planner
11
JButton
jButtonGenerate
start de uitvoering van de verzoeken naar de synchrone of asynchrone service

De componenten van de grafische interface van het tabblad [Response] zijn de volgende:

 
nr.
type
naam
rol
1
JLabel
jLabelDuree
de totale uitvoeringstijd in milliseconden van de verzoeken
2
JLabel
jLabelNbReponses
het totale aantal waargenomen antwoorden (kan afwijken van het aantal verzoeken, aangezien elk verzoek meerdere waargenomen waarden kan opleveren)
3
JList
jListNumbers
weergave van de waargenomen (ontvangen) waarden
4
JButton
jButtonAnnuler
Annuleert de lopende verzoeken

8.7. Het grafische gebruikersinterface wordt geïnitialiseerd

  

De klasse [JFrameAleasEvents] beheert de gebeurtenissen van de grafische interface, met name het klikken op de knop [Générer]. Het is een uitvoerbare klasse die in de volgende context wordt gestart:


public class JFrameAleasEvents extends AbstractJFrameAleas {

    private static final long serialVersionUID = 1L;
    // synchrone generatiedienst
    private IService service;
    // asynchrone opwekkingsdienst
    private IRxService rxService;

    // de invoergegevens
    private int nbRequests;
    private int a;
    private int b;
    private int minDelay;
    private int maxDelay;
    private int minCount;
    private int maxCount;

    // foutmeldingen
    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 ";

    // abonnementen op observables
    protected List<Subscription> subscriptions = new ArrayList<Subscription>();
    // begin-einde van de uitvoering
    private long debut;
    // mapper jSON
    private ObjectMapper jsonMapper;
    // antwoordmodel
    private DefaultListModel<String> model;

    // constructor
    public JFrameAleasEvents() {
        // bovenliggend
        super();
        // lokaal
        initJFrame();
        // diensten
        service = new Service();
        rxService = new RxService(service);
        // mapper jSON
        jsonMapper = new ObjectMapper();
    }

    private void initJFrame() {
        // foutmeldingen worden verborgen
        jLabelCountError.setText("");
        jLabelDelayError.setText("");
        jLabelIntervalError.setText("");
        jLabelNbValuesError.setText("");
        // standaardteksten worden verborgen
        jTextFieldA.setText("100");
        jTextFieldB.setText("200");
        jTextFieldMinCount.setText("5");
        jTextFieldMaxCount.setText("10");
        jTextFieldMinDelay.setText("100");
        jTextFieldMaxDelay.setText("500");
        jTextFieldNbValeurs.setText("10");
        jLabelDuree.setText("");
        // antwoordmodel
        model = new DefaultListModel<>();
        jListNumbers.setModel(model);
        // aantal kernen
        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);
        }

        /* Het formulier aanmaken en weergeven */
        java.awt.EventQueue.invokeLater(() -> {
            new JFrameAleasEvents().setVisible(true);
        });
    }
  • regel 1: de klasse [JFrameAleasEvents] is een uitbreiding van de klasse [AbstractJFrameAleas], die op haar beurt weer een uitbreiding is van de Swing-klasse [JFrame]. De klasse [JFrameAleasEvents] is dus een Swing-venster;
  • regels 68-75: de methode [main] die zal worden uitgevoerd;
  • regel 70: stelt de look-and-feel van de grafische interface in;
  • regel 79: de constructor van de klasse [JFrameAleasEvents] wordt aangeroepen: de grafische interface wordt opgebouwd en geïnitialiseerd. Zodra dit is gebeurd, wordt deze zichtbaar gemaakt;
  • regels 34-44: de constructor;
  • regel 36: de aanroep van de bovenliggende constructor initialiseert de grafische interface. Op dat moment ziet deze eruit zoals de ontwikkelaar deze heeft ontworpen. De interface is nog niet zichtbaar;
  • regel 38: bepaalde componenten van de grafische interface worden geïnitialiseerd;
  • regel 40: instantiatie van de synchrone service;
  • regel 41: instantiëren van de asynchrone service;

8.8. Uitvoering van synchrone verzoeken

Als u op de knop [Générer] klikt, wordt de volgende methode [doGenerate] uitgevoerd:


    @Override
    protected void doGenerate() {
        // geldige invoer?
        if (!isPageValid()) {
            return;
        }
        // rx of niet?
        if (jCheckBoxRxSwing.isSelected()) {
            // asynchrone verzoeken
            doGenerateWithRxService();
        } else {
            // synchrone verzoeken
            doGenerateWithService();
        }
}
  • regels 4-6: er wordt gecontroleerd of de invoer van de gebruiker geldig is. We zullen geen toelichting geven bij de methode [isPageValid]. Deze is eenvoudig;
  • regel 8: de status van het selectievakje RxSwing wordt getest;
  • regel 13: de query's worden synchroon uitgevoerd;

De methode [doGenerateWithService] is als volgt:


    // synchrone generatie
    private void doGenerateWithService() {
        // begin wachttijd
        beginWaiting();
        try {
            for (int i = 0; i < nbRequests; i++) {
                // voorbereiding van het antwoord
                UiResponse uiResponse = new UiResponse();
                // klantnummer
                uiResponse.setIdClient(i);
                // synchrone oproep
                uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
                // tijdstip van antwoord
                uiResponse.setResponseAt();
                // het model van JList bijwerken met de ontvangen antwoorden
                model.add(0, jsonMapper.writeValueAsString(uiResponse));
                // het aantal antwoorden bijgewerkt
                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);
        }
        // afgewerkt
        endWaiting();
}
  • regel 12: synchrone aanroep van de dienst voor het genereren van willekeurige getallen;
  • de uitvoering van de methode [doGenerateWithService] vindt volledig plaats in de thread van de Swing-eventloop. Zolang de methode niet is voltooid, verwerkt de grafische interface geen nieuwe gebeurtenissen. Deze is bevroren (frozen). Zo zullen bijvoorbeeld de updates van de grafische interface in de regels 16 en 18 nooit te zien zijn. Ze zullen pas zichtbaar zijn met hun uiteindelijke waarden, en wel aan het einde van de uitvoering van alle verzoeken;

De methode [beginWaiting] (regel 4) is als volgt:


    private void beginWaiting() {
        // knoppen
        jButtonGenerate.setVisible(false);
        jButtonCancel.setVisible(true);
        // wachtcursor
        jTabbedPane1.setCursor(Cursor.getPredefinedCursor(Cursor.WAIT_CURSOR));
        jButtonCancel.setCursor(Cursor.getDefaultCursor());
        // antwoorden wissen
        model.clear();
        // Rx-abonnementen
        subscriptions.clear();
        // weergave van antwoorden
        jTabbedPane1.setSelectedIndex(1);
        jLabelNbReponses.setText("0");
        jLabelDuree.setText("");
        // start uitvoering
        debut = new Date().getTime();
}
  • regel 3: de knop [Générer] is verborgen. Dit creëert een gebeurtenis die eveneens pas kan worden uitgevoerd nadat alle verzoeken zijn uitgevoerd. We zien deze knop dus nooit verborgen, omdat de methode [endWaiting] op regel 25 van de methode [doGenerateWithService] hem weer weergeeft;
  • regel 13: we selecteren het tabblad [Response] om de antwoorden binnen te zien komen. Ook deze gebeurtenis wordt pas uitgevoerd nadat alle verzoeken zijn afgehandeld, waarna we alle antwoorden tegelijk te zien krijgen, terwijl we ze juist één voor één wilden zien binnenkomen;

De synchrone interface vertoont duidelijk tekortkomingen. Deze worden overwonnen dankzij de asynchrone interface.

8.9. Uitvoering van asynchrone verzoeken

De code voor het uitvoeren van asynchrone verzoeken is als volgt:


private void doGenerateWithRxService() {
        // wachtperiode begint
        beginWaiting();
        // we gaan de willekeurige getallen verkrijgen in de vorm van een observable
        Observable<UiResponse> observable = Observable.empty();
        // Uitvoeringsschema voor de verschillende observables
        Scheduler[] schedulers = { Schedulers.io(), Schedulers.computation(), Schedulers.newThread(),
                Schedulers.trampoline(), Schedulers.immediate() };
        Scheduler scheduler = schedulers[jComboBoxSchedulers.getSelectedIndex()];
        // configuratie van de observables
        for (int i = 0; i < nbRequests; i++) {
            // voorbereiding van het antwoord
            UiResponse uiResponse = new UiResponse();
            uiResponse.setIdClient(i);
            // de observable is geconfigureerd om te worden uitgevoerd op de door de gebruiker gekozen planner
            // vervolgens wordt de verkregen observable samengevoegd met de totale observable
            observable = observable.mergeWith(
                    rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(scheduler));
        }
        // waarnemer
        observable = observable.observeOn(SwingScheduler.getInstance());
        // tot nu toe hebben we alleen de configuratie uitgevoerd
        // er is nog geen verzoek gedaan aan de synchrone dienst voor het genereren van willekeurige getallen
        // we abonneren ons op de observable – dit zal de aanroep naar de synchrone dienst voor het genereren van willekeurige getallen activeren
        try {
            // hier is er slechts één abonnement – het resultaat is een inschrijving
            subscriptions.add(observable.subscribe(
                    // uitzendmelding
                    uiResponse -> {
                        // de UI wordt bijgewerkt met het antwoord
                        // dit is mogelijk omdat de bewerking plaatsvindt in de UI-thread
                        updateUi(uiResponse);
                    } ,
                    // foutmelding
                    th -> {
                        // foutgeval – wordt weergegeven
                        String message = getInfoForThrowable("L'erreur suivante s'est produite", th);
                        JOptionPane.showMessageDialog(this, message, "Informations", JOptionPane.PLAIN_MESSAGE);
                        // verzoeken annuleren
                        doCancel();
                    } ,
                    // melding [onCompleted]
                    // einde van het wachten
                    this::endWaiting));
        } catch (Throwable th) {
            // uitzonderingsgeval + algemeen - wordt weergegeven
            String message = getInfoForThrowable("L'erreur suivante s'est produite", th);
            JOptionPane.showMessageDialog(this, message, "Informations", JOptionPane.PLAIN_MESSAGE);
            // verzoeken worden geannuleerd
            doCancel();
        }
    }
  • regel 3: de grafische interface wordt aangepast om aan te geven dat er een mogelijk langdurige bewerking aan de gang is;
  • regel 5: er wordt een lege observable aangemaakt. Deze observable wordt door de laag [swing] geobserveerd;
  • regel 7: de lijst met mogelijke schedulers;
  • regel 9: we hebben de gebruiker de mogelijkheid gegeven om de scheduler te kiezen waarop de verzoeken moeten worden uitgevoerd. We halen de door hem gekozen scheduler op;
  • regels 11-19: elk van de query’s retourneert een observable waarvan de elementen worden samengevoegd (mergeWith) (regel 17) in de observable van regel 5;
  • regels 13-14: het object [UiResponse] wordt aangemaakt. Ter herinnering: dit object is zowel de invoerparameter van de methode [RxService.getAleas] als het resultaat ervan (regels 17-18);
  • regel 14: elk verzoek wordt geïdentificeerd aan de hand van zijn nummer, hier [idClient] genoemd. Dit is noodzakelijk omdat in een asynchrone omgeving de volgorde waarin de antwoorden worden ontvangen, kan afwijken van de volgorde waarin de verzoeken worden verzonden. Aan de hand van [idClient] kan worden vastgesteld bij welk verzoek het antwoord hoort;
  • regels 17-18: het asynchrone verzoek wordt gedaan: [rxService.getAleas]. Het wordt uitgevoerd op de door de gebruiker gekozen scheduler. Het resultaat, van het type Observable<UiResponse>, wordt samengevoegd met de observable uit regel 5. Men moet zich goed realiseren dat de methode [rxService.getAleas] hier wordt uitgevoerd en een observable retourneert. Dit betekent echter niet dat er willekeurige getallen zijn gegenereerd. Een observable wordt namelijk pas uitgevoerd wanneer men zich erop abonneert. Dat is nog niet het geval;
  • regel 21: dit is de belangrijke instructie: er wordt gevraagd om de waarneming van de elementen die door de observable van regel 5 worden uitgezonden, uit te voeren op de UI-thread. Hier wordt gebruikgemaakt van een eigen scheduler van de bibliotheek RxSwing;
  • regels 25-51: we abonneren ons op de observable uit regel 5. Pas nu worden de willekeurige getallen opgevraagd bij de synchrone dienst die deze getallen genereert. Het belangrijkste zit in de instructies van de regels 29-33. De rest houdt zich voornamelijk bezig met foutafhandeling en de [onCompleted]-melding van de observable;
  • regels 28-44: we moeten niet vergeten dat we hebben gevraagd om het proces van regel 5 te observeren in de UI-thread. De code in de regels 28-44 wordt dus uitgevoerd in de UI-thread;
  • regels 29-33: we verwerken de melding [onNext] van de observable. We ontvangen een type [UiResponse] dat is verzonden door het geobserveerde proces. Dit is het resultaat van een van de asynchrone verzoeken. We werken de grafische interface bij met dit antwoord;
  • regels 34-41: de melding [onError] van de observable wordt verwerkt. Er wordt een dialoogvenster weergegeven met de foutmelding (regels 37-38) en vervolgens worden de verzoeken geannuleerd (regel 40);
  • regels 42-44: de melding [onCompleted] van de observable wordt verwerkt. De grafische interface wordt bijgewerkt om aan te geven dat de aangevraagde dienst is voltooid. Regel 44 had ook als volgt geschreven kunnen worden
 ()->{endWaiting();}

Hier is gekozen voor het gebruik van een methodeverwijzing;

  • regels 45-51: bepaalde uitzonderingen lopen niet via de regels 34-41. Dit is het geval wanneer er te veel verzoeken worden ingediend. Zodra een bepaalde limiet wordt overschreden – die afhankelijk is van de werkomgeving op het moment van uitvoering – ontstaat er een [StackOverflowError] die wordt opgevangen door de regels 45-51;
  • regel 27: het abonnement genereert een type [Subscription] dat wordt toegevoegd aan een lijst met abonnementen. Deze lijst bevat hier slechts één element;

In regel 32 wordt de grafische interface bijgewerkt met de volgende methode [updateUi]:


    private void updateUi(UiResponse uiResponse) {
        // tijdstip van antwoord
        uiResponse.setResponseAt();
        // observatiethread
        uiResponse.setObservedOn();
        // aantal antwoorden
        jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
        // uitvoeringstijd
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
        // toevoeging van de tekenreeks jSON uit het antwoord aan het sjabloon van JList voor antwoorden
        try {
            model.add(0, jsonMapper.writeValueAsString(uiResponse));
        } catch (JsonProcessingException e) {
            e.printStackTrace();
        }
}

We zien hier dat componenten van de grafische interface worden bijgewerkt (regels 7, 9, 12). Om dit mogelijk te maken, moet men zich verplicht in de UI-thread (event loop) bevinden.

De methode [endWaiting] is als volgt:


    private void endWaiting() {
        // knop [Générer] zichtbaar
        jButtonGenerate.setVisible(true);
        // knop [Annuler] verborgen
        jButtonCancel.setVisible(false);
        // wachtscursor verborgen
        jTabbedPane1.setCursor(Cursor.getDefaultCursor());
        // tabblad 'Antwoorden' geselecteerd
        jTabbedPane1.setSelectedIndex(1);
        // laatste update
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
}

De methode [doCancel] wordt aangeroepen wanneer er een fout optreedt bij de uitvoering van asynchrone verzoeken of wanneer de gebruiker op de knop [Annuler] klikt. De code ervan is als volgt:


// abonnementen op observables
    private List<Subscription> subscriptions = new ArrayList<Subscription>();
....

    @Override
    protected void doCancel() {
        // einde wachttijd
        endWaiting();
        // in het geval van abonnementen
        if (jCheckBoxRxSwing.isSelected() && subscriptions != null) {
            subscriptions.forEach(Subscription::unsubscribe);
            //subscriptions.forEach(s -> s.unsubscribe());
        }
    }

  • regel 2: [subscriptions] is een lijst met een abonnement;
  • regel 11: alle abonnementen worden opgezegd;
  • regel 12: een andere weergave van regel 11. De methode [forEach] verwacht hier een instantie van het type Consumer<Subscription> (zie paragraaf 4.4);

Laten we terugkeren naar de code van de methode [doGenerateWithService]: deze kan worden opgesplitst in twee stappen:

  1. de stap voor het configureren van de observables. Dit gebeurt in de thread van de aanroeper van de methode [doGenerateWithService], d.w.z. de thread van de UI;
  2. het abonnement dat ervoor zorgt dat de observables worden uitgevoerd;

Als de observables worden ingepland door een van de schedulers [Schedulers.computation(), Scheduler.io(), Schedulers.newThread()], dan worden ze buiten de UI-thread uitgevoerd. Deze verschillende threads zullen strijden om de processor(en) van de machine. Aangezien de verzoeken langdurige bewerkingen zijn (enkele honderden milliseconden), zal de methode [doGenerateWithService] die in de UI-thread wordt uitgevoerd, worden voltooid voordat de verzoeken hun antwoorden hebben teruggestuurd. Deze methode was echter uitgevoerd bij de klik op de knop [Générer]. Nu deze gebeurtenis is verwerkt, kan de UI-thread overgaan tot de verwerking van de volgende gebeurtenissen. Er zijn er meerdere. Zo had de methode [beginWaiting] er meerdere ingesteld:


    private void beginWaiting() {
        // knoppen
        jButtonGenerate.setVisible(false);
        jButtonCancel.setVisible(true);
        // wachtschaduw
        jTabbedPane1.setCursor(Cursor.getPredefinedCursor(Cursor.WAIT_CURSOR));
        jButtonCancel.setCursor(Cursor.getDefaultCursor());
        // antwoorden wissen
        model.clear();
        // Rx-abonnementen
        subscriptions.clear();
        // de weergave van de antwoorden wordt getoond
        jTabbedPane1.setSelectedIndex(1);
        jLabelNbReponses.setText("0");
        jLabelDuree.setText("");
        // start uitvoering
        debut = new Date().getTime();
}

Vrijwel alle regels van deze code hebben invloed op de grafische interface. Deze update vindt niet onmiddellijk plaats: er worden gebeurtenissen in de wachtrij van de event loop geplaatst. Zodra de klikgebeurtenis op de knop [Générer] is verwerkt, worden deze gebeurtenissen op hun beurt uitgevoerd en kan de gebruiker zien hoe de grafische interface verandert:

  • het tabblad [Response] wordt weergegeven (regel 13) en er wordt een wachtcursor aan toegewezen (regel 6)
  • de bijbehorende knop [Annuler] wordt weergegeven (regel 4) en de gebruiker kan erop klikken;
  • het JList-veld voor antwoorden wordt leeggemaakt (regel 9);
  • de JLabel voor het aantal antwoorden geeft 0 weer;
  • de JLabel voor de uitvoeringstijd geeft een lege tekenreeks weer;

Gedurende de gehele uitvoering van de verzoeken heeft de thread van de UI regelmatig toegang tot de processor. Hij kan dan de wachtende gebeurtenissen verwerken. Hieronder vallen onder andere de gebeurtenissen die door de methode [updateUi] zijn ingesteld:


    private void updateUi(UiResponse uiResponse) {
        // tijd van antwoord
        uiResponse.setResponseAt();
        // observatiethread
        uiResponse.setObservedOn();
        // aantal antwoorden
        jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
        // uitvoeringstijd
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
        // toevoeging van de tekenreeks jSON uit het antwoord aan het sjabloon van JList van de antwoorden
        try {
            model.add(0, jsonMapper.writeValueAsString(uiResponse));
        } catch (JsonProcessingException e) {
            e.printStackTrace();
        }
}

Wanneer de UI-thread aan de beurt is:

  • wordt de JLabel van het aantal reacties bijgewerkt (regel 7);
  • wordt de JLabel van de uitvoeringstijd bijgewerkt (regel 9);
  • de JList van de antwoorden wordt bijgewerkt via het bijbehorende model (regel 12);

Zo ziet de gebruiker de voortgang van de uitvoering van de verzoeken. Bovendien kan hij ze annuleren via de knop [Annuler]. Dit is precies het voordeel van asynchrone services voor de [swing]-laag, en RxJava is de technologie bij uitstek om deze te implementeren.

Tot slot merken we op dat, als de gebruiker een van de [Schedulers.immediate(), Schedulers.trampoline()]-planners kiest, de observables worden uitgevoerd op dezelfde thread als de aanroeper, d.w.z. de thread van de UI. We keren dan terug naar een synchrone werking.

De resultaten die met de verschillende schedulers zijn verkregen, zijn weergegeven in de paragrafen 2.8.1, 2.8.2, 2.8.3 en 2.8.4.