Skip to content

8. Swing环境中的RxJava

8.1. Introduction

在此我们将回顾第2节中介绍的Swing应用程序。

  

要在 Swing 环境中使用 RxJava,我们将使用 RxSwing 库,该库为 RxJava 添加了在 Swing 环境中非常有用的类和接口。为此,Swing 示例的 Gradle 文件如下:

  

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'
}
  • 第15行:对RxSwing的依赖;

我们将仅使用一个专属于 RxSwing 的对象:调度器 [SwingScheduler.getInstance()],它负责在 Swing 事件循环线程上执行/监听可观察对象。 我们将专门使用它来监听事件循环线程之外的其他线程上运行的可观察对象。回顾一下示例应用程序的架构:

Image

  • 异步服务层提供了返回可观察对象的方法。我们在与事件循环线程不同的线程中执行这些可观察对象。因此,图形界面不会冻结,能够响应用户的操作。 最直观的示例是允许用户点击按钮 [Annuler] 来中断耗时过长的异步操作。要实现这一点,只需确保图形界面处于静止状态(frozen);
  • Swing 层需要利用异步操作返回的结果,并据此更新图形界面。但这只能在事件循环线程中进行。为此,这些结果会在调度器 [SwingScheduler.getInstance()] 中被监听;

因此,在图形用户界面的事件处理代码中,与异步层 [rxService] 的交互以如下形式进行:


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

其中,调度器 [Schedulers.computation()] 可根据具体使用场景替换为其他调度器。

建议读者重读第2段。现在读者已具备完全理解该段落所需的知识。

8.2. 代码结构

该代码实现了以下架构:

Image

实现该架构的 IntelliJ IDEA 项目如下:

  
  • 包 [rxswing.service] 实现了同步服务层(IService,Service)和异步服务层(IRxService、RxService);
  • 包 [rxswing.ui] 实现了 Swing 接口;

8.3. 项目运行

要在 IntelliJ IDEA 中运行该项目,请按以下步骤操作:

 

8.4. 同步服务

Image

  

同步服务层具有以下接口 [IService]:


package dvp.rxswing.service;

public interface IService {
  // 随机数位于区间 [a,b]
  // 生成 n 个随机数,其中 n 本身也是 [minCount, maxCount] 区间内的随机数
  // 在等待 delay 毫秒后生成这些数,
  // 其中 [delay] 本身是区间 [minDelay, maxDelay] 内的随机数
  public ServiceResponse getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay);
}

服务响应的类型 [ServiceResponse] 如下:


package dvp.rxswing.service;

import java.util.List;

public class ServiceResponse {

  // 服务等待时间
  private int delay;
  // 随机数
  private List<Integer> aleas;
  // 执行线程
  private String executedOn;

  // 构造函数

  public ServiceResponse() {
      // 执行线程
    executedOn = Thread.currentThread().getName();
  }

  public ServiceResponse(int delay, List<Integer> aleas) {
      // 局部构造函数
    this();
    // 其他初始化
    this.delay = delay;
    this.aleas = aleas;
  }

  // 获取器和设置器
...
}

接口 [IService] 由以下类 [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) {
    // 区间内的随机数 [a,b]
    // 生成 n 个随机数,其中 n 本身也是该区间内的随机数 [minCount, maxCount]
    // 在等待 delay 毫秒后生成这些数,
    // 其中 [delay] 本身是区间 [minDelay, maxDelay] 内的随机数

    // 一些验证
    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;
    }
    // 错误?
    if (!messages.isEmpty()) {
      throw new AleasException(String.join(" [---] ", messages), erreur);
    }
    // 随机数生成器
    Random random = new Random();
    // 等待?
    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);
      }
    }
    // 生成结果
    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));
    }
    // 返回结果
    return new ServiceResponse(delay,nombres);
  }

}

服务使用的异常类 [AleasException] 如下:


package dvp.rxswing.service;

public class AleasException extends RuntimeException {

    private static final long serialVersionUID = 1L;
    // 错误代码
  private int code;

  // 构造函数
  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;
  }

  // 获取器和设置器
...
}
  • 第 3 行:它继承了 [RuntimeException] 类。因此这是一个未捕获的异常;
  • 第 7 行:它为其父类添加了一个错误代码(0=无错误);

8.5. 异步服务

Image

  

异步服务层提供以下接口 [IRxService]:


package dvp.rxswing.service;

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

public interface IRxService {
  // 区间内的随机数 [a,b]
  // 生成 n 个数,其中 n 本身是区间内的随机数 [minCount, maxCount]
  // 在等待 delay 毫秒后生成这些数,
  // 其中 [delay] 本身是区间 [minDelay, maxDelay] 内的随机数
  public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse);
}
  • 第 11 行:该服务的方法 [getAleas] 现在返回一个可观察对象;

方法 [getAleas] 返回类型为 [UiResponse] 的响应,该响应面向 [Ui] 层。该类型定义如下:


package dvp.rxswing.ui;

import dvp.rxswing.service.ServiceResponse;

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

public class UiResponse {

  // 客户ID
  private int idClient;
  // 服务响应
  private ServiceResponse serviceResponse;
  // 观察线程名称
  private String observedOn;
  // 请求时间
  private String requestAt;
  // 响应时间
  private String responseAt;

  // 构造函数

  public UiResponse() {
      // 监视线程
    observedOn = Thread.currentThread().getName();
    // 请求时间
    requestAt = getTimeStamp();
  }

  // 私有方法

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

  // getter 和 setter
...
}
  • 随机数位于第13行的字段中;
  • 其余字段用于指定异步服务可观察对象的执行线程和观察线程,以及向服务发起请求和获得响应的时间;

异步接口由以下类 [RxService] 实现:


package dvp.rxswing.service;

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

public class RxService implements IRxService {

  // 同步服务
  private IService service;

  // 构造函数
  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) {
      // 创建一个可观察对象,用于发布同步服务返回的值
    return Observable.create(subscriber -> {
      try {
          // 同步调用
        uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
        // 将结果传递给观察者
        subscriber.onNext(uiResponse);
      } catch (Exception e) {
          // 将错误传递给观察者
        subscriber.onError(e);
      } finally {
          // 通知观察者已停止发布
        subscriber.onCompleted();
      }
    });
  }
}
  • 第12-14行:异步服务的[RxService]类基于同步接口[IService]的实例构建;
  • 第20-33行:构建可观察对象,即方法[getAleas]的返回结果;
  • 第 22 行:调用同步方法 [service.getAleas]。其类型为 [ServiceResponse] 的结果被包含在类型为 [UiResponse] 的对象中,该对象将提供给 [swing] 层。 该对象最初是在方法调用参数中传递的(最后一个参数,第17行);
  • 第24行:响应[UiResponse]被发送给观察者([swing]层)。 对象 [UiResponse] 不仅包含第 22 行同步服务构建的信息,还包含第 17 行方法 [getAleas] 的调用方构建的其他信息。 正因如此,该调用方法将对象 [UiResponse] 作为参数传递给了方法 [getAleas](第 17 行中的最后一个参数);
  • 第30行:别忘了标记数据传输结束。这里有一个仅返回一个值的可观察对象:即同步服务返回的值;
  • 第27行:向观察者报告可能出现的错误;

8.6. 图形用户界面

Image

  
  • 图形界面是使用IDE [Netbeans]构建的,该工具拥有一个优秀的图形编辑器。该编辑器生成了文件[AbstractJFrameAleas.form],该文件仅能由IDE使用;
  • [AbstractJFrameAleas]类同样由NetBeans的图形编辑器生成。 随后对其进行了如下重构:需要处理的图形界面事件由 [AbstractJFrameAleas] 类通过抽象方法进行处理,这些方法在子类 [JFrameAleasEvents] 中实现。最终,
    • 抽象类 [AbstractJFrameAleas] 负责构建和显示图形用户界面;
    • 子类 [JFrameAleasEvents] 负责管理其事件;

[Request] 选项卡的图形用户界面组件如下:

 
编号
类型
名称
作用
1
JTabbedPane
jTabbedPane1
一个标签页容器。包含两个标签页(JPanel)[jPanelRequest]用于请求,[jPanelresponse]用于响应;
2
JTextField
jTextFieldNbValeurs
向随机数服务发送的请求数量。对于在调度器 [Schedulers.io] 上运行的异步服务,这些请求将共享一个处理器;
3
JTextField
jTextFieldA
[a,b]区间的a端点
4
JTextField
jTextFieldB
区间中的端子 b [a,b]
5
JTextField
jTextFieldMinCount
区间 [minCount, maxCount] 中的端子 minCount
6
JTextField
jTextFieldMaxCount
[minCount, maxCount]区间内的maxCount端子
7
JTextField
jTextFieldMinDelay
minDelay 端子,区间 [minDelay, maxDelay]
8
JTextField
jTextFieldMaxDelay
maxDelay 端子,区间 [minDelay, maxDelay]
9
JCheckBox
jCheckBoxRxSwing
如果勾选此框,请求将通过异步接口发送。否则将通过同步接口发送
10
JComboBox
jComboBoxSchedulers
对于异步请求,将使用此处选择的调度程序进行执行
11
JButton
jButtonGenerate
将请求的执行交由同步或异步服务处理

[Response] 选项卡的图形界面组件如下:

 
编号
类型
名称
作用
1
JLabel
jLabelDuree
请求的总执行时间(以毫秒为单位)
2
JLabel
jLabelNbReponses
观察到的响应总数(可能与请求数不同,因为每个请求可能提供多个待观察的值)
3
JList
jListNumbers
显示观察到的(接收到的)值
4
JButton
jButtonAnnuler
取消正在执行的请求

8.7. 实例化图形用户界面

  

类 [JFrameAleasEvents] 管理图形用户界面的事件,特别是对按钮 [Générer] 的点击。这是一个可执行类,在以下上下文中启动:


public class JFrameAleasEvents extends AbstractJFrameAleas {

    private static final long serialVersionUID = 1L;
    // 同步生成服务
    private IService service;
    // 异步生成服务
    private IRxService rxService;

    // 数据录入
    private int nbRequests;
    private int a;
    private int b;
    private int minDelay;
    private int maxDelay;
    private int minCount;
    private int maxCount;

    // 错误消息
    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 ";

    // 可观察对象的订阅
    protected List<Subscription> subscriptions = new ArrayList<Subscription>();
    // 执行开始-结束
    private long debut;
    // 映射器 jSON
    private ObjectMapper jsonMapper;
    // 响应模型
    private DefaultListModel<String> model;

    // 构造器
    public JFrameAleasEvents() {
        // 父级
        super();
        // 本地
        initJFrame();
        // 服务
        service = new Service();
        rxService = new RxService(service);
        // 映射器 jSON
        jsonMapper = new ObjectMapper();
    }

    private void initJFrame() {
        // 隐藏错误消息
        jLabelCountError.setText("");
        jLabelDelayError.setText("");
        jLabelIntervalError.setText("");
        jLabelNbValuesError.setText("");
        // 隐藏默认文本
        jTextFieldA.setText("100");
        jTextFieldB.setText("200");
        jTextFieldMinCount.setText("5");
        jTextFieldMaxCount.setText("10");
        jTextFieldMinDelay.setText("100");
        jTextFieldMaxDelay.setText("500");
        jTextFieldNbValeurs.setText("10");
        jLabelDuree.setText("");
        // 响应模板
        model = new DefaultListModel<>();
        jListNumbers.setModel(model);
        // 核心数
        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);
        }

        /* 创建并显示表单 */
        java.awt.EventQueue.invokeLater(() -> {
            new JFrameAleasEvents().setVisible(true);
        });
    }
  • 第 1 行:类 [JFrameAleasEvents] 继承自类 [AbstractJFrameAleas],而后者又继承自 Swing 类 [JFrame]。因此,类 [JFrameAleasEvents] 是一个 Swing 窗口;
  • 第68-75行:即将执行的[main]方法;
  • 第 70 行:设置图形界面的外观和感觉
  • 第 79 行:调用类 [JFrameAleasEvents] 的构造函数:图形界面将被构建并初始化。完成后,界面被显示出来;
  • 第34-44行:构造函数;
  • 第 36 行:调用父构造函数将初始化图形界面。此时,界面呈现为开发人员设计的样式,但尚未可见;
  • 第 38 行:图形界面的某些组件被初始化;
  • 第 40 行:同步服务的实例化;
  • 第 41 行:实例化异步服务;

8.8. 同步请求的执行

单击按钮 [Générer] 将触发以下方法 [doGenerate] 的执行:


    @Override
    protected void doGenerate() {
        // 输入有效吗?
        if (!isPageValid()) {
            return;
        }
        // 是否接收?
        if (jCheckBoxRxSwing.isSelected()) {
            // 异步请求
            doGenerateWithRxService();
        } else {
            // 同步请求
            doGenerateWithService();
        }
}
  • 第 4-6 行:验证用户输入是否有效。我们不再对方法 [isPageValid] 进行说明,该方法较为基础;
  • 第 8 行:检测复选框 RxSwing 的状态;
  • 第13行:同步执行请求;

方法 [doGenerateWithService] 如下:


    // 同步生成
    private void doGenerateWithService() {
        // 开始等待
        beginWaiting();
        try {
            for (int i = 0; i < nbRequests; i++) {
                // 响应准备
                UiResponse uiResponse = new UiResponse();
                // 客户编号
                uiResponse.setIdClient(i);
                // 同步调用
                uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
                // 响应时间
                uiResponse.setResponseAt();
                // 使用收到的响应更新 JList 模型
                model.add(0, jsonMapper.writeValueAsString(uiResponse));
                // 更新响应数量
                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);
        }
        // 等待结束
        endWaiting();
}
  • 第 12 行:同步调用随机数生成服务;
  • 方法 [doGenerateWithService] 的执行完全在 Swing 的事件循环线程中进行。只要该方法未完成,图形界面就不会处理任何新事件。界面处于冻结状态(frozen)。 因此,例如第16行和第18行的图形界面更新将永远不会被看到。只有在所有请求执行完毕后,它们才会显示其最终值;

方法 [beginWaiting](第 4 行)如下:


    private void beginWaiting() {
        // 按钮
        jButtonGenerate.setVisible(false);
        jButtonCancel.setVisible(true);
        // 等待光标
        jTabbedPane1.setCursor(Cursor.getPredefinedCursor(Cursor.WAIT_CURSOR));
        jButtonCancel.setCursor(Cursor.getDefaultCursor());
        // 清空回复
        model.clear();
        // 接收订阅
        subscriptions.clear();
        // 显示响应视图
        jTabbedPane1.setSelectedIndex(1);
        jLabelNbReponses.setText("0");
        jLabelDuree.setText("");
        // 开始执行
        debut = new Date().getTime();
}
  • 第 3 行:按钮 [Générer] 被隐藏。这会触发一个事件,该事件同样只能在所有请求执行完毕后才能执行。 因此我们从未看到它被隐藏,因为方法 [doGenerateWithService] 的第 25 行中的方法 [endWaiting] 会将其重新显示;
  • 第13行:选择[Response]选项卡以查看响应的到达情况。同样,该事件也仅会在所有请求执行完毕后才被触发,届时将显示全部响应,而我们原本希望看到它们依次到达;

同步接口显然存在不足。这些问题通过异步接口得以解决。

8.9. 异步请求的执行

异步请求的执行代码如下:


private void doGenerateWithRxService() {
        // 开始等待
        beginWaiting();
        // 将以可观察对象的形式获取随机数
        Observable<UiResponse> observable = Observable.empty();
        // 不同可观察对象的执行调度器
        Scheduler[] schedulers = { Schedulers.io(), Schedulers.computation(), Schedulers.newThread(),
                Schedulers.trampoline(), Schedulers.immediate() };
        Scheduler scheduler = schedulers[jComboBoxSchedulers.getSelectedIndex()];
        // 可观察对象的配置
        for (int i = 0; i < nbRequests; i++) {
            // 准备响应
            UiResponse uiResponse = new UiResponse();
            uiResponse.setIdClient(i);
            // 将可观察量配置为在用户选择的调度器上运行
            // 随后将所得观测值累加到总观测值中
            observable = observable.mergeWith(
                    rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(scheduler));
        }
        // 观察者
        observable = observable.observeOn(SwingScheduler.getInstance());
        // 目前仅完成了配置
        // 尚未向同步随机数生成服务发出任何请求
        // 订阅该可观察对象——这将触发对同步随机数生成服务的调用
        try {
            // 这里只有一个订阅——结果是一个订阅
            subscriptions.add(observable.subscribe(
                    // 发布通知
                    uiResponse -> {
                        // 使用响应更新 UI
                        // 这是可能的,因为观察是在 UI 线程中进行的
                        updateUi(uiResponse);
                    } ,
                    // 错误通知
                    th -> {
                        // 出现错误 - 显示错误
                        String message = getInfoForThrowable("L'erreur suivante s'est produite", th);
                        JOptionPane.showMessageDialog(this, message, "Informations", JOptionPane.PLAIN_MESSAGE);
                        // 取消请求
                        doCancel();
                    } ,
                    // 通知 [onCompleted]
                    // 等待结束
                    this::endWaiting));
        } catch (Throwable th) {
            // 异常情况 + 通用 - 显示
            String message = getInfoForThrowable("L'erreur suivante s'est produite", th);
            JOptionPane.showMessageDialog(this, message, "Informations", JOptionPane.PLAIN_MESSAGE);
            // 取消请求
            doCancel();
        }
    }
  • 第3行:修改图形界面以显示正在进行一项可能耗时的操作;
  • 第 5 行:创建一个空的可观察对象。[swing] 层将监听该可观察对象;
  • 第 7 行:可用的调度器列表;
  • 第 9 行:我们允许用户选择用于执行请求的调度器。我们获取用户选择的调度器;
  • 第11-19行:每个查询都会返回一个可观察对象,其元素会被累加(mergeWith)(第17行)到第5行的可观察对象中;
  • 第 13-14 行:构建对象 [UiResponse]。需要说明的是,该对象既是方法 [RxService.getAleas] 的输入参数,也是其返回结果(第 17-18 行);
  • 第14行:每个请求都通过其编号进行标识,此处称为[idClient]。这是必要的,因为在异步环境中,响应的接收顺序可能与请求的发送顺序不同。[idClient]用于确定响应属于哪个请求;
  • 第17-18行:发出异步请求 [rxService.getAleas]。该请求在用户选定的调度器上执行。其结果类型为 Observable<UiResponse>,并将该结果与第5行的可观察对象合并。 必须清楚,此处的 [rxService.getAleas] 方法已执行并返回了一个可观察对象。但这并不意味着已生成随机数。实际上,可观察对象只有在被订阅时才会被执行。目前尚未发生订阅;
  • 第21行:这是关键指令:要求在UI线程上监听第5行可观察对象发出的元素。此处使用了RxSwing库特有的调度器;
  • 第25-51行:订阅第5行的可观察对象。直到此时,才会向同步随机数生成服务请求随机数。关键部分在于第29-33行的指令。 其余部分主要处理错误情况以及可观察对象的[onCompleted]通知;
  • 第28-44行:需注意,我们是在UI线程上订阅了第5行的观察对象。因此第28-44行的代码将在UI线程中执行;
  • 第29-33行:处理可观察对象的[onNext]通知。接收由被观察进程发出的[UiResponse]类型数据,这是某次异步请求的结果。使用该响应更新图形界面;
  • 第 34-41 行:处理可观察对象的 [onError] 通知。显示一个显示错误的对话框(第 37-38 行),然后取消请求(第 40 行);
  • 第 42-44 行:处理可观察对象的 [onCompleted] 通知。更新图形界面以显示所请求的服务已完成。第 44 行也可以如下编写
 ()->{endWaiting();}

此处我们更倾向于使用方法引用;

  • 第45-51行:某些异常不会经过第34-41行。例如,当请求过多时即属此类情况。一旦超过某个阈值(该阈值取决于运行时的环境),就会产生一个[StackOverflowError]异常,该异常会被第45-51行拦截;
  • 第27行:订阅生成类型为[Subscription]的对象,并将其添加到订阅列表中。该列表在此处仅包含一个元素;

第32行,使用以下方法[updateUi]更新图形界面:


    private void updateUi(UiResponse uiResponse) {
        // 响应时间
        uiResponse.setResponseAt();
        // 监视线程
        uiResponse.setObservedOn();
        // 响应数量
        jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
        // 执行时间
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
        // 将响应中的字符串 jSON 添加到响应模板 JList 中
        try {
            model.add(0, jsonMapper.writeValueAsString(uiResponse));
        } catch (JsonProcessingException e) {
            e.printStackTrace();
        }
}

此处可见图形界面的组件已被更新(第7、9、12行)。要实现这一点,必须位于Ui线程(事件循环)中。

方法 [endWaiting] 如下所示:


    private void endWaiting() {
        // 可见按钮 [Générer]
        jButtonGenerate.setVisible(true);
        // 隐藏按钮 [Annuler]
        jButtonCancel.setVisible(false);
        // 等待光标隐藏
        jTabbedPane1.setCursor(Cursor.getDefaultCursor());
        // 选中“回答”选项卡
        jTabbedPane1.setSelectedIndex(1);
        // 最后一次更新时间
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
}

当异步请求执行过程中发生错误,或者用户点击按钮 [Annuler] 时,会调用方法 [doCancel]。其代码如下:


// 可观察对象的订阅
    private List<Subscription> subscriptions = new ArrayList<Subscription>();
....

    @Override
    protected void doCancel() {
        // 等待结束
        endWaiting();
        // 订阅情况
        if (jCheckBoxRxSwing.isSelected() && subscriptions != null) {
            subscriptions.forEach(Subscription::unsubscribe);
            //subscriptions.forEach(s -> s.unsubscribe());
        }
    }

  • 第 2 行:[subscriptions] 是一个订阅列表;
  • 第 11 行:所有订阅均被取消;
  • 第 12 行:第 11 行的另一种写法。此处 [forEach] 方法期待一个 Consumer<Subscription> 类型的实例(参见第 4.4 节);

让我们回到方法 [doGenerateWithService] 的代码:它可分解为两个步骤:

  1. 可观察对象的配置阶段。该操作在方法 [doGenerateWithService] 的调用方线程中进行,即 UI 线程;
  2. 订阅操作,该操作将触发可观察对象的执行;

如果可观察对象的调度器是 [Schedulers.computation(), Scheduler.io(), Schedulers.newThread()] 中的某个调度器,那么它们将在 UI 线程之外执行。这些不同的线程将争夺机器的一个或多个核心。 由于请求属于耗时较长的操作(需数百毫秒),在 UI 线程中执行的 [doGenerateWithService] 方法将在请求返回响应之前结束。 然而,该方法是在按钮 [Générer] 的点击事件上触发的。该事件处理完毕后,UI 线程将能够转而处理后续事件。此类事件有多个。因此,方法 [beginWaiting] 已设置了多个:


    private void beginWaiting() {
        // 按钮
        jButtonGenerate.setVisible(false);
        jButtonCancel.setVisible(true);
        // 加载指示器
        jTabbedPane1.setCursor(Cursor.getPredefinedCursor(Cursor.WAIT_CURSOR));
        jButtonCancel.setCursor(Cursor.getDefaultCursor());
        // 清空回复
        model.clear();
        // Rx 订阅
        subscriptions.clear();
        // 显示答案视图
        jTabbedPane1.setSelectedIndex(1);
        jLabelNbReponses.setText("0");
        jLabelDuree.setText("");
        // 开始执行
        debut = new Date().getTime();
}

该代码中的几乎每一行都会对图形界面产生影响。这些更新并非立即发生:相关事件会被放入事件循环的队列中。一旦处理完对按钮 [Générer] 的点击事件,这些事件就会依次执行,用户便能看到图形界面的变化:

  • 显示 [Response] 选项卡(第 13 行),并为其关联一个等待光标(第 6 行)
  • 其按钮 [Annuler] 显示出来(第 4 行),用户可以点击它;
  • 响应字段 JList 被清空(第 9 行);
  • 响应数量的 JLabel 显示 0;
  • 执行时长的 JLabel 显示为空字符串;

在整个查询执行期间,UI线程会定期获得处理器访问权限。此时,它可处理待处理事件。其中包括由[updateUi]方法触发的事件:


    private void updateUi(UiResponse uiResponse) {
        // 响应时间
        uiResponse.setResponseAt();
        // 观察线程
        uiResponse.setObservedOn();
        // 响应数
        jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
        // 执行时长
        jLabelDuree.setText(String.valueOf(new Date().getTime() - debut));
        // 将响应中的字符串 jSON 添加到响应模板 JList 中
        try {
            model.add(0, jsonMapper.writeValueAsString(uiResponse));
        } catch (JsonProcessingException e) {
            e.printStackTrace();
        }
}

当 UI 线程获得控制权时:

  • 响应数量的 JLabel 会被更新(第 7 行);
  • 更新执行时长的 JLabel(第 9 行);
  • 通过其模板更新响应的 JList(第 12 行);

因此,用户可以查看请求执行的进度。 此外,用户可通过 [Annuler] 按钮取消这些请求。这正是 [swing] 层前端采用异步服务的意义所在,而 RxJava 则是实现这些服务的首选技术。

最后需要指出的是,如果用户选择了 [Schedulers.immediate(), Schedulers.trampoline()] 调度器之一,则可观察对象将在与调用者相同的线程上执行,即 UI 线程。此时系统将回归同步运行模式。

第 2.8.12.8.22.8.32.8.4 节展示了使用不同调度器所获得的结果。