【问题标题】:How to get the Count of an IObservable如何获取 IObservable 的计数
【发布时间】:2014-12-03 07:50:14
【问题描述】:

我迷路了。

我尝试从一长串数据库记录中获取一些聚合值(主要是Count)。我们曾经使用常规的 linq,但是数据量已经变得太大而无法放入内存。我想简单地将查询转换为 IObservable 以获得“流式传输”结果。但我一定是错过了什么。大多数示例和文档似乎都没有考虑到这种情况。也许 Rx 不是正确的工具集。

因此,为了重现问题,我只是生成了一些随机数据。然后我将它分组并期望每个组中的项目数。

static void Main(string[] args)
{
    // lots of DB records; from IEnumerable to IObservable
    var records= //GetRecords();
    ///*
        Observable.Interval(TimeSpan.FromMilliseconds(11))
                     .Take(100)
                     .Select(i => new { Group = DateTime.Now.Millisecond % 10 });
    //*/
    var result = from r in records
                 group r by r.Group into g
                 select new
                 {
                     Key = g.Key,
                     Count = g.Count()
                 };

    foreach (var item in result.ToEnumerable())
    {
        Console.WriteLine("{0} - {1}", item.Key, item.Count.Wait());
    }
}

结果只会给我第一项的值:

8 - 12
4 - 0
0 - 0
5 - 0
1 - 0
7 - 0
2 - 0
3 - 0
9 - 0
6 - 0

我在这里做错了什么?

【问题讨论】:

  • count.Wait() 是做什么的?
  • Rx 似乎是一个奇怪的选择。您的应用程序不是关于事件,而是关于非常适合IEnumerable<T> 和 LINQ 的记录。您不必加载所有记录来计算它们。创建一个迭代器块,一次只迭代一个记录,然后将 LINQ 应用于生成的IEnumerable<T>。您甚至可以使用 P-LINQ 来加速聚合。
  • 我们目前正在使用 LINQ,但我们所做的聚合需要我们将所有数据加载到内存中。它曾经按预期工作。这是关于内存使用,而不是性能。顺便说一句,我们的记录代表事件,我将修改名称以避免与 Rx 混淆,例如鼠标事件
  • 如果您必须将所有数据加载到内存中以执行聚合,那么我看不出如何通过使用 Rx 来减少内存使用量?或任何其他解决方案。
  • 我们加载的数据作为 blob 存储在 DB 中(CQRS/ES 场景中的事件),因此最简单的解决方案是将所有数据加载到内存中,然后使用 LINQ to 对象执行聚合。使用 Rx,我们可以加载数据并让聚合以流式方式计算值。

标签: c# linq system.reactive


【解决方案1】:

Count() 返回一个IObservable<int>,并且您订阅它的时间较晚,也就是说,当您从events 观察到所有值时。我不是 100% 确定 Group 的行为,但您似乎需要提前订阅 Count() 以避免丢失元素。将.ToTask() 添加到.Count() 看看会发生什么。这样一来,调用Wait() 还是有意义的。

【讨论】:

  • 好的,因此将其转换为任务会执行订阅并允许我返回结果。谢谢
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多