20. Асинхронне програмування з RxJava
Матеріал для ознайомлення: [Introduction à RxJava. Application aux environnements Swing et Android.]
У цьому розділі ми повернемося до розділу 17.6, де ми створили клієнт-серверний додаток із такою архітектурою:
![]() |
Деякі дії користувача в інтерфейсі Swing у [1] запускають дії аж у базі даних у [3] через мережу HTTP [2]. Через це відповідь на дію користувача може надходити з більшою чи меншою затримкою. Було б добре мати можливість розмістити індикатор очікування на користувацькому інтерфейсі з можливістю скасування запущеної операції, якщо вона затягнеться надто довго. У розділі 17.6 кожна дія користувача, що вимагає обміну інформацією з сервером, є синхронною. Обробник події, що виконується кодом, завершується лише після отримання відповіді. Протягом усього цього часу графічний інтерфейс «зависає»: він не реагує на нові дії користувача. Ці дії просто ставляться в чергу для обробки після завершення подійного обробника, що виконується в даний момент. Таким чином, якби з’явилася кнопка скасування, користувач міг би натиснути на неї, але нічого не відбувалося б, доки поточна операція не завершилася. Тоді кнопка скасування не мала б жодного сенсу.
Щоб натискання кнопки скасування дало результат, поточна операція має бути завершена. Для цього вона повинна запустити операцію, яка може тривати довго, в асинхронному режимі:
- обробник події запускає тривалу операцію, але не чекає на її результат і передає управління потоку UI, який обробляє події графічного інтерфейсу. Тривала операція запускається в потоці, відмінному від потоку UI, що не блокує останній;
- якщо користувач натисне кнопку скасування до завершення тривалої операції, порожній потік UI може обробити цю подію. У такому разі можна припинити тривалу операцію, проігнорувавши її результат;
- якщо тривала операція не була скасована, надходження відповіді викличе подію в потоці UI. Якщо цей потік вільний, він виконає код, пов’язаний із цією подією, який обробить відповідь;
Інтерфейс користувача працюватиме так само, як і раніше. Якщо час відгуку сервера невеликий, користувач не помітить різниці. Якщо ж затримка помітна, користувачеві з’явиться кнопка скасування, і він матиме можливість перервати поточну операцію.
Бібліотека [Rx] дозволяє здійснювати асинхронне програмування. Її велика цінність полягає в тому, що вона була портирована на багато середовищ (Java, .NET, JS, ...) і що навички роботи з нею в одному середовищі можна легко застосувати в іншому. Тут ми будемо спиратися на розділ 2 документа [Introduction à RxJava. Application aux environnements Swing et Android]. Читачеві рекомендується ознайомитися з ним. Далі ми використаємо код із прикладів цього розділу.
Ми будемо розвивати архітектуру додатка таким чином:
![]() |
- у [1] ми вставляємо рівень [RxJava] між рівнями [swing] та [métier]. Методи цього рівня відтепер викликатимуться асинхронно;
Ми будемо діяти в кілька етапів:
- етап 1: на даний момент шар [metier, DAO] має синхронний інтерфейс до шару [ui]. Ми перетворимо його на асинхронний шар [RxJava, metier, DAO];
- етап 2: ми перенесемо синхронний консольний додаток у додаток, який залишиться синхронним, але використовуватиме асинхронний інтерфейс [RxJava, metier, DAO];
- етап 3: ми перенесемо синхронний додаток Swing на асинхронний додаток Swing;
20.1. етап 1
Ми перетворюємо поточний синхронний рівень [metier, DAO] на асинхронний рівень [RxJava, metier, DAO].
20.1.1. Створення
Ми беремо за основу проект Maven із розділу 17.4, який відкриваємо за допомогою NetBeans:
![]() | ![]() |
Ми дублюємо цей проєкт [1] (копіюємо та вставляємо) у новий проєкт [elections-rxjava-metier-dao-security-webjson] [2].
20.1.2. Налаштування Maven
Ми оновлюємо файл [pom.xml] нового проєкту, щоб додати залежність від бібліотеки [RxJava]:
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>istia.st.elections</groupId>
<artifactId>elections-metier-dao-security-rxjava-webjson</artifactId>
<version>0.0.1-SNAPSHOT</version>
<description>Client jUnit du serveur web / jSON</description>
<name>elections-metier-dao-security-rxjava-webjson</name>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<java.version>1.8</java.version>
</properties>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>1.2.7.RELEASE</version>
</parent>
<dependencies>
<!-- Spring -->
<dependency>
<groupId>org.springframework</groupId>
<artifactId>spring-web</artifactId>
</dependency>
<!-- бібліотека jSON, яку використовує Spring -->
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-core</artifactId>
</dependency>
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
</dependency>
<!-- компонент, що використовується Spring RestTemplate -->
<dependency>
<groupId>org.apache.httpcomponents</groupId>
<artifactId>httpclient</artifactId>
</dependency>
<!-- Google Guava -->
<dependency>
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
<version>16.0.1</version>
<scope>test</scope>
</dependency>
<!-- бібліотека логів -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-logging</artifactId>
</dependency>
<!-- Spring Boot Test -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<!-- Spring Boot -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot</artifactId>
<scope>test</scope>
</dependency>
<!-- https://mvnrepository.com/artifact/io.reactivex/rxjava -->
<dependency>
<groupId>io.reactivex</groupId>
<artifactId>rxjava</artifactId>
<version>1.2.0</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-surefire-plugin</artifactId>
<version>2.18.1</version>
</plugin>
</plugins>
</build>
</project>
- рядки 65–70: ми додали залежність від бібліотеки RxJava;
20.1.3. Асинхронна реалізація шару [métier]
![]() |
Для реалізації шару [RxJava, métier] ми додаємо асинхронний інтерфейс [IRxElectionsMetier] [1] та його реалізацію [RxElectionsMetier] [2] до проєкту:
![]() |
Інтерфейс [IRxElectionsMetier] є асинхронним інтерфейсом шару [RxJava, métier]. Його код такий:
package elections.security.client.metier;
import elections.security.client.entities.ListeElectorale;
import elections.security.client.entities.User;
import rx.Observable;
public interface IRxElectionsMetier {
// аутентифікація
Observable<Void> authenticate(User user);
// отримати списки кандидатів
Observable<ListeElectorale[]> getListesElectorales(User user);
// кількість місць, що підлягають заповненню
Observable<Integer> getNbSiegesAPourvoir(User user);
// виборчий поріг
Observable<Double> getSeuilElectoral(User user);
// реєстрація результатів
Observable<Void> recordResultats(User user, ListeElectorale[] listesElectorales);
// розподіл місць
Observable<ListeElectorale[]> calculerSieges(User user, ListeElectorale[] listesElectorales);
}
Інтерфейс [IRxElectionsMetier] переймає методи інтерфейсу [IElectionsMetier], але там, де метод M інтерфейсу [IElectionsMetier] повертав результат типу T, метод M інтерфейсу [IRxElectionsMetier] повертає результат типу Observable<T>. Тип [Observable] надається бібліотекою RxJava. Тип Observable<T> надає метод [subscribe], який асинхронно отримує тип T. З цим методом пов’язані три події:
- onSuccess(T result), яке повідомляє про наявність результату типу T. Асинхронна операція може повернути кілька результатів;
- onError(Throwable th), що повідомляє про виникнення помилки під час асинхронної операції;
- onCompleted(), що повідомляє про завершення асинхронної операції;
Доки метод [Observable.subscribe] не викликано, асинхронна операція, пов’язана з об’єктом Observable, не запускається. Код, який викликає метод M інтерфейсу [IRxElectionsMetier], отримує не очікуваний результат T, а тип Observable<T>, який згодом дозволить йому отримати результат T шляхом виклику методу [Observable.subscribe].
Реалізація [RxElectionsMetier] інтерфейсу [IRxElectionsMetier] виглядає наступним чином:
package elections.security.client.metier;
import elections.security.client.entities.ListeElectorale;
import elections.security.client.entities.User;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import rx.Observable;
@Component
public class RxElectionsMetier implements IRxElectionsMetier {
@Autowired
private IElectionsMetier metier;
@Override
public Observable<Void> authenticate(User user) {
...
}
@Override
public Observable<ListeElectorale[]> getListesElectorales(User user) {
return Observable.create(subscriber -> {
try {
// виклик синхронним методом, а потім відповідь абоненту
subscriber.onNext(metier.getListesElectorales(user));
// повідомляється про завершення спостережуваного процесу
subscriber.onCompleted();
} catch (Exception e) {
// передача винятку
subscriber.onError(e);
}
});
}
@Override
public Observable<Integer> getNbSiegesAPourvoir(User user) {
...
}
@Override
public Observable<Double> getSeuilElectoral(User user) {
...
}
@Override
public Observable<Void> recordResultats(User user, ListeElectorale[] listesElectorales) {
...
}
@Override
public Observable<ListeElectorale[]> calculerSieges(User user, ListeElectorale[] listesElectorales) {
...
}
}
- рядки 12–13: ін'єкція Spring у синхронний бізнес-шар;
- рядки 20–34: ми прокоментуємо метод [getListesElectorales], який замість того, щоб повертати тип [ListeElectorale[]], повертає тип [Observable<ListeElectorale[]>];
- рядки 22–32: статичний метод [Observable.create] дозволяє створити Observable на основі типу [Subscriber]. Тип [Subscriber] представляє підписника потоків результатів, що генеруються спостережуваним процесом (Observable). Він надає три методи:
- [Subscriber.onNext] (рядок 25) для отримання результату від спостережуваного процесу;
- [Subscriber.onError] (рядок 30) — для отримання винятку від спостережуваного процесу. Після виникнення винятку тип [Observable] більше не надсилає результатів;
- [Subscriber.onCompleted] (рядок 27) для отримання сигналу про завершення передачі даних від спостережуваного процесу. У цьому випадку спостережуваний процес передає лише один елемент. Слід зауважити, що цей сигнал не надсилається, якщо відбувається виняток. Це стандартна поведінка Observables: виникнення винятку також сигналізує про завершення передачі даних. Підписники знають про це;
- рядки 22–34: метод [Observable.create] приймає як параметр тип [Observable.OnSubscribe]. Цей тип є функціональним інтерфейсом. Це поняття було введено в Java 8 і позначає інтерфейс, що має єдиний метод. У даному випадку єдиний метод інтерфейсу [Observable.OnSubscribe] має такий вигляд:
Для реалізації функціонального інтерфейсу з єдиним методом m(param1, param2, ..., paramn) можна використовувати такий спрощений синтаксис:
Саме це зроблено в рядках 22–34:
- [subscriber] є параметром методу [Observable.OnSubscribe.call];
- рядки 23–32: код, який ми хочемо присвоїти методу [call];
- рядок 25: виборчі списки запитуються синхронно у шару [métier], введеного в рядку 13. Отже, відбуватиметься очікування результату. Коли результат буде отримано, його передають до методу [onNext] підпискувача;
- рядок 28: у разі помилки виняток передається до методу [onError] підпискувача;
- рядок 31: очікується лише один результат. Коли його отримано (виборчі списки або виняток), підписнику повідомляється, що спостережуваний процес завершив видачу результатів;
Слід пам’ятати, що метод [RxElectionsMetier] повертає тип Observable<ListeElectorale[]>, а не сам тип ListeElectorale[]. Код, що викликає цей метод, повинен викликати метод Observable<ListeElectorale[]>.subscribe, щоб код рядків 23–33 було виконано та повернуто списки виборців за допомогою рядка 25.
Код інших методів є аналогічним:
package elections.security.client.metier;
import elections.security.client.entities.ListeElectorale;
import elections.security.client.entities.User;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import rx.Observable;
@Component
public class RxElectionsMetier implements IRxElectionsMetier {
@Autowired
private IElectionsMetier metier;
@Override
public Observable<Void> authenticate(User user) {
return Observable.create(subscriber -> {
try {
// виклик синхронного методу
metier.authenticate(user);
// повідомляється про завершення спостережуваного об’єкта
subscriber.onCompleted();
} catch (Exception e) {
// передається виняток
subscriber.onError(e);
}
});
}
@Override
public Observable<ListeElectorale[]> getListesElectorales(User user) {
return Observable.create(subscriber -> {
try {
// виклик синхронного методу, а потім відповідь підписнику
subscriber.onNext(metier.getListesElectorales(user));
// повідомляється про завершення спостережуваного об’єкта
subscriber.onCompleted();
} catch (Exception e) {
// передається виняток
subscriber.onError(e);
}
});
}
@Override
public Observable<Integer> getNbSiegesAPourvoir(User user) {
return Observable.create(subscriber -> {
try {
// виклик синхронного методу, а потім відповідь підписнику
subscriber.onNext(metier.getNbSiegesAPourvoir(user));
// повідомляється про завершення спостережуваного об’єкта
subscriber.onCompleted();
} catch (Exception e) {
// передається виняток
subscriber.onError(e);
}
});
}
@Override
public Observable<Double> getSeuilElectoral(User user) {
return Observable.create(subscriber -> {
try {
// виклик синхронного методу, а потім відповідь підписнику
subscriber.onNext(metier.getSeuilElectoral(user));
// повідомляється про завершення спостережуваного об’єкта
subscriber.onCompleted();
} catch (Exception e) {
// передається виняток
subscriber.onError(e);
}
});
}
@Override
public Observable<Void> recordResultats(User user, ListeElectorale[] listesElectorales) {
return Observable.create(subscriber -> {
try {
// виклик синхронного методу
metier.recordResultats(user, listesElectorales);
// повідомляється про завершення спостережуваного об’єкта
subscriber.onCompleted();
} catch (Exception e) {
// передається виняток
subscriber.onError(e);
}
});
}
@Override
public Observable<ListeElectorale[]> calculerSieges(User user, ListeElectorale[] listesElectorales) {
return Observable.create(subscriber -> {
try {
// виклик синхронного методу, а потім відповідь підписнику
subscriber.onNext(metier.calculerSieges(user, listesElectorales));
// повідомляється про завершення спостережуваного об’єкта
subscriber.onCompleted();
} catch (Exception e) {
// передається виняток
subscriber.onError(e);
}
});
}
}
- рядки 20 та 81: метод [onNext] підписника не викликається, оскільки він не очікує результатів;
20.1.4. Тести JUnit шару [métier]
![]() |
20.1.4.1. Test01
Повернемося до модульного тесту [Test01], розглянутого в розділі 17.4.4. Він був розроблений для здійснення синхронних викликів до інтерфейсу [IElectionsMetier]. Ми модифікуємо його так, щоб він здійснював синхронні виклики до нового інтерфейсу [IRxElectionsMetier]. Адже можна здійснювати синхронні виклики до асинхронного інтерфейсу RxJava. Код виглядає наступним чином:
package elections.security.client.metier.junit;
import org.junit.Assert;
import org.junit.BeforeClass;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.SpringApplicationConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import elections.security.client.config.MetierConfig;
import elections.security.client.entities.ElectionsException;
import elections.security.client.entities.ListeElectorale;
import elections.security.client.entities.User;
import elections.security.client.metier.IRxElectionsMetier;
import rx.observables.BlockingObservable;
@SpringApplicationConfiguration(classes = MetierConfig.class)
@RunWith(SpringJUnit4ClassRunner.class)
public class Test01 {
// шар [electionsMetier]
@Autowired
private IRxElectionsMetier electionsMetier;
// мапер jSON
private final ObjectMapper mapper = new ObjectMapper();
// користувачі
static private User admin;
static private User user;
static private User unknown;
@BeforeClass
public static void initTest() {
admin = new User("admin", "admin");
user = new User("user", "user");
unknown = new User("x", "y");
}
@Test()
public void checkUserUser() {
ElectionsException se = null;
try {
BlockingObservable.from(electionsMetier.authenticate(user)).firstOrDefault(null);
} catch (ElectionsException e) {
se = e;
}
Assert.assertNotNull(se);
Assert.assertEquals("403 Forbidden", se.getErreurs().get(0));
}
@Test()
public void checkUserUnknown() {
ElectionsException se = null;
try {
BlockingObservable.from(electionsMetier.authenticate(unknown)).firstOrDefault(null);
} catch (ElectionsException e) {
se = e;
}
Assert.assertNotNull(se);
Assert.assertEquals("401 Unauthorized", se.getErreurs().get(0));
}
@Test()
public void checkUserAdmin() {
ElectionsException se = null;
try {
BlockingObservable.from(electionsMetier.authenticate(admin)).firstOrDefault(null);
} catch (ElectionsException e) {
se = e;
}
Assert.assertNull(se);
}
/**
* vérification 1 : méthode de calcul des sièges on fixe en dur les listes
*/
@Test
public void calculSieges1() {
// створюється таблиця із 7 списків кандидатів
ListeElectorale[] listes = new ListeElectorale[7];
listes[0] = new ListeElectorale("A", 32000, 0, false);
listes[1] = new ListeElectorale("B", 25000, 0, false);
listes[2] = new ListeElectorale("C", 16000, 0, false);
listes[3] = new ListeElectorale("D", 12000, 0, false);
listes[4] = new ListeElectorale("E", 8000, 0, false);
listes[5] = new ListeElectorale("F", 4500, 0, false);
listes[6] = new ListeElectorale("G", 2500, 0, false);
// розраховуються місця для кожного зі списків
listes = BlockingObservable.from(electionsMetier.calculerSieges(admin, listes)).first();
// перевіряються результати
Assert.assertEquals(2, listes[0].getSieges());
Assert.assertFalse(listes[0].isElimine());
Assert.assertEquals(2, listes[1].getSieges());
Assert.assertFalse(listes[1].isElimine());
Assert.assertEquals(1, listes[2].getSieges());
Assert.assertFalse(listes[2].isElimine());
Assert.assertEquals(1, listes[3].getSieges());
Assert.assertFalse(listes[3].isElimine());
Assert.assertEquals(0, listes[4].getSieges());
Assert.assertFalse(listes[4].isElimine());
Assert.assertEquals(0, listes[5].getSieges());
Assert.assertTrue(listes[5].isElimine());
Assert.assertEquals(0, listes[6].getSieges());
Assert.assertTrue(listes[6].isElimine());
}
/**
* vérification 2 : méthode de calcul des sièges on demande les listes à la couche [metier] puis on fixe en dur les
* voix
*/
@Test
public void calculSieges2() {
// створюється таблиця із 7 списків кандидатів
ListeElectorale[] listes = BlockingObservable.from(electionsMetier.getListesElectorales(admin)).first();
// фіксуються голоси
listes[0].setVoix(32000);
listes[1].setVoix(25000);
listes[2].setVoix(16000);
listes[3].setVoix(12000);
listes[4].setVoix(8000);
listes[5].setVoix(4500);
listes[6].setVoix(2500);
// розраховуються місця, отримані кожним зі списків
listes = BlockingObservable.from(electionsMetier.calculerSieges(admin, listes)).first();
// перевіряються результати
Assert.assertEquals(2, listes[0].getSieges());
Assert.assertFalse(listes[0].isElimine());
Assert.assertEquals(2, listes[1].getSieges());
Assert.assertFalse(listes[1].isElimine());
Assert.assertEquals(1, listes[2].getSieges());
Assert.assertFalse(listes[2].isElimine());
Assert.assertEquals(1, listes[3].getSieges());
Assert.assertFalse(listes[3].isElimine());
Assert.assertEquals(0, listes[4].getSieges());
Assert.assertFalse(listes[4].isElimine());
Assert.assertEquals(0, listes[5].getSieges());
Assert.assertTrue(listes[5].isElimine());
Assert.assertEquals(0, listes[6].getSieges());
Assert.assertTrue(listes[6].isElimine());
}
/**
* vérification 3 méthode de calcul des sièges on provoque une exception
*/
@Test(expected = ElectionsException.class)
public void calculSieges3() {
// створюється таблиця з 24 списками кандидатів, кожен із яких має 1 голос
ListeElectorale[] listes = new ListeElectorale[25];
// усі 25 списків матимуть однакову кількість голосів (4%)
for (int i = 0; i < listes.length; i++) {
listes[i] = new ListeElectorale("Liste" + (i + 1), 1, 0, false);
}
// розподіл місць — зазвичай має вийти ElectionsException
// з виборчим бар’єром у 5%
BlockingObservable.from(electionsMetier.calculerSieges(admin, listes)).first();
}
/**
* enregistrement des résultats de l'élection
*
* @throws JsonProcessingException
*/
@Test
public void ecritureResultatsElections() throws JsonProcessingException {
// створюється таблиця із 7 списків кандидатів
ListeElectorale[] listes = BlockingObservable.from(electionsMetier.getListesElectorales(admin)).first();
// встановлюємо фіксовану кількість голосів
listes[0].setVoix(32000);
listes[1].setVoix(25000);
listes[2].setVoix(16000);
listes[3].setVoix(12000);
listes[4].setVoix(8000);
listes[5].setVoix(4500);
listes[6].setVoix(2500);
// розраховуються місця, отримані кожним зі списків
listes = BlockingObservable.from(electionsMetier.calculerSieges(admin, listes)).first();
// виводимо результати
for (int i = 0; i < listes.length; i++) {
System.out.println(mapper.writeValueAsString(listes[i]));
}
// результати записуються в базу даних
BlockingObservable.from(electionsMetier.recordResultats(admin, listes)).firstOrDefault(null);
// перевіряються результати
listes = BlockingObservable.from(electionsMetier.getListesElectorales(admin)).first();
// виводить результати
for (int i = 0; i < listes.length; i++) {
System.out.println(mapper.writeValueAsString(listes[i]));
}
Assert.assertEquals(2, listes[0].getSieges());
Assert.assertFalse(listes[0].isElimine());
Assert.assertEquals(2, listes[1].getSieges());
Assert.assertFalse(listes[1].isElimine());
Assert.assertEquals(1, listes[2].getSieges());
Assert.assertFalse(listes[2].isElimine());
Assert.assertEquals(1, listes[3].getSieges());
Assert.assertFalse(listes[3].isElimine());
Assert.assertEquals(0, listes[4].getSieges());
Assert.assertFalse(listes[4].isElimine());
Assert.assertEquals(0, listes[5].getSieges());
Assert.assertTrue(listes[5].isElimine());
Assert.assertEquals(0, listes[6].getSieges());
Assert.assertTrue(listes[6].isElimine());
}
}
Розглянемо зміни:
- рядок 48: статичний метод [BlockingObservable.from(Observable).first]:
- підписується на спостережуваний параметр [from];
- запускає виконання коду, пов’язаного з об’єктом спостереження;
- очікує на отримання першого результату. Отже, це синхронна операція;
Тут ми використовуємо метод [firstOrDefault(null)], оскільки об’єкт спостереження [metier.authenticate] не повертає результату під час виконання. Отже, результатом методу [firstOrDefault(null)] буде null — значення, яке тут не використовується;
Ми повторюємо цю схему в решті коду щоразу, коли хочемо звернутися до шару [métier].
Юнітарний тест [Test01] повинен пройти успішно:
![]() |
Завдання: перевірити, чи проходить тест [Test01].
20.1.4.2. Test02
Ми модифікуємо тест [Test01], щоб тепер перевіряти асинхронний інтерфейс [IRxElectionsMetier], виконуючи асинхронні виклики його методів.
Розглянемо перший тест:
// семафор синхронізації потоків
private CountDownLatch latch;
// -----------------------------------
private ElectionsException checkUserUserException;
@Test()
public void checkUserUser() throws InterruptedException {
// семафор з значенням 1
latch = new CountDownLatch((1));
// асинхронна операція
electionsMetier.authenticate(user).subscribeOn(Schedulers.io())
.subscribe((result) -> {
},
(th) -> {
checkUserUserException = (ElectionsException) th;
latch.countDown();
},
() -> {
latch.countDown();
});
// очікування семафора
latch.await();
// перевірка результатів
Assert.assertNotNull(checkUserUserException);
Assert.assertEquals("403 Forbidden", checkUserUserException.getErreurs().get(0));
}
- рядок 2: семафор — це інструмент, що використовується для синхронізації потоків між собою. Потоки — це потоки виконання, що виконуються паралельно. Щоб виконати завдання T1, потоку [Thread1] може знадобитися, щоб завдання T2, яке виконується потоком [Thread2], було завершено. Тоді він очікує, поки потік [Thread2] надішле йому сигнал, що вказує на завершення завдання T2. Існують різні способи управління цією синхронізацією між двома потоками. Метод, що використовується тут, є таким:
- рядок 10: потік [Thread1] створює семафор із значенням 1;
- рядок 12: потік [Thread1] створює та запускає потік [Thread2]. Це досягається за допомогою такого синтаксису:
electionsMetier.authenticate(user).subscribeOn(Schedulers.io())
Метод [Observable.subscribeOn] визначає потік, на якому буде виконуватися спостережуваний процес. Параметром методу [subscribeOn] є пул потоків. Бібліотека RxJava надає кілька таких пулів, пристосованих до різних ситуацій. Пул [Schedulers.io()] рекомендується для мережевих операцій;
- (продовження)
- рядки 12–13: операція
electionsMetier.authenticate(user).subscribeOn(Schedulers.io()).subscribe(...)
виконує синхронну операцію, інкапсульовану в об’єкт спостереження [authenticate(user)]. Однак, оскільки ця синхронна операція запускається в іншому потоці, ніж потік [Thread1], останній не чекає на відповідь від методу [subscribe] і переходить до наступної інструкції;
- (продовження)
- рядок 23: потік [Thread1] зупиняється й очікує, поки семафор не прийме значення 0 (наразі його значення — 1);
- рядки 13–21: метод [subscribe] приймає як параметри три лямбда-функції:
- перша, [(result)->{...}], викликається щоразу, коли спостережувана величина [authenticate(user)] видає результат [result]. Тут ми маємо спостережувану величину [authenticate(user)], яка щось робить, але не видає жодного результату. Отже, лямбда-функція [(result)->{}] ніколи не буде викликана. Саме тому її код тут порожній — [{}];
- друга — [(th)->{...}] — отримує як параметр тип [Throwable]. Вона викликається, коли під час виконання обсервабеля виникає виняток. Тут ми обробляємо параметр [Throwable th] наступним чином:
- рядок 16: ми зберігаємо його в полі тестового класу типу [ElectionsException], оскільки виконувана спостережувана величина генерує лише цей тип винятку;
- рядок 17: ми встановлюємо семафор у значення 0, щоб вказати, що потік [Thread2] завершив свою роботу;
- третій [()->{...}] викликається, коли обсервабель більше не має елементів для виведення. Ми обробляємо цю подію наступним чином:
- рядок 20: ми встановлюємо семафор на 0, щоб вказати, що потік [Thread2] завершив свою роботу;
Слід зауважити, що третя лямбда-функція не викликається, якщо виникає виняток. Саме тому ми були змушені також встановити семафор на 0 у рядку 17;
- рядок 25: коли ми доходимо до цього рядка, обсервабель завершив свою роботу. Тоді можна виконати ті самі перевірки, що й у тесті [Test01];
Розглянемо інший тест:
// -----------------------------------
private ElectionsException calculSieges1Exception;
private ListeElectorale[] listesCalculSieges1;
@Test
public void calculSieges1() throws InterruptedException {
// створюється масив із 7 списків-кандидатів
ListeElectorale[] listes = new ListeElectorale[7];
listes[0] = new ListeElectorale("A", 32000, 0, false);
listes[1] = new ListeElectorale("B", 25000, 0, false);
listes[2] = new ListeElectorale("C", 16000, 0, false);
listes[3] = new ListeElectorale("D", 12000, 0, false);
listes[4] = new ListeElectorale("E", 8000, 0, false);
listes[5] = new ListeElectorale("F", 4500, 0, false);
listes[6] = new ListeElectorale("G", 2500, 0, false);
// семафор на 1
latch = new CountDownLatch((1));
// асинхронна операція
// розраховуємо кількість місць для кожного зі списків
electionsMetier.calculerSieges(admin, listes).subscribeOn(Schedulers.io())
.subscribe((result) -> {
listesCalculSieges1 = result;
},
(th) -> {
calculSieges1Exception = (ElectionsException) th;
latch.countDown();
},
() -> {
latch.countDown();
});
// очікування семафора
latch.await();
// перевірка результатів
Assert.assertNull(calculSieges1Exception);
Assert.assertEquals(2, listesCalculSieges1[0].getSieges());
Assert.assertFalse(listesCalculSieges1[0].isElimine());
Assert.assertEquals(2, listesCalculSieges1[1].getSieges());
Assert.assertFalse(listesCalculSieges1[1].isElimine());
Assert.assertEquals(1, listesCalculSieges1[2].getSieges());
Assert.assertFalse(listesCalculSieges1[2].isElimine());
Assert.assertEquals(1, listesCalculSieges1[3].getSieges());
Assert.assertFalse(listesCalculSieges1[3].isElimine());
Assert.assertEquals(0, listesCalculSieges1[4].getSieges());
Assert.assertFalse(listesCalculSieges1[4].isElimine());
Assert.assertEquals(0, listesCalculSieges1[5].getSieges());
Assert.assertTrue(listesCalculSieges1[5].isElimine());
Assert.assertEquals(0, listesCalculSieges1[6].getSieges());
Assert.assertTrue(listesCalculSieges1[6].isElimine());
}
- рядки 20–30: асинхронне виконання обсервабеля [electionsMetier.calculerSieges(admin, listes)];
- рядки 21–23: виконання обсервабеля повертає тип [ListeElectorale[]], який зберігається у полі класу тесту, рядок 3;
- рядки 34–48: ці перевірки відповідають тесту [Test01], до яких додано перевірку в рядку 34, що гарантує відсутність винятків;
Повний текст тесту [Test02] доступний у навчальних матеріалах.
Завдання: запустити тест [Test02] і переконатися, що він пройшов успішно.
20.1.4.3. Test03
Тест [Test03] виконує те саме, що й тест [Test01]: він перевіряє інтерфейс [IRxElectionsMetier] за допомогою синхронних викликів цього інтерфейсу. Це копія тесту [Test02] з двома відмінностями:
- спостережувані об’єкти більше не виконуються в потоці, відмінному від того, що виконує тести. Коли потік [Thread1] виконує метод [subscribe] об’єкта спостереження, цей метод запускає операцію HTTP до сервера, також у потоці [Thread1]. Тоді весь метод [subscribe] стає синхронним;
- оскільки залишається лише один потік, синхронізація потоків стає непотрібною, і семафор зникає;
Ось два приклади тестів:
// -----------------------------------
private ElectionsException checkUserUserException;
@Test()
public void checkUserUser() throws InterruptedException {
// синхронна операція
electionsMetier.authenticate(user)
.subscribe((result) -> {
},
(th) -> {
checkUserUserException = (ElectionsException) th;
},
() -> {
});
// перевірка результатів
Assert.assertNotNull(checkUserUserException);
Assert.assertEquals("403 Forbidden", checkUserUserException.getErreurs().get(0));
}
- рядок 7: за замовчуванням метод [electionsMetier.authenticate(user).subscribe] виконується у потоці коду, що його викликає. Отже, ми маємо синхронну операцію;
// -----------------------------------
private ElectionsException calculSieges1Exception;
private ListeElectorale[] listesCalculSieges1;
@Test
public void calculSieges1() throws InterruptedException {
// створюється таблиця із 7 списків кандидатів
ListeElectorale[] listes = new ListeElectorale[7];
listes[0] = new ListeElectorale("A", 32000, 0, false);
listes[1] = new ListeElectorale("B", 25000, 0, false);
listes[2] = new ListeElectorale("C", 16000, 0, false);
listes[3] = new ListeElectorale("D", 12000, 0, false);
listes[4] = new ListeElectorale("E", 8000, 0, false);
listes[5] = new ListeElectorale("F", 4500, 0, false);
listes[6] = new ListeElectorale("G", 2500, 0, false);
// синхронна операція
// розраховуються місця для кожного зі списків
electionsMetier.calculerSieges(admin, listes)
.subscribe((result) -> {
listesCalculSieges1 = result;
},
(th) -> {
calculSieges1Exception = (ElectionsException) th;
},
() -> {
});
// перевірка результатів
Assert.assertNull(calculSieges1Exception);
Assert.assertEquals(2, listesCalculSieges1[0].getSieges());
Assert.assertFalse(listesCalculSieges1[0].isElimine());
Assert.assertEquals(2, listesCalculSieges1[1].getSieges());
Assert.assertFalse(listesCalculSieges1[1].isElimine());
Assert.assertEquals(1, listesCalculSieges1[2].getSieges());
Assert.assertFalse(listesCalculSieges1[2].isElimine());
Assert.assertEquals(1, listesCalculSieges1[3].getSieges());
Assert.assertFalse(listesCalculSieges1[3].isElimine());
Assert.assertEquals(0, listesCalculSieges1[4].getSieges());
Assert.assertFalse(listesCalculSieges1[4].isElimine());
Assert.assertEquals(0, listesCalculSieges1[5].getSieges());
Assert.assertTrue(listesCalculSieges1[5].isElimine());
Assert.assertEquals(0, listesCalculSieges1[6].getSieges());
Assert.assertTrue(listesCalculSieges1[6].isElimine());
}
Завдання: виконати тест [Test03] і перевірити, чи він проходить успішно.
20.2. крок 2
Тепер ми переносимо синхронну консольну програму з розділу 17.5 у програму, яка також залишається синхронною, але використовує асинхронний інтерфейс [RxJava, metier, DAO];
![]() |
Ми беремо за основу проєкт [elections-console-metier-dao-security-webjson] [1] з розділу 17.5, який дублюємо в новому проєкті [elections-console-rxjava- metier-dao-security-webjson] [2]:
![]() | ![]() |
- у [3-4]; у новому проєкті видаляємо залежність від старого синхронного шару [métier];
![]() | ![]() |
- у [5-9] додається залежність від нового асинхронного шару [métier];
![]() | ![]() |
- у [10-14] перейменовуємо клас [ElectionsConsole] на [ElectionsConsole01];
Аналогічно, клас [BootElectionsConsole] перейменовується на [BootElectionsConsole01]:
![]() |
Поточний код класу [BootElectionsConsole01] такий:
package elections.security.client.boot;
import elections.security.client.console.IElectionsUI;
public class BootElectionsConsole01 extends AbstractBootElections{
public static void main(String[] arguments) {
new BootElectionsConsole01().run();
}
@Override
protected IElectionsUI getUI() {
return ctx.getBean("electionsConsole",IElectionsUI.class);
}
}
- рядок 13: оскільки ми змінили назву класу [ElectionsConsole] на [ElectionsConsole01], тепер потрібно написати:
return ctx.getBean("electionsConsole01",IElectionsUI.class);
Повернемося до коду класу [ElectionsConsole01]:
@Component
public class ElectionsConsole01 implements IElectionsUI {
@Autowired
private IElectionsMetier electionsMetier;
@Autowired
private User admin;
@Override
public void run() {
// списки, що беруть участь у виборах
ListeElectorale[] listes;
// введення даних
try (Scanner clavier = new Scanner(System.in)) {
// запитуються списки, що беруть участь у виборах, у шарі [metier]
listes = electionsMetier.getListesElectorales(admin);
...
// проводиться підрахунок місць
listes=electionsMetier.calculerSieges(admin,listes);
// результати зберігаються
electionsMetier.recordResultats(admin,listes);
...
}
Якщо слідувати прикладу тесту [Test01] з параграфа 20.1.4.1, рядки 5, 17, 20 та 22 зміняться наступним чином:
@Component
public class ElectionsConsole01 implements IElectionsUI {
@Autowired
private IRxElectionsMetier electionsMetier;
@Autowired
private User admin;
@Override
public void run() {
// списки кандидатів
ListeElectorale[] listes;
// введення даних
try (Scanner clavier = new Scanner(System.in)) {
// запит списків, що беруть участь у виборах, до шару [metier]
listes = BlockingObservable.from(electionsMetier.getListesElectorales(admin)).first();
...
// проводиться підрахунок місць
listes = BlockingObservable.from(electionsMetier.calculerSieges(admin, listes)).first();
// результати зберігаються
BlockingObservable.from(electionsMetier.recordResultats(admin, listes));
...
}
Завдання: налаштуйте проект для виконання класу [BootElectionsConsole01] із трьома параметрами [SS, Heures travaillées, Jours travaillés] та переконайтеся, що виконання налаштованого таким чином проекту дає очікувані результати.
Завдання: налаштуйте проект для виконання пари [BootElectionsConsole02, ElectionsConsole02], де клас [ElectionsConsole02] буде написано за зразком тесту [Test02] з параграфа 20.1.4.2.
Завдання: налаштуйте проект для виконання пари [BootElectionsConsole03, ElectionsConsole03], де клас [ElectionsConsole03] має бути написаний за зразком тесту [Test03] з параграфа 20.1.4.3.
20.3. Етап 3
Тепер перейдемо до перенесення Swing-додатку в асинхронне середовище.
![]() |
Почнемо з дублювання проєкту [elections-swing-metier-dao-security-webjson] [1] з розділу 17.6 у новий проєкт [elections-swing-rxjava-metier-dao-security-webjson] [2]:
![]() | ![]() |
- у [3, 4] ми усуваємо залежність від синхронного шару [console];
![]() | ![]() |
- у [5-9] ми додаємо залежність від асинхронного рівня консолі;
Рівень [swing] буде здійснювати справжні асинхронні виклики до рівня [métier]. Під час виклику методу цього рівня будуть задіяні два потоки:
- потік UI, який обробляє події;
- потік вводу-виводу, який виконає виклик HTTP до сервера;
Протягом усього часу асинхронного виклику слід відображати зображення очікування та кнопку скасування. Ми цього не робитимемо в даному випадку, і це буде запропоновано вам як вдосконалення додатка. Зміни відбуваються в обох класах, які здійснюють виклики до шару [métier]:
![]() |
20.3.1. Налаштування Maven
Тут ми будемо використовувати бібліотеку [RxSwing], яка доповнює бібліотеку [RxJava] функціоналом, доступним лише в середовищі Swing. Для цього ми змінюємо файл [pom.xml] наступним чином:
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>istia.st.elections</groupId>
<artifactId>elections-swing-rxjava-metier-dao-security-webjson</artifactId>
<version>0.0.1-SNAPSHOT</version>
<name>elections-swing-rxjava-metier-dao-security-webjson</name>
<description>couche swing asynchrone du client web / jSON</description>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<java.version>1.8</java.version>
<maven.compiler.source>1.8</maven.compiler.source>
<maven.compiler.target>1.8</maven.compiler.target>
</properties>
<dependencies>
<!-- RxSwing -->
<!-- https://mvnrepository.com/artifact/io.reactivex/rxswing -->
<dependency>
<groupId>io.reactivex</groupId>
<artifactId>rxswing</artifactId>
<version>0.27.0</version>
</dependency>
<!-- нижні шари -->
<dependency>
<groupId>istia.st.elections</groupId>
<artifactId>elections-console-rxjava-metier-dao-security-webjson</artifactId>
<version>0.0.1-SNAPSHOT</version>
</dependency>
</dependencies>
</project>
20.3.2. Клас [ElectionsConnectForm]
У режимі асинхронної роботи клас [ElectionsConnectForm] набуває такого вигляду:
package elections.security.client.swing;
import elections.security.client.console.IElectionsUI;
import elections.security.client.entities.User;
import java.awt.Dimension;
import java.awt.Toolkit;
import javax.swing.SwingUtilities;
import elections.security.client.metier.IRxElectionsMetier;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import rx.schedulers.Schedulers;
import rx.schedulers.SwingScheduler;
@Component
public class ElectionsConnectForm extends AbstractElectionsConnectForm implements IElectionsUI {
private static final long serialVersionUID = 1L;
// посилання на асинхронний рівень [métier]
@Autowired
private IRxElectionsMetier metier;
// зареєстрований користувач
private User user;
// головна форма
@Autowired
private ElectionsMainForm electionsMainForm;
// сесія UI
@Autowired
private UiSession uiSession;
@Override
protected void doConnect() {
if (isPageValid()) {
// аутентифікація користувача
metier.authenticate(user).subscribeOn(Schedulers.io()).observeOn(SwingScheduler.getInstance()).subscribe(
// відповіді немає
(result) -> {
},
// обробка винятку
(th) -> {
// помилка зафіксована
String info = getInfoForException("Les erreurs suivantes se sont produites :", th);
// відображення інформації
jTextPaneErreurs.setText(info);
jTextPaneErreurs.setCaretPosition(0);
},
// аутентифікація завершена
() -> {
// користувач зберігається в сесії
uiSession.setUser(user);
// сторінка входу прихована
setVisible(false);
// відображається головна сторінка
electionsMainForm.run();
});
}
}
// ініціалізація
@Override
protected void init() {
...
}
@Override
public void run() {
// відображається графічний інтерфейс
SwingUtilities.invokeLater(new Runnable() {
public void run() {
init();
setVisible(true);
}
});
}
private boolean isPageValid() {
...
}
private String getInfoForException(String message, Throwable ex) {
...
}
}
- рядки 36–63: метод [doConnect] виконується, коли користувач натискає пункт меню [Connexion]:
![]() |
Все вирішується у рядку 40:
metier.authenticate(user).subscribeOn(Schedulers.io()).observeOn(SwingScheduler.getInstance()).subscribe(...)
- спостережуваний процес — [metier.authenticate(user)];
- він буде виконуватися у потоці вводу-виводу, взятому з пулу [Schedulers.io()];
- його спостереження відбуватиметься у потоці UI, який обробляє події інтерфейсу Swing [observeOn(SwingScheduler.getInstance())]. Цей потік отримується за допомогою методу [SwingScheduler.getInstance()], де [SwingScheduler] — це клас, наданий бібліотекою [RxSwing]. Це є обов’язковим. Після отримання результату асинхронної операції його часто використовують для зміни елементів інтерфейсу Swing. Однак цей інтерфейс можна змінювати лише у потоці UI, інакше виникає виняток. Отже, рядки 41–61 мають виконуватися у потоці UI. Це забезпечується методом [observeOn(SwingScheduler.getInstance())];
Прокоментуємо решту коду:
- рядки 42–43: ці рядки потрібні для дотримання синтаксису методу [subscribe]. Вони ніколи не будуть виконані, оскільки процес [metier.authenticate(user)] не повертає жодного результату;
- рядки 35–52: при отриманні винятку його відображають;
- рядки 54–61: виконуються, коли процес [metier.authenticate(user)] повідомляє про завершення передачі даних;
20.3.3. Клас [ElectionsMainForm]
![]() |
20.3.3.1. Ініціалізація графічного інтерфейсу
package elections.security.client.swing;
import elections.security.client.console.IElectionsUI;
import elections.security.client.entities.ListeElectorale;
import elections.security.client.entities.User;
import elections.security.client.metier.IRxElectionsMetier;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import rx.schedulers.Schedulers;
import rx.schedulers.SwingScheduler;
import javax.swing.*;
import java.awt.*;
import java.util.ArrayList;
import java.util.List;
@Component
public class ElectionsMainForm extends AbstractElectionsMainForm implements IElectionsUI {
private static final long serialVersionUID = 1L;
// посилання на асинхронний рівень [métier]
@Autowired
private IRxElectionsMetier metier;
// сесія UI
@Autowired
private UiSession uiSession;
// користувач увійшов у систему
private User user;
// шаблони списків JList
private DefaultListModel<String> modèleNomsVoix = null;
private DefaultListModel<String> modèleRésultats = null;
// списки, що беруть участь у змаганні
private ListeElectorale[] listes;
// списки, введені користувачем
private final List<ListeElectorale> listesSaisies = new ArrayList<>();
private ListeElectorale[] tListesSaisies;
// ініціалізації
@Override
protected void init() {
// генерація компонентів батьківським класом
super.init();
// стан форми
Utilitaires.setEnabled(new JLabel[]{jLabelAjouter, jLabelCalculer, jLabelEnregistrer, jLabelSupprimer}, false);
Utilitaires.setEnabled(
new JMenuItem[]{jMenuItemAjouter, jMenuItemCalculer, jMenuItemEnregistrer, jMenuItemSupprimer}, false);
// вирівнювання вікна по центру
Dimension screenSize = Toolkit.getDefaultToolkit().getScreenSize();
Dimension frameSize = getSize();
if (frameSize.height > screenSize.height) {
frameSize.height = screenSize.height;
}
if (frameSize.width > screenSize.width) {
frameSize.width = screenSize.width;
}
setLocation((screenSize.width - frameSize.width) / 2, (screenSize.height - frameSize.height) / 2);
// користувач увійшов у систему
user = uiSession.getUser();
// локальні ініціалізації
modèleNomsVoix = new DefaultListModel<>();
jListNomsVoix.setModel(modèleNomsVoix);
modèleRésultats = new DefaultListModel<>();
jListResultats.setModel(modèleRésultats);
// запит списків до рівня [métier]
metier.getListesElectorales(user).subscribeOn(Schedulers.io()).observeOn(SwingScheduler.getInstance())
.subscribe(
// відповідь
listesElectorales -> {
// списки зберігаються
listes = listesElectorales;
},
// виняток
(th) -> showException(th),
// кінець спостереження
() -> {
// наступний крок
doInitStep2();
});
}
...
- рядок 46: метод [init] виконується, коли має з’явитися відповідне вікно. Його метою є ініціалізація наведених нижче компонентів [1-3]:
![]() |
- рядки 71–85: асинхронно запитуються списки кандидатів (компонент [1]);
- рядок 71: спостережуваний процес — [metier.getListesElectorales(user)]. Він виконується у потоці вводу-виводу [subscribeOn(Schedulers.io())] і спостережується у потоці UI [observeOn(SwingScheduler.getInstance()];
- рядки 74–77: результат, повернений спостережуваним процесом, зберігається у полі [listes] у рядку 38;
- рядок 79: можливе виключення обробляється за допомогою наступного методу:
private void showException(Throwable th) {
// виводиться виняток
jTextPaneMessages.setText(getInfoForException("Les erreurs suivantes se sont produites : ", th));
jTextPaneMessages.setCaretPosition(0);
}
- рядки 81–84: після завершення спостережуваного процесу виконуються рядки 81–84. Ці рядки не виконуються, якщо стався виняток. Метод [doInitStep2] забезпечує виконання етапу 2 ініціалізації наступним чином:
private void doInitStep2() {
// прив'язування назв списків до комбінованого списку jComboBoxNomsListes
for (int i = 0; i < listes.length; i++) {
jComboBoxNomsListes.addItem(String.format("%s - %s", listes[i].getId(), listes[i].getNom()));
}
// кількість вакантних місць
metier.getNbSiegesAPourvoir(user).subscribeOn(Schedulers.io()).observeOn(SwingScheduler.getInstance())
.subscribe(
// відповідь
nbSiegesAPourvoir -> {
// ініціалізується мітка, пов’язана з цією інформацією
jLabelSAP.setText(jLabelSAP.getText() + nbSiegesAPourvoir);
},
// виняток
(th) -> showException(th),
// кінець спостереження
() -> {
// наступний етап
doInitStep3();
});
}
- рядки 3–5: використовується результат попереднього етапу для заповнення списку, що розгортається, назвами списків кандидатів;
- рядки 7–20: асинхронно запитується кількість вакантних місць;
- рядок 7: спостережуваний процес — [metier.getNbSiegesAPourvoir(user)]. Він виконується у потоці вводу-виводу [subscribeOn(Schedulers.io())] і спостережується у потоці UI [observeOn(SwingScheduler.getInstance()];
- рядки 10–13: результат, повернутий процесом, використовується для оновлення графічного інтерфейсу;
- рядок 15: відображається можливе виключення;
- рядки 17–20: після отримання сигналу про завершення спостережуваного об’єкта відбувається перехід до етапу 3 процесу ініціалізації;
Етап 3 ініціалізації забезпечується наступним кодом:
private void doInitStep3() {
// виборчий поріг
metier.getSeuilElectoral(user).subscribeOn(Schedulers.io()).observeOn(SwingScheduler.getInstance())
.subscribe(
// відповідь
seuilElectoral -> {
// ініціалізується мітка, пов'язана з цією інформацією
jLabelSE.setText(jLabelSE.getText() + seuilElectoral);
},
// виняток
(th) -> showException(th),
// кінець спостереження
() -> {
});
}
- рядки 3–4: асинхронно запитується виборчий поріг;
- рядок 3: спостережуваний процес — [metier.getSeuilElectoral(user)]. Він виконується у потоці вводу-виводу [subscribeOn(Schedulers.io())] і спостерігається у потоці UI [observeOn(SwingScheduler.getInstance()];
- рядки 6–9: результат, повернутий процесом, використовується для оновлення графічного інтерфейсу;
- рядок 11: відображається можливе виключення;
- рядки 13–14: після отримання сигналу про завершення спостережуваного об’єкта ніяких дій не виконується: процес ініціалізації графічного інтерфейсу завершено;
20.3.3.2. Розрахунок кількості місць, отриманих різними списками
Метод [doCalculer] призначений для обчислення кількості місць, отриманих різними списками:
@Override
protected void doCalculer() {
tListesSaisies = listesSaisies.toArray(new ListeElectorale[0]);
// розрахунок місць
String info = null;
metier.calculerSieges(user, tListesSaisies).subscribeOn(Schedulers.io()).observeOn(SwingScheduler.getInstance())
.subscribe(
// обробка результату
result -> consumeResultSieges(result),
// обробка винятку
th -> showException(th),
// кінцева точка
() -> {
}
);
}
- рядки 6–15: асинхронно обчислюються місця, отримані різними списками;
- рядок 6: спостережуваний процес — [metier.calculerSieges(user, tListesSaisies)]. Він виконується у потоці вводу-виводу [subscribeOn(Schedulers.io())] і спостережується у потоці UI [observeOn(SwingScheduler.getInstance()];
- рядок 9: результат, повернений процесом, використовується методом [consumeResultSieges];
- рядок 11: відображається можливе виключення;
- рядки 13–14: після отримання сигналу про завершення роботи об’єкта спостереження ніяких дій не виконується;
У рядку 9 метод [consumeResultSieges] обробляє результат, повернений спостережуваним процесом, списки кандидатів із оновленими полями [sieges, elimine]:
private void consumeResultSieges(ListeElectorale[] tListesSaisies) {
// зберігання результату
this.tListesSaisies = tListesSaisies;
// виведення результатів
modèleRésultats.clear();
for (int i = 0; i < tListesSaisies.length; i++) {
modèleRésultats.addElement(tListesSaisies[i].toString());
}
// оновлення стану форми
Utilitaires.setEnabled(new JLabel[]{jLabelEnregistrer}, true);
Utilitaires.setEnabled(new JLabel[]{jLabelCalculer}, false);
Utilitaires.setEnabled(new JMenuItem[]{jMenuItemEnregistrer}, true);
Utilitaires.setEnabled(new JMenuItem[]{jMenuItemCalculer}, false);
jTextPaneMessages.setText("Calcul terminé");
}
- рядки 4–14: отриманий результат використовується для оновлення графічного інтерфейсу;
20.3.3.3. Реєстрація результатів виборів
Запис результатів виборів здійснюється за допомогою наступного методу [doEnregistrer]:
@Override
protected void doEnregistrer() {
// запит на запис до шару [métier]
metier.recordResultats(user, tListesSaisies).subscribeOn(Schedulers.io()).observeOn(SwingScheduler.getInstance())
.subscribe(
// обробка результату — тут його немає
(param) -> {
},
// обробка винятку
(th) -> showException(th),
// кінцевий результат
() -> {
// оновлення форми
Utilitaires.setEnabled(new JLabel[]{jLabelEnregistrer}, false);
Utilitaires.setEnabled(new JMenuItem[]{jMenuItemEnregistrer}, false);
jTextPaneMessages.setText("Enregistrement des résultats réalisé");
}
);
}
- рядки 4–17: результати виборів записуються асинхронно;
- рядок 4: спостережуваний процес — [metier.recordResultats(user, tListesSaisies)]. Він виконується у потоці вводу-виводу [subscribeOn(Schedulers.io())] і спостерігається у потоці UI [observeOn(SwingScheduler.getInstance()];
- рядки 7–8: ці рядки ніколи не будуть виконані, оскільки спостережуваний процес не повертає результату;
- рядок 10: відображається можливе виключення;
- рядки 14–16: після отримання сигналу про завершення спостережуваного об’єкта оновлюється графічний інтерфейс;
Завдання: переконайтеся, що додаток Swing працює. Потім вдоскональте графічний інтерфейс та код так, щоб під час асинхронної операції з веб-сервером / jSON з’являлося зображення очікування, а також опція скасування поточної операції.























