【发布时间】:2017-04-03 11:03:42
【问题描述】:
我正在使用 IObservable 推送数据更新/更改,我有一个方法可以从数据库 GetLatestElement 获取最新数据,每当有人调用 UpdateElement 并且数据得到更新时,都会通过消息传递消息系统。
所以我正在创建一个发出最新值的可观察对象,然后在它从消息系统接收到更新事件时重新发出新值:
public IObservable<IElement> GetElement(Guid id)
{
return Observable.Create<T>((observer) =>
{
observer.OnNext(GetLatestElement(id));
// subscribe to internal or external update notifications
var messageCallback = (message) =>
{
// new update message recieved,
observer.OnNext(GetLatestElement(id));
}
messageService.SubscribeToTopic(id, messageCallback);
return Disposable.Create(() => Console.Writeline("Observer Disposed"));
});
}
我的问题是这是不确定的。这些更新可能会永远发生。由于我试图让系统尽可能无状态,因此为GetElementType 的每个请求创建一个新的 Observable。这意味着生命周期由 订阅者 决定,而不是数据的来源。
我永远不会在 Observable 中调用 OnComplete(),我想在 Observer/User 完成后完成。
但是,我需要在某个时间点致电messageService.Unsubscribe(messageCallback);,以便在 Observable 完成后取消订阅消息。
我可以在订阅被释放时执行此操作,但我只能订阅一次,这似乎可能会引入错误。
Observables 应该如何做到这一点?
【问题讨论】:
-
也许在你的
Disposable.Create处理程序中完成它?顺便说一句,您的意思可能是“无限”,而不是“无限”。 -
啊,我明白了。我认为这是在处理 IObservable(当我认为 IObservable 是一次性的),所以它有效地完成了它。这将是一个解决方案,但它会将我的 observables 限制为单个订阅,不必这样做会很好。
-
老实说,我看不出有什么理由来完成这个序列。当订阅者取消订阅时-您从
messageService中删除messageCallback,因此没有资源会泄漏。当您完成序列时 - 您通知订阅者它已完成。但是在您所说的情况下-您不需要这个,因为订户在完成时知道自己。所以不要完成它。 -
这意味着我只能订阅一次。因此,如果我在任何时候使用 Take(1) 或任何其他取消订阅的方法,它都会破坏任何活动订阅。我已经稍微编辑了这个问题。
-
只有当您的
messageService实施不正确时才会如此。当您使用相同的 id 但不同的回调调用Subscribe多次时 - 它应该将这些回调添加到列表中,并在更新到达时将它们全部调用。当你取消订阅时 - 你应该传递 id 和回调,它应该只从列表中删除那个回调。
标签: c# system.reactive observable