【问题标题】:Android: Polling a server with RetrofitAndroid:使用 Retrofit 轮询服务器
【发布时间】:2015-04-06 19:50:22
【问题描述】:

我正在 Android 上构建 2 人游戏。游戏轮流进行,所以玩家 1 等到玩家 2 输入,反之亦然。我有一个网络服务器,我在其中运行带有 Slim 框架的 API。在我使用 Retrofit 的客户端上。因此,在客户端上,我想每隔 X 秒轮询一次我的网络服务器(我知道这不是最好的方法),以检查是否有来自玩家 2 的输入,如果是,请更改 UI(游戏板)。

处理 Retrofit 我遇到了 RxJava。我的问题是弄清楚我是否需​​要使用 RxJava?如果是,是否有任何非常简单的轮询示例进行改造? (因为我只发送了几个键/值对)如果不是如何通过改造来做到这一点?

我在这里找到了this 线程,但它也没有帮助我,因为我仍然不知道我是否需要 Retrofit + RxJava,有没有更简单的方法?

【问题讨论】:

  • 你不需要 RxJava。但它确实很适合改装。
  • 那么如何在没有 RxJava 的情况下使用 Retrofit 实现轮询?

标签: android rest polling retrofit rx-java


【解决方案1】:

假设您为 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 处理错误的完整示例。

【讨论】:

  • 稍后我将不得不看看错误处理。但是我们不使用doOnError 进行错误处理吗?
  • 这取决于您所说的“处理”。 doOnError 只是在将错误传递到 Observable 管道时做一些事情。但我认为,这并不是真正意义上的解决错误的地方。默认情况下,错误可以在onError 中处理——但这也会完全终止 Observable。还有其他选项可以对错误做出反应,以保持订阅不变:发出其他内容 (onErrorResumeNext)、retry 等。
【解决方案2】:

谢谢,我终于以我在问题中提到的post 的类似方式做到了。这是我现在的代码:

Subscriber sub =  new Subscriber<Long>() {
        @Override
        public void onNext(Long _EmittedNumber)
        {
            GameTurn Turn =  Api.ReceiveGameTurn(mGameInfo.GetGameID(), mGameInfo.GetPlayerOneID());
            Log.d("Polling", "onNext: GameID - " + Turn.GetGameID());
        }

        @Override
        public void onCompleted() {
            Log.d("Polling", "Completed!");
        }

        @Override
        public void onError(Throwable e) {
            Log.d("Polling", "Error: " + e);
        }
    };

    Observable.interval(3, TimeUnit.SECONDS, Schedulers.io())
            // .map(tick -> Api.ReceiveGameTurn())
            // .doOnError(err -> Log.e("Polling", "Error retrieving messages" + err))
            .retry()
            .subscribe(sub);

现在的问题是,当我得到肯定的答案(GameTurn)时,我需要终止发射。我读到了takeUntil 方法,我需要传递另一个Observable,它会发出一些东西,这会触发我的轮询终止。但我不确定如何实现这一点。 根据您的解决方案,您的 API 方法会返回 Observable,就像 Retrofit 网站上显示的那样。也许这是解决方案?那么它是如何工作的呢?

更新: 我考虑了@david.miholas 的建议,并通过重试和过滤尝试了他的建议。您可以在下面找到游戏初始化的代码。轮询应该是一样的:玩家 1 开始新游戏 -> 轮询对手,玩家 2 加入游戏 -> 服务器发送给玩家 1 对手的 ID -> 轮询终止。

    Subscriber sub =  new Subscriber<String>() {
        @Override
        public void onNext(String _SearchOpponentResult) {}

        @Override
        public void onCompleted() {
            Log.d("Polling", "Completed!");
        }

        @Override
        public void onError(Throwable e) {
            Log.d("Polling", "Error: " + e);
        }
    };

    Observable.interval(3, TimeUnit.SECONDS, Schedulers.io())
            .map(tick -> mApiService.SearchForOpponent(mGameInfo.GetGameID()))
            .doOnError(err -> Log.e("Polling", "Error retrieving messages: " + err))
            .retry()
            .filter(new Func1<String, Boolean>()
            {
                @Override
                public Boolean call(String _SearchOpponentResult)
                {
                    Boolean OpponentExists;
                    if (_SearchOpponentResult != "0")
                    {
                        Log.e("Polling", "Filter " + _SearchOpponentResult);
                        OpponentExists = true;
                    }
                    else
                    {
                        OpponentExists = false;
                    }
                    return OpponentExists;

                }
            })
            .take(1)
            .subscribe(sub);

发射是正确的,但是每次发射时我都会收到以下日志消息:

E/Polling﹕ Error retrieving messages: java.lang.NullPointerException

显然doOnError 会在每次发射时触发。通常我会在每次发射时得到一些改造调试日志,这意味着mApiService.SearchForOpponent 不会被调用。我做错了什么?

【讨论】:

  • 也许您可以更详细地解释您想要实现的目标:起初我以为您想在程序运行的整个过程中每两秒轮询一次服务器,但现在您说想要在您从服务器成功回复后终止项目的发射...您当然可以在 subscribe 之前插入一个 take(1) 以将发射限制为一个项目,但您也可以看看 @ 987654333@ - 也许这更接近你真正想要的......
  • 这是一个回合制游戏,所以玩家等待/轮询服务器以进行更改。当其他玩家设置输入时,他会发送轮到他的输入,该输入存储在数据库中。然后,第一个玩家收到带有轮次信息的正面 MySQL 结果,这是更新 UI 和终止轮询所必需的,因为已收到信息。官方 Rx wiki 说:“如果源 Observable 发出错误,则将该错误传递给另一个 Observable 以确定是否重新订阅源” 这对我没有帮助
  • 你发布的版本做了你想要的,只是它没有在一个发射项目后停止?在这种情况下,您实际上可以将take(1) 放在retry() 和subscribe() 之间——这样take 只会“通过”一项然后取消订阅,从而停止在interval 中生成刻度。但是,我认为在例如的情况下网络错误retry 将立即重新订阅interval - 并且由于interval 在订阅后立即发出第一个整数,因此下一个请求也将立即发送 - 而不是您可能想要的 3 秒后...
  • 它必须在一个有效的GameTurn 之后停止发射,而不是在发射一个之后。 retry 将如何帮助我?我应该在onCompleted 而不是onNext 中调用我的请求吗?如果出现网络错误,它可能会有所帮助。但是对我的请求的无效答案不是程序错误,它表明还没有发生任何事情(其他玩家还没有轮到他)。因此,错误和否定响应之间存在差异。在否定响应时 -> 继续发射,出错时 -> 告诉 UI“发生错误”并针对它做一些事情(也终止发射,重新开始,等等)
  • 好的,是的,这是有道理的 - 那么,我可以通过查看结果本身来确定 MySQL 结果是否有效,即其他玩家进行的新回合,还是我需要与之前存储在本地的游戏状态进行比较?如果我只需要响应本身,您可以例如 filter 删除任何不需要的/无效/非新项目 - 它们只会从流中消失,take(1) 会再次起作用。
猜你喜欢
  • 2012-07-31
  • 2021-12-29
  • 2013-05-25
  • 2011-04-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-07-03
  • 1970-01-01
相关资源
最近更新 更多