【问题标题】:Proxy an Observable and connect it in callback代理一个 Observable 并在回调中连接它
【发布时间】:2018-11-04 13:58:59
【问题描述】:

我正在尝试返回一个在回调中异步创建的Observable

const mkAsync = (observer, delay) =>
  setTimeout(() => Observable.of('some result').subscribe(observer), delay)

const create = arg => {
  const ret = new Subject()
  mkAsync(ret, arg)
  return ret
}

因此,我使用Subject 作为单播代理,它在回调中订阅了底层Observable

这个解决方案的问题是,当我取消订阅Subject 的订阅时,取消订阅不会转发到底层Observable。看起来我需要某种类型的引用计数来使 Subject 在没有更多订阅者时取消订阅,但是在这种命令式回调样式中使用它时我无法弄清楚。

我必须将mkAsync 保留为空,并且正在寻找替代实现。

这是正确的做法吗?有没有使用Subject 的替代解决方案?

如何确保在取消订阅主题时取消创建的Observable(在Subscription 上调用unsubscribe)?

【问题讨论】:

    标签: rxjs observable subscription subject-observer


    【解决方案1】:

    这是一个相当广泛的问题,很难说你想用这个实现什么。我有两个想法:

    第一件事是refCount() 运算符仅存在于ConnectableObservable 类上,该类根据您传递的参数从multicast(或publish)返回。更多细节见实现(基本上如果你没有设置任何selector函数):https://github.com/ReactiveX/rxjs/blob/5.5.11/src/operators/multicast.ts

    我能想到的第二个问题是你基本上是这样做的:

    const ret = new Subject()
    Observable.of(...).subscribe(ret);
    

    问题在于.of 会立即发出next 项目,然后发送complete 通知。主体有内部状态,当Subject 收到complete 通知时,它会将自己标记为stopped,并且它永远不会发出任何东西

    我怀疑这就是发生在你身上的事情。即使您返回带有return ret 的Subject 实例并且稍后可能订阅它,您仍然不会收到任何东西,因为该Subject 已经收到complete 通知。

    【讨论】:

    • 似乎refCount 给了我一个新的 Observable 但我已经返回了主题。目前的方法基本上是有效的,对于退订问题是安全的。我可以使用 ReplaySubject 来不错过任何事件,但现在我会立即订阅该主题。
    猜你喜欢
    • 2021-02-05
    • 1970-01-01
    • 1970-01-01
    • 2018-08-04
    • 1970-01-01
    • 2017-08-25
    • 1970-01-01
    • 2021-05-24
    • 1970-01-01
    相关资源
    最近更新 更多