【问题标题】:How to manage a lifetime of an infinite observable when it's dictated by the subscriber not the source当订阅者而不是源指定时,如何管理无限可观察的生命周期
【发布时间】: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


【解决方案1】:

似乎对Observable.Create 的工作方式存在一些误解。每当您在 GetElement() 的结果上调用 Subscribe 时 - Observable.Create 的主体就会被执行。因此,对于每个订阅者,您都单独订阅您的messageService,并执行单独的回调。如果您取消订阅 - 您只会删除 那个 订阅者的订阅。所有其他人保持活跃,因为他们有自己的messageCallback。这当然是假设messageService 已正确实施。以下是说明这一点的示例应用程序:

static IElement  GetLatestElement(Guid id) {
    return new Element();
}

public class Element : IElement {

}

public interface IElement {

}

class MessageService {
    private Dictionary<Guid, Dictionary<Action<IElement>, CancellationTokenSource>> _subs = new Dictionary<Guid, Dictionary<Action<IElement>, CancellationTokenSource>>();
    public void SubscribeToTopic(Guid id, Action<IElement> callback) {
        var ct = new CancellationTokenSource();
        if (!_subs.ContainsKey(id))
            _subs[id] = new Dictionary<Action<IElement>, CancellationTokenSource>();
        _subs[id].Add(callback, ct);
        Task.Run(() =>
        {
            while (!ct.IsCancellationRequested) {
                callback(new Element());
                Thread.Sleep(500);
            }
        });
    }

    public void Unsubscribe(Guid id, Action<IElement> callback) {
        _subs[id][callback].Cancel();
        _subs[id].Remove(callback);
    }
}

public static IObservable<IElement> GetElement(Guid id)
{
    var messageService = new MessageService();
    return Observable.Create<IElement>((observer) =>
    {
        observer.OnNext(GetLatestElement(id));

        // subscribe to internal or external update notifications
        Action<IElement> messageCallback = (message) =>
        {
            // new update message recieved,
            observer.OnNext(GetLatestElement(id));
        };
        messageService.SubscribeToTopic(id, messageCallback);

        return Disposable.Create(() => {
            messageService.Unsubscribe(id, messageCallback);
            Console.WriteLine("Observer Disposed");
        });
    });
}

public static void Main(string[] args) {
    var ob = GetElement(Guid.NewGuid());
    var sub1 = ob.Subscribe(c =>
    {
        Console.WriteLine("got element");
    });

    var sub2 = ob.Subscribe(c =>
    {
        Console.WriteLine("got element 2");
    });
    // at this point we see both subscribers receive messages
    Console.ReadKey();
    sub1.Dispose();
    // first one is unsubscribed, but second one is still alive
    Console.ReadKey();
}

所以正如我所说的 cmets - 在这种情况下,我认为没有理由完成您的 observable。

【讨论】:

  • 我明白了,是的,我从根本上误解了可观察的创造。我假设 Observable.Create 中的回调在创建时被调用一次,但我现在明白了。每个订阅调用一次。有道理,因为 Create 方法 get 是 Observer 的副本。感谢您帮助解决这个问题。
  • 您认为在这种情况下,使用主题(或带有单个项目缓冲区的 RecordSubject)会更明智吗?在这种情况下,不会为每个订阅创建一个新的观察者(而只是提供最新的值),它允许我减少额外的订阅(从而减少消息负载)。
  • 我认为阅读本文会在这方面对您有所启发:introtorx.com/content/v1.0.10621.0/…。简而言之,如果您想共享您的服务订阅 - 您需要将您的冷可观察变为热,例如使用 Publish 和 RefCount,正如另一个答案所建议的那样。
  • 谢谢。我会看到那个页面,但我想我现在有了更多理解它的上下文。 Publish 和 RefCount 看起来正是我所需要的。
【解决方案2】:

正如 Evk 指出的,Observable.Create 运行然后几乎立即处理。如果您想保持messageService 订阅开放,Rx 可以帮助您。看看MessageObservableProvider。剩下的就是编译:

public class MessageObservableProvider
{
    private MessageService messageService;
    private Dictionary<Guid, IObservable<Unit>> _messageNotifications = new Dictionary<Guid, IObservable<Unit>>();
    private IObservable<Unit> GetMessageNotifications(Guid id)
    {
        return Observable.Create<Unit>((observer) =>
        {
            Action<Message> messageCallback = _ => observer.OnNext(Unit.Default);
            messageService.SubscribeToTopic(id, messageCallback);

            return Disposable.Create(() =>
            {
                messageService.Unsubscribe(messageCallback);
                Console.WriteLine("Observer Disposed");
            });
        });
    }

    public IObservable<IElement> GetElement(Guid id)
    {
        if(!_messageNotifications.ContainsKey(id))
            _messageNotifications[id] = GetMessageNotifications(id).Publish().RefCount();

        return _messageNotifications[id]
            .Select(_ => GetLatestElement(id))
            .StartWith(GetLatestElement(id));
    }

    private IElement GetLatestElement(Guid id)
    {
        throw new NotImplementedException();
    }
}

public class IElement { }
public class Message { }
public class MessageService
{
    public void SubscribeToTopic(Guid id, Action<Message> callback)
    {
        throw new NotImplementedException();
    }

    public void Unsubscribe(Action<Message> callback)
    {
        throw new NotImplementedException();
    }
}

您最初的Create 实现包含StartWithSelect 的功能。我把它们移走了,所以现在Observable.Create 只会在有新消息可用时返回通知。

更重要的是,在GetElement 中现在有一个.Publish().RefCount() 调用。这将使messageService 订阅保持打开状态(通过不调用.Dispose()),只要至少有一个孩子可观察(订阅)在附近闲逛。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-04-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-06-28
    • 2016-06-26
    • 1970-01-01
    相关资源
    最近更新 更多