【问题标题】:Iterate IEnumerable selected from IGroupedObservable in RX迭代从 RX 中的 IGroupedObservable 中选择的 IEnumerable
【发布时间】:2018-03-21 05:42:29
【问题描述】:

我有一个IObservable<T> 序列,其中T 是一个KeyValuePair<TKey, TValue>,我使用来自System.Reactive.LinqGroupBy 对其进行分组。

我想对每个 IGroupedObservable<TKey, KeyValuePair<TKey, TValue>> 执行聚合操作,但该聚合被定义为 Func<IEnumerable<TValue>, TValue>

例如,在这里我想计算每个不同单词出现的次数并将其打印到控制台:

Func<IEnumerable<int>, int> aggregate = x => x.Count();

using (new[] { "one", "fish", "two", "fish" }
    .Select(x => new KeyValuePair<string, int>(x, 1))
    .ToObservable()
    .GroupBy(x => x.Key)
    .Select(x => new KeyValuePair<string, IEnumerable<int>>(
                x.Key,
                x.Select(y => y.Value).ToEnumerable()))
    //.SubscribeOn(Scheduler.Default)
    .Subscribe(x => Console.WriteLine($"{x.Key} [{aggregate(x.Value)}]")))
{
}

我希望输出与此类似(顺序不重要):

one [1]
fish [2]
two [1]

但相反,它要么阻塞(可能是死锁),要么根本不提供任何输出(当我取消注释 LINQ 语句的 SubscribeOn 子句时)。

我已尝试从实际使用场景中减少上述代码,该场景尝试链接两个 TPL Dataflow 块但遇到类似行为:

Func<IEnumerable<int>, int> aggregate = x => x.Sum();

var sourceBlock = new TransformBlock<string, KeyValuePair<string, int>>(x => new KeyValuePair<string, int>(x, 1));
var targetBlock = new ActionBlock<KeyValuePair<string, IEnumerable<int>>>(x => Console.WriteLine($"{x.Key} [{aggregate(x.Value)}]"));
using (sourceBlock.AsObservable()
    .GroupBy(x => x.Key)
    .Select(x => new KeyValuePair<string, IEnumerable<int>>(x.Key, x.Select(y => y.Value).ToEnumerable()))
    .Subscribe(targetBlock.AsObserver()))
{
    foreach (var kvp in new[] { "one", "fish", "two", "fish" })
    {
        sourceBlock.Post(kvp);
    }
    sourceBlock.Complete();
    targetBlock.Completion.Wait();
}

我知道有框架提供了 SumCount 方法,可在 IObservable&lt;T&gt; 上运行,但我受限于 IEnumerable&lt;T&gt; 聚合函数。

我误解了ToEnumerable,我该怎么做才能解决它?

编辑: IEnumerable&lt;T&gt; 的约束是由我试图链接的两个数据流块的 target 引入的,其签名不是我可以更改的。

【问题讨论】:

    标签: c# linq system.reactive tpl-dataflow


    【解决方案1】:

    GroupBy 的工作原理是这样的:当新元素到达时,它会提取一个键并查看该键之前是否已经被观察过。如果不是 - 它会创建新组(新的可观察对象)并将密钥和可观察对象推送给您。关键点是 - 当您订阅 GroupBy 并将项目推送到您的订阅时 - 序列尚未分组。推送的是组键和另一个 observable (IGroupedObservable),该组中的元素将被推送到。

    您在代码中所做的基本上是订阅GroupBy,然后在GroupBy 订阅中阻塞以尝试枚举IGroupingObservable。但是此时您无法枚举它,因为分组不完整。为了完成它 - GroupBy 应该处理整个序列,但它不能,因为它被阻止等待您的订阅处理程序完成。并且您的订阅处理程序等待 GroupBy 完成(阻止尝试枚举尚未准备好的序列)。因此你有一个死锁。

    如果您尝试引入 ObserveOn(Scheduler.Default) 在线程池线程上运行您的订阅处理程序 - 这将无济于事。它将消除死锁,但会引入竞争条件并且您将丢失项目,因为您在开始枚举ToEnumerable 的结果时只订阅单个组。此时可能为时已晚,并且在您订阅它之前(通过开始枚举),一些和一些项目已经被推送到单独的组 observable。这些项目不会重播,因此会丢失。

    什么将有助于它确实使用为IObservable 提供的Count(),但由于某种原因你说你不能这样做。

    在您使用数据流块的情况下,您可以尝试以下操作:

    sourceBlock.AsObservable()
        .GroupBy(x => x.Key)
        .Select(x => {
            var res = new { x.Key, Value = x.Select(y => y.Value).Replay() };
            // subscribe right here
            // Replay will ensure that no items are missed
            res.Value.Connect();
            return res;
        })                
        // observe on thread pool threads to not deadlock if necessary
        // in the example with datablock in your question - it is not
        //.ObserveOn(Scheduler.Default)
        // now no deadlock and no missing items
        .Select(x => new KeyValuePair<string, IEnumerable<int>>(x.Key, x.Value.ToEnumerable()))
        .Subscribe(targetBlock.AsObserver())
    

    【讨论】:

    • 如果您查看数据流示例,您会看到块签名将 IEnumerable&lt;int&gt; 作为每个输入的一部分,并且块应用聚合函数(不一定是 CountSum :这些是我的 MVCE 的玩具示例)。希望我的编辑能更清楚地说明这一点。
    • @Jono 我已经用想到的一个选项更新了答案。
    • 感谢您的帮助。我希望我们可以突出显示答案的一部分,因为当您说“您只订阅......开始枚举时”时,它开始变得有意义。我尝试了自己的快速Replay,但我错过了Connect,所以它仍然存在竞争条件。你的回答解决了这个问题。
    猜你喜欢
    • 1970-01-01
    • 2017-11-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-08-22
    相关资源
    最近更新 更多