【问题标题】:rxjava ObservableCreate with emitter vs rxrelay PublishRelayrxjava ObservableCreate 与发射器 vs rxrelay PublishRelay
【发布时间】:2021-12-24 16:50:52
【问题描述】:

我对响应式编程有点陌生,所以我想先解释我的用例和我做了什么,然后是我的解决方案的问题以及我认为的答案,但我不确定。

所以这是我的用例: 我有一个尝试连接到远程服务器的 android 应用程序。 当用户单击一个按钮时,我正在尝试使用响应式编程连接到该服务器,并且我尝试连接多个“连接器”,因此我遍历它们,直到其中一个连接器成功。

所以我就是这样做的:

Disposable disposable = currentConnector.connect(ip)
        .observeOn(AndroidSchedulers.mainThread())
        .flatMap(fallbackDefaultConnector())
        .subscribeWith(connectionSubscriber);

这是connectionSubscriber

new DisposableObserver<Boolean>() {
            @Override
            public void onNext(@NonNull Boolean item) {
                Log.d(TAG, "Observable emits: " + item);
                if (item) {
                    Log.d(TAG, "Connected!");
                } else {
                    handleConnectionFailed();
                }
            }

            @Override
            public void onError(@NonNull Throwable e) {
                Log.e(TAG, "On Error" + Log.getStackTraceString(e));
                handleConnectionFailed();
            }

            @Override
            public void onComplete() {
                closeConnectionProcedure();
            }
        };

这是currentConnector.connect(ip) 的示例(请注意超时):

@Override
public Observable<Boolean> connect(String ipAddress) {
    isConnected = false;
    client.setHost(ipAddress);

    Observable<Boolean> observable = client.connect
            .timeout(CONNECTION_TIMEOUT, TimeUnit.SECONDS)
            .subscribeOn(Schedulers.io());
    connectionEmitter = client.getConnectionEmitter();

    return observable;
}

现在是棘手的部分 -> 在这个连接器中,我保存了客户端的 emitter。为什么?例如,因为客户端将我重定向到另一个配对片段。 这就是client.connect客户端的构造函数中的样子:

connect = new ObservableCreate<>(emitter -> {
    connectionEmitter = emitter;
    doConnect(host);
});

doConnect(host) 在客户端成功时,我会发送这样的消息:

connectionEmitter.onNext(true);
connectionEmitter.onComplete();

它有效。当客户端连接时,一切都很好。 我的问题是当 所有连接器上的连接都失败 或用户收到 timeoutexception 时。如果失败,我会发送这样的失败:

connectionEmitter.onNext(false);
connectionEmitter.onComplete();

确实,用户看到“连接失败”消息,没关系。

那么问题出在哪里? 如果失败并且用户再次尝试连接,则不会发生任何事情! 我对其进行了研究,我发现在 observable 发生 onComplete 事件或 onError 事件后,发射器无法发射更多事件

所以我对此有两个想法-

  1. 也许我没有使用适合我的用例的反应式编程?
  2. 也许我需要用rxrelayPublishRealy 来实现它?如果是这样,我该如何保存 emitter 并执行更多逻辑而不是 accept(value)

【问题讨论】:

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


    【解决方案1】:

    保存发射器的设计选择似乎是错误的。一旦发射器完成发射,就应该为 GC 释放它。这就是 Rx 的设计方式,如果您需要重做工作,那么您应该创建一个发射器或可观察的新实例,然后您可以观察这个新实例。

    【讨论】:

    • 即使我在失败时将 observable 设置为 null,它仍然无法再次发射
    • 当你重新连接时你会创建一个新的 observable 吗??
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-06-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多