【问题标题】:With Rx, How to get the latest value per unique key?使用 Rx,如何获取每个唯一键的最新值?
【发布时间】:2021-06-02 00:57:55
【问题描述】:

如何获取每个唯一键的最新值?

在这个示例中,我有股票代码和价格。我如何组合它以输出具有最新价格的独特股票?

在现实生活中,这种事件流将持续数年,我只需要股票的当前价格。

using System;
using System.Reactive.Concurrency;
using System.Reactive.Linq;

namespace LearnRx
{
    class Program
    {
        public class Stock
        {
            public string Symbol { get; set; }
            public float Price { get; set; }
            public DateTime ModifiedAt { get; set; }

            public override string ToString()
            {
                return "Stock: " + Symbol + " " + Price + " " + ModifiedAt;
            }
        }

        static void Main(string[] args)
        {
            var baseTime = new DateTime(2000, 1, 1, 1, 1, 1);
            var fiveMinutes = TimeSpan.FromMinutes(5);
            var stockEvents = new Stock[]
            {
                new Stock
                {
                    Symbol = "MSFT",
                    Price = 100,
                    ModifiedAt = baseTime
                },
                new Stock
                {
                    Symbol = "GOOG",
                    Price = 200,
                    ModifiedAt = baseTime + fiveMinutes
                },
                new Stock
                {
                    Symbol = "MSFT",
                    Price = 150,
                    ModifiedAt = baseTime + fiveMinutes + fiveMinutes
                },
                new Stock
                {
                    Symbol = "AAPL",
                    Price = 300,
                    ModifiedAt = baseTime + fiveMinutes + fiveMinutes + fiveMinutes
                },
            }.ToObservable();

            var scheduler = new HistoricalScheduler();
            var replay = Observable.Generate(
                stockEvents.GetEnumerator(),
                events => events.MoveNext(),
                events => events,
                events => events.Current,
                events => events.Current.ModifiedAt,
                scheduler);


            replay
                .Subscribe( i => Console.WriteLine($"Event: {i} happened at {scheduler.Now}"));

            scheduler.Start();
        }
    }
}

【问题讨论】:

  • 你能贴一张你想要的大理石图吗?如果您总是对最新价格感兴趣......这只是最新的活动。
  • 上述代码的输出有 4 个条目,预期为 3。MSFT 应该只有一个输出。
  • “在现实生活中,这种事件流将持续数年,我只需要股票的当前价格。”在任何给定时刻,最近的事件是当前价格。你想发出所有事件吗?在测试用例中有结束,在现实生活中没有。
  • 假设我想要计算 90 美元以上的股票数量,上面的代码将返回 4 而不是 3,因为它会重复计算 MSFT。就像我需要一个 combineLastest ,它需要一个键,在这种情况下是股票代码。 stock.CombineLastest(r => r.Symbol).Where(r => Price > 99)
  • 作为一名在日常工作中实际执行此操作的开发人员,我建议您考虑使用IDictionary<Stock, IObservable<Price>> 而不是IObservable<IDictionary<Stock, Price>>。后者使您很难从订阅中删除股票...

标签: system.reactive


【解决方案1】:

听起来你想要一个 GroupBy 后跟一个 CombineLatest 应该存在,但不存在。

这是缺少的CombineLatest

public static class X
{
    
    public static IObservable<ImmutableDictionary<TKey, TValue>> CombineLatest<TKey, TValue>(this IObservable<IGroupedObservable<TKey, TValue>> source)
    {
        return source
            .SelectMany(o => o.Materialize().Select(n => (Notification: n, Key: o.Key)))
            .Scan(default(Notification<ImmutableDictionary<TKey, TValue>>), (stateNotification, t) => {
                var state = stateNotification?.Value ?? ImmutableDictionary<TKey, TValue>.Empty;
                switch(t.Notification.Kind) {
                    case NotificationKind.OnError:
                        return Notification.CreateOnError<ImmutableDictionary<TKey, TValue>>(t.Notification.Exception);
                    case NotificationKind.OnCompleted:
                        return Notification.CreateOnNext(state.Remove(t.Key));
                    case NotificationKind.OnNext:
                        return Notification.CreateOnNext(state.SetItem(t.Key, t.Notification.Value));
                    default:
                        throw new NotImplementedException();
                }
            })
            .Dematerialize();
    }
}

然后这是Main中的最终查询:

var prices = replay
    .GroupBy(s => s.Symbol)
    .CombineLatest()

这将为您提供每个价格点的Observable&lt;Dictionary&lt;string, Stock&gt;&gt;。这应该使您能够在下游执行任何操作。

【讨论】:

  • 不错。它正在获取最新事件。使用带有 upsert 的字典是我正在考虑的解决方案。该代码确实有一个有趣的输出。我了解迭代 #4 之前的输出,但在那之后,为什么这些项目会一一消失? #1 字典计数 1 MSFT 100 #2 字典计数 2 GOOG 200 MSFT 100 #3 字典计数 2 GOOG 200 MSFT 150 #4 字典计数 3 GOOG 200 MSFT 150 AAPL 300 #5 字典计数 2 GOOG 200 AAPL 300 #6 字典计数 1 AAPL 300 #7 字典计数 0
  • 该代码支持分组结束,模仿被退市的股票或其他东西。如果其中一个子 observables 结束,它就会从字典中消失。在测试代​​码中,四个子可观察对象中的每一个都在父对象之前结束,从而创建您所看到的输出。
猜你喜欢
  • 1970-01-01
  • 2017-06-11
  • 2021-11-15
  • 1970-01-01
  • 2022-09-23
  • 2010-12-17
  • 2017-02-09
  • 1970-01-01
  • 2016-05-06
相关资源
最近更新 更多