【问题标题】:RxJava - Trigger another Observable/Completable before complete and complete only after itRxJava - 在完成之前触发另一个 Observable/Completable 并且仅在它之后完成
【发布时间】:2016-10-31 17:59:15
【问题描述】:

我对 RxJava 并不完全陌生,但我被看似简单的任务所困。

我有一个公开反应式 API 的数据源,我要做的就是获取一些数据,返回它并在没有其他要发出的内容时自动关闭连接。

这是我的代码:

public Observable<Object> execute(String query) {

    Single<RxConnection> rxConnection = getRxDB().getConnection();

    return rxConnection.flatMapObservable(conn -> {
        Observable<Object> rxResult = conn.query(query);

        return rxResult.doOnCompleted(() -> {
            conn.close(); // THIS DOES NOT WORK. I would like to close the connection and to wait without blocking.
        });

    });

}

conn.query() 和 conn.close() 在不同的 Scheduler 中异步执行。 此代码不起作用,因为 conn.close() 返回一个没有订阅者的 Completable。另外,如果我在 doOnCompleted 方法本身中手动订阅,则 rxResult Observable 会在不等待连接关闭的情况下完成。

我希望“执行(字符串查询)”方法返回一个 Observable: - 发出 conn.query() 调用获取的所有项目 - 当没有更多的项目要发射时,它会触发 conn.close() - 仅在 conn.close() Completable 之后完成;

谢谢。

【问题讨论】:

    标签: java rx-java


    【解决方案1】:

    资源关闭由Observable.using管理:

    Observable<T> obs = 
        Observable.using(
            resourceFactory,
            observableFactory,
            disposeAction)
    

    这种创建方法确保resourceFactory 创建的资源在终止(完成或错误)或取消订阅时被释放。

    【讨论】:

    • 这是一个有趣的解决方案;但是,我不能让它工作。实际上,resourceFactory 不允许返回 Observable,而只允许返回一个资源实例,这不是我的情况。
    • resourceFactory 提供资源,observableFactory 从资源(由resourceFactory 创建)创建一个可观察对象。 @akarnokd 正在涵盖您的具体问题的详细信息,但您应该使用 Observable.using 来获得适当的资源关闭。
    【解决方案2】:

    怎么样:

    Observable<Object> rxResult = conn.query(query)
    .concatWith(conn.close().toObservable())
    .onErrorResumeNext(e -> 
        conn.close().toObservable().concatWith(Observable.error(e)));
    

    【讨论】:

    • 这似乎以一种非常简单的方式解决了我的问题。此外,它还处理可能的错误。然而,可怕的是 rx 迫使您将所有内容转换为可观察的,从而创建大量无用的对象。
    【解决方案3】:

    因此,您只想在订阅 observable 时执行查询:

    private Observable<Object> executeDeferred(String query) {
       // your original execute() method here, including doOnComplete
    }
    
    
    public Observable<Object> execute(String query) {
      return Observable.defer(() -> executeDeferred(query));
    }
    

    这将确保连接、查询和释放仅在 observable 获得订阅时发生,而不是在设置期间发生。

    编辑:您还需要一个 doOnUnsubscribe / doOnTerminate,以便在发生过早取消订阅时关闭连接。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-12-13
      • 2016-11-20
      • 1970-01-01
      • 2020-07-06
      • 2020-09-27
      相关资源
      最近更新 更多