假设您为 Retrofit 定义的接口包含这样的方法:
public Observable<GameState> loadGameState(@Query("id") String gameId);
可以通过以下三种方式之一定义改造方法:
1.) 一个简单的同步:
public GameState loadGameState(@Query("id") String gameId);
2.) 采用Callback 进行异步处理:
public void loadGameState(@Query("id") String gameId, Callback<GameState> callback);
3.) 和返回 rxjava Observable 的那个,见上文。我认为如果你打算将 Retrofit 与 rxjava 结合使用,那么使用这个版本是最有意义的。
这样你就可以像这样直接将 Observable 用于单个请求:
mApiService.loadGameState(mGameId)
.observeOn(AndroidSchedulers.mainThread())
.subscribe(new Subscriber<GameState>() {
@Override
public void onNext(GameState gameState) {
// use the current game state here
}
// onError and onCompleted are also here
});
如果您想使用timer() 或interval() 的版本重复轮询服务器,可以提供“脉冲”:
Observable.timer(0, 2000, TimeUnit.MILLISECONDS)
.flatMap(mApiService.loadGameState(mGameId))
.observeOn(AndroidSchedulers.mainThread())
.subscribe(new Subscriber<GameState>() {
@Override
public void onNext(GameState gameState) {
// use the current game state here
}
// onError and onCompleted are also here
}).
需要注意的是,我在这里使用的是flatMap 而不是map - 这是因为loadGameState(mGameId) 的返回值本身就是一个Observable。
但是您在更新中使用的版本也应该可以工作:
Observable.interval(2, TimeUnit.SECONDS, Schedulers.io())
.map(tick -> Api.ReceiveGameTurn())
.doOnError(err -> Log.e("Polling", "Error retrieving messages" + err))
.retry()
.observeOn(AndroidSchedulers.mainThread())
.subscribe(sub);
也就是说,如果 ReceiveGameTurn() 像我上面的 1.) 一样被同步定义,您将使用 map 而不是 flatMap。
在这两种情况下,Subscriber 的 onNext 将每两秒调用一次,并提供来自服务器的最新游戏状态。您可以通过在subscribe() 之前插入take(1) 来逐个处理它们,将发射限制为单个项目。
然而,关于第一个版本:一个单一的网络错误将首先传递给onError,然后 Observable 将停止发射任何更多的项目,使您的订阅者无用且没有输入(请记住,onError 只能被调用一次)。要解决此问题,您可以使用 rxjava 的任何 onError* 方法将故障“重定向”到 onNext。
例如:
Observable.timer(0, 2000, TimeUnit.MILLISECONDS)
.flatMap(new Func1<Long, Observable<GameState>>(){
@Override
public Observable<GameState> call(Long tick) {
return mApiService.loadGameState(mGameId)
.doOnError(err -> Log.e("Polling", "Error retrieving messages" + err))
.onErrorResumeNext(new Func1<Throwable, Observable<GameState>(){
@Override
public Observable<GameState> call(Throwable throwable) {
return Observable.emtpy());
}
});
}
})
.filter(/* check if it is a valid new game state */)
.take(1)
.observeOn(AndroidSchedulers.mainThread())
.subscribe(new Subscriber<GameState>() {
@Override
public void onNext(GameState gameState) {
// use the current game state here
}
// onError and onCompleted are also here
}).
这将每两秒一次:
* 使用 Retrofit 从服务器获取当前游戏状态
* 过滤掉无效的
* 取第一个有效的
* 和退订
如果出现错误:
* 它将在doOnNext 中打印一条错误消息
* 否则忽略错误:onErrorResumeNext 将“消耗”onError-Event(即不会调用您的 Subscriber 的 onError)并将其替换为空(Observable.empty())。
并且,关于第二个版本:如果出现网络错误,retry 将立即重新订阅该时间间隔 - 由于interval 在订阅后立即发出第一个整数,因此下一个请求也将立即发送 - 而不是3 秒后,你可能想要...
最后说明:另外,如果您的游戏状态非常大,您也可以先轮询服务器以询问是否有新状态可用,只有在得到肯定答案的情况下才重新加载新的游戏状态。
如果您需要更详细的示例,请询问。
更新:我重写了这篇文章的部分内容,并在其间添加了更多信息。
更新 2:我添加了一个使用 onErrorResumeNext 处理错误的完整示例。