【问题标题】:Dispose inner subscription of merge处理合并的内部订阅
【发布时间】:2019-03-15 14:06:15
【问题描述】:

!!警告:Rx 新手!!

我们有多个价格信息。要求是订阅所有这些提要,并且每 1 秒只输出最新的滴答声(节流)

 public static class FeedHandler
{
        private static IObservable<PriceTick> _combinedPriceFeed = null;

          private static double _throttleFrequency = 1000;

        public static void AddToCombinedFeed(IObservable<PriceTick> feed)
        {
            _combinedPriceFeed = _combinedPriceFeed != null ? _combinedPriceFeed.Merge(feed) : feed;
            AddFeed(_combinedPriceFeed);
        }

              private static IDisposable _subscriber;

        private static void AddFeed(IObservable<PriceTick> feed)
        {
            _subscriber?.Dispose();
            _subscriber = feed.Buffer(TimeSpan.FromMilliseconds(_throttleFrequency)).Subscribe(buffer => buffer.GroupBy(x => x.InstrumentId, (key, result) => result.First()).ToObservable().Subscribe(NotifyClient));
        }

         public static void NotifyClient(PriceTick tick)
        {
        //Do some action
        }

}

代码有多个问题。如果我多次使用相同的提要调用 AddToCombinedFeed,则流将开始重复。例如。下面

IObservable<PriceTick> feed1;

FeedHandler.AddToCombinedFeed(feed1);//1 stream
FeedHandler.AddToCombinedFeed(feed1);//2 streams(even though the groupby and first() functions will prevent this effect to propagate further

这让我想到了这个问题。如果我想从合并的流中删除一个价格流,我该怎么做?

【问题讨论】:

  • 你一开始就不应该有内部订阅。
  • 避免内订阅的原因是什么?
  • 因为您最终可能会遇到竞争条件并且很难正确处理内部订阅。可观察管道被设计为自行清理,因此您应该始终编写单个查询而不是多个内部查询。

标签: c# reactive-programming system.reactive


【解决方案1】:

更新 - 新解决方案

来自 RolandPheasant 和 Nuget 的 Dynamic-Data(MIT-License)。

  1. 使用 SourceList 而不是 List
  2. 使用 MergeMany 运算符

代码:

public class FeedHandler
{
    private readonly IDisposable _subscriber;
    private readonly SourceList<IObservable<PriceTick>> _feeds = new SourceList<IObservable<PriceTick>>();
    private readonly double _throttleFrequency = 1000;

    public FeedHandler()
    {
        var combinedPriceFeed = _feeds.Connect().MergeMany(x => x).Buffer(TimeSpan.FromMilliseconds(_throttleFrequency)).SelectMany(buffer => buffer.GroupBy(x => x.InstrumentId, (key, result) => result.First()));
        _subscriber = combinedPriceFeed.Subscribe(NotifyClient);
    }

    public void AddFeed(IObservable<PriceTick> feed) => _feeds.Add(feed);

    public void NotifyClient(PriceTick tick)
    {
        //Do some action
    }
}

旧解决方案

  1. 通过应用 Switch() 技术消除重新订阅的需要。 您的 _combinedPriceFeed 只是切换到下一个可观察的 将由 _combinePriceFeedChange 提供。
    1. 保留一个列表以管理您的多个提要。每当列表更改时创建新的 observable,并通过 _combinePriceFeedChange 提供。
    2. 你应该得到对应的remove方法的逻辑。

代码:

public class FeedHandler
{
    private readonly IDisposable _subscriber;
    private readonly IObservable<PriceTick> _combinedPriceFeed;
    private readonly List<IObservable<PriceTick>> _feeds = new List<IObservable<PriceTick>>();
    private readonly BehaviorSubject<IObservable<PriceTick>> _combinedPriceFeedChange = new BehaviorSubject<IObservable<PriceTick>>(Observable.Never<PriceTick>());
    private readonly double _throttleFrequency = 1000;

    public FeedHandler()
    {
        _combinedPriceFeed = _combinedPriceFeedChange.Switch().Buffer(TimeSpan.FromMilliseconds(_throttleFrequency)).SelectMany(buffer => buffer.GroupBy(x => x.InstrumentId, (key, result) => result.First()));
        _subscriber = _combinedPriceFeed.Subscribe(NotifyClient);
    }

    public void AddFeed(IObservable<PriceTick> feed)
    {
        _feeds.Add(feed);
        _combinedPriceFeedChange.OnNext(_feeds.Merge());
    }


    public void NotifyClient(PriceTick tick)
    {
        //Do some action
    }
}

【讨论】:

  • 比起重新订阅更喜欢切换有什么实际优势吗? switch 不是在幕后做同样的事情吗?
  • 不需要重新订阅正是优势。它是一种处理处置/重新订阅的声明方式。 _combinedPriceFeed-Observable 始终有效。在我的代码中,大部分时间 IObservable 是一个公共属性或者它是由一个方法生成的,我不想破坏提供方的现有订阅,而不是错误/完成。
【解决方案2】:

这是您需要的代码:

private static SerialDisposable _subscriber = new SerialDisposable();

private static void AddFeed(IObservable<PriceTick> feed)
{
    _subscriber.Disposable =
        feed
            .Buffer(TimeSpan.FromMilliseconds(_throttleFrequency))
            .SelectMany(buffer =>
                buffer
                    .GroupBy(x => x.InstrumentId, (key, result) => result.First()))
            .Subscribe(NotifyClient);
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2012-09-07
    • 1970-01-01
    • 2023-03-14
    • 1970-01-01
    • 1970-01-01
    • 2018-09-20
    • 2021-11-25
    • 2021-08-15
    相关资源
    最近更新 更多