【问题标题】:RxJava error handling for hot observable热可观察的 RxJava 错误处理
【发布时间】:2017-04-11 07:40:36
【问题描述】:

我是 RxJava 的新手,对模式等有一些疑问。 我正在使用下面的代码创建一个 observable:

    public Observable<Volume> getVolumeObservable(Epic epic) {
        return Observable.create(event -> {
            try {
                listeners.add(streamingAPI.subscribeForChartCandles(epic.getName(), MINUTE, new HandyTableListenerAdapter() {
                    @Override
                    public void onUpdate(int i, String s, UpdateInfo updateInfo) {
                        if (updateInfo.getNewValue(CONS_END).equals(ONE)) {
                            event.onNext(new Volume(Integer.parseInt(updateInfo.getNewValue(LAST_TRADED_VOLUME))));
                        }
                    }
                }));
            } catch (Exception e) {
                LOG.error("Error from volume observable", e);
            }
        });
    }

一切都按预期工作,但我对错误处理有一些疑问。 如果我理解正确,这将被视为“热观察”,即无论是否有订阅,事件都会发生(onUpdate 是我无法控制的远程服务器使用的回调)。

我选择不在这里调用 onError,因为我不希望 observable 在发生单个异常时停止发出事件。有没有更好的模式可以使用? .retry() 浮现在脑海中,但我不确定它对于热可观察是否有意义?

另外,在创建订阅但在调用第一个 onNext 之前,observable 是如何表示的?它只是一个 Observable.empty()

【问题讨论】:

  • 您认为错误来自哪里?来自listeners.add() 还是来自onUpdate()?您能否举例说明您希望通知订阅者的错误情况。
  • 我猜你有点误解 hot/cold Observable 。这并不热,每个订阅者都有自己的监听器来发出事件。甚至您也没有在 dispose 中取消注册您的听众。由于 Observable.create 机制,可观察对象在您处置后不会发出事件。
  • 可以是 listeners.add() 和 onUpdate()。不幸的是,我使用的 API 指定得很差。
  • 谢谢,我现在意识到我的生产者必须在 Observable.create() 之外创建才能被认为是热的

标签: java rx-java reactive-programming rx-java2


【解决方案1】:

1) 你的 observable 不是 hot。区别因素是多个订阅者是否共享同一个订阅。 Observable.create() 为每个订阅者调用订阅函数,即它是

虽然很容易。只需添加share() 运算符。它将订阅第一个订阅者并取消订阅最后一个订阅者。不要忘记用这样的方式实现 unsubscribe 功能:

event.setCancellable(() -> listeners.remove(...));

2) 错误可能是可恢复的不可恢复的

如果您认为错误是可自行恢复的(您无需采取任何措施),则不应调用 onError,因为这会杀死您的 observable(不会发出更多事件)。您可以通过发送带有错误详细信息的特殊 Volume 消息来通知您的订阅者。

如果错误是致命的,例如您未能添加侦听器,因此可能没有更多消息,您不应该默默地忽略这一点。发出 onError 因为你的 observable 无论如何都不起作用。

如果错误需要您采取措施,通常是重试或超时重试,您可以添加 retryXxx() 运算符之一。在create() 之后但share() 之前执行此操作。

3) Observable 是一个带有subscribe() 方法的对象。它的精确表示方式取决于您创建它的方法。以create()的源代码为例。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多