【发布时间】: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 事件后,发射器无法发射更多事件。
所以我对此有两个想法-
- 也许我没有使用适合我的用例的反应式编程?
- 也许我需要用
rxrelay和PublishRealy来实现它?如果是这样,我该如何保存 emitter 并执行更多逻辑而不是accept(value)?
【问题讨论】:
标签: java android rx-java reactive-programming rx-android