【问题标题】:How does task cancellation work in RxJava?RxJava 中的任务取消是如何工作的?
【发布时间】:2014-08-16 21:18:48
【问题描述】:

我不清楚如何在 RXJava 中实现任务取消。

我有兴趣移植使用 Guava 的 ListenableFuture 构建的现有 API。我的用例如下:

  • 我有一个由Futures.transform() 连接的期货序列组成的单一操作
  • 多个订阅者观察操作的最终未来。
  • 每个观察者都可以取消最终的未来,并且所有观察者都会见证取消事件。
  • 取消最终未来会导致取消其依赖项,例如按顺序1->2->33 的取消传播到2,依此类推。

RxJava wiki 中关于此的信息很少;我能找到取消订阅的唯一参考提到 Subscription 相当于 .NET 的 Disposable,但据我所知,订阅仅提供取消订阅序列中后续值的功能。

我不清楚如何通过此 API 实现“任何订阅者都可以取消”语义。我是否以错误的方式思考这个问题?

我们将不胜感激。

【问题讨论】:

    标签: java rx-java


    【解决方案1】:

    了解Cold vs Hot Observables 很重要。如果您的 Observable 是冷的,那么如果您没有订阅者,它们的操作将不会执行。因此要“取消”,只需确保所有观察者都取消订阅源 Observable。

    但是,如果源中只有一个 Observer 取消订阅,并且还有其他 Observer 仍然订阅源,这不会导致“取消”。在这种情况下,您可以使用(但这不是唯一的解决方案)ConnectableObservables。另见this link about Rx.NET

    使用 ConnectableObservables 的一种实用方法是在任何冷的 Observable 上简单地调用 .publish().refCount()。这样做是创建一个单一的“代理”观察者,它将事件从源中继到实际的观察者。代理观察者在最后一个实际观察者退订时退订。

    要手动控制 ConnectableObservable,只需调用 coldSource.publish(),您将获得 ConnectableObservable 的实例。然后您可以致电.connect(),它将返回您“代理”观察者的订阅。要手动“取消”源,您只需取消订阅代理观察者的订阅即可。


    对于您的具体问题,您也可以使用.takeUntil() 运算符。

    假设你的“最终未来”在 RxJava 中被移植为 finalStream,并且假设“取消事件”是 Observables cancelStream1cancelStream2 等,那么“取消”由 @ 产生的操作变得相当简单987654333@:

    Observable<FooBar> finalAndCancelableStream = finalStream
        .takeUntil( Observable.merge(cancelStream1, cancelStream2) );
    

    在图表中,this is how takeUntil worksthis is how merge works

    简单来说,您可以将其解读为“finalAndCancelableStream 是 finalStream,直到 cancelStream1 或 cancelStream2 发出事件”。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-06-25
      • 2015-08-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多