【问题标题】:How to execute many RxJava2 flux in a row如何连续执行多个 RxJava2 Flux
【发布时间】:2019-12-20 19:41:34
【问题描述】:

我正在向自己介绍 RxJava2,但我觉得我做错了什么。就我而言,我想做一些以下异步操作。

在此示例中,第一个操作是检查设备是否已连接(wifi 或数据,让我们承认这需要时间),然后我想连接到 api,然后我想进行 http 调用以获取列表(可观察的),然后使用它。如果其中一项操作失败,则应在订阅中引发和处理 onError 或异常。

我有这个工作的代码:

Single.create((SingleEmitter<Boolean> e) -> e.onSuccess(Connectivity.isDeviceConnected(MainActivity.this)) )
    .subscribeOn(Schedulers.io())
    .flatMap(isDeviceConnected -> {
        Log.i("LOG", "isDeviceConnected : "+ isDeviceConnected);
        if(!isDeviceConnected)
            throw new Exception("whatever"); // TODO : Chercher vrai erreur

        return awRepository.getFluxAuthenticate(host, port, user, password); // Single<DisfeApiAirWatch>
    })
    .toObservable()
    .flatMap(awRepository::getFluxManagedApps)  // List of apps : Observable<AirwatchApp>

    .observeOn(AndroidSchedulers.mainThread())
    .doFinally(this::hideProgressDialog)
    .subscribe(
            app -> Log.i("LOG", "OnNext : "+ app),
            error -> Log.i("LOG", "Error : " + error),
            () -> Log.i("LOG", "Complete : ")
);

但是做一个为简单的“如果”发出布尔值的人听起来是错误的。 Completable 似乎更合乎逻辑(工作与否,继续或停止)。我尝试使用以下代码,但它不起作用。

Completable.create((CompletableEmitter e) -> {
    if(Connectivity.isDeviceConnected(MainActivity.this))
        e.onComplete(); // Guess not good, should call the complete of subscribe ?
    else
        e.onError(new Exception("whatever"));
} ).toObservable()
    .subscribeOn(Schedulers.io())
    .flatMap(awRepository.getFluxAuthenticate(host, port, user, password)) //Single<DisfeApiAirWatch>
    .toObservable()
    .flatMap(awRepository::getFluxManagedApps) // List of apps : Observable<AirwatchApp>

    .observeOn(AndroidSchedulers.mainThread())
    .doFinally(this::hideProgressDialog)
    .subscribe(
            app -> Log.i("LOG", "OnNext : "+ app),
            error -> Log.i("LOG", "Error : " + error),
            () -> Log.i("LOG", "Complete : ")
);

如何使这段代码工作?

我知道我可以先订阅可兼容的,然后在这个的“onSuccess”中编写另一个通量/其余代码。但我不认为堆栈在彼此内部流动是一个好的解决方案。

最好的问候

【问题讨论】:

    标签: android asynchronous rx-java2


    【解决方案1】:

    Completable 没有任何值,因此永远不会调用 flatMap。您必须使用andThen 并将身份验证成功值作为后续flatMap 的输入:

    Completable.create((CompletableEmitter e) -> {
        if(Connectivity.isDeviceConnected(MainActivity.this))
            e.onComplete();
        else
           e.onError(new Exception("whatever"));
    })
    .subscribeOn(Schedulers.io())
    .andThen(awRepository.getFluxAuthenticate(host, port, user, password)) // <-----------
    .flatMapObservable(awRepository::getFluxManagedApps)
    .observeOn(AndroidSchedulers.mainThread())
    .doFinally(this::hideProgressDialog)
    .subscribe(
            app -> Log.i("LOG", "OnNext : "+ app),
            error -> Log.i("LOG", "Error : " + error),
            () -> Log.i("LOG", "Complete : ")
    

    );

    【讨论】:

    • 我认为它工作正常,但是 .andThen() 中的方法在 Completable.create() 中的代码之前被调用...
    • 使用此行解决:.andThen(Single.defer(() -> awRepository.getFluxAuthenticate(host, port, user, password)))
    猜你喜欢
    • 2018-01-28
    • 2021-05-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-08-10
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多