C# Rx管道中GroupBy仅首次执行,无法持续运行求助
Rx管道分组仅执行一次的问题分析与解决
问题原因
你当前的GroupBy是在全局数据流上执行的,Rx的GroupBy特性是复用已存在的分组Observable:第一次出现某个Channel时,会创建对应的分组并触发你订阅的Console.WriteLine逻辑;后续同一Channel的事件只会流入已创建的分组,不会再触发新的订阅回调。而你需要的是每次处理时间段数据时,输出该批次数据的所有分组,全局GroupBy显然不符合这个需求。
解决方案
把GroupBy移到SelectMany内部,针对每个时间段的事件流单独做分组处理,这样每个批次的事件都会独立生成分组并输出Channel:
修改后的代码示例(简洁版)
var observable = Observable .Interval(TimeSpan.FromSeconds(10)) .Scan(seed, NextPeriod) .SelectMany(period => { Console.WriteLine(); Console.WriteLine("fetching events for"); Console.WriteLine(period.Begin.ToString("yyyy-MM-dd hh:mm:ss.ffff")); Console.WriteLine(period.End.ToString("yyyy-MM-dd hh:mm:ss.ffff")); var events = FetchEvents(period, eventDataProvider); Console.WriteLine($"{events.Count} events fetched"); // 针对当前时间段的事件单独分组,提取分组Key输出 return events .GroupBy(x => x.Channel) .Select(group => group.Key); }) .Subscribe(channelKey => { Console.WriteLine($"Channel - {channelKey}"); });
如需处理分组内事件的版本
如果需要对每个分组内的事件做进一步处理,可以在每个时间段的逻辑内部订阅分组:
var observable = Observable .Interval(TimeSpan.FromSeconds(10)) .Scan(seed, NextPeriod) .Subscribe(period => { Console.WriteLine(); Console.WriteLine("fetching events for"); Console.WriteLine(period.Begin.ToString("yyyy-MM-dd hh:mm:ss.ffff")); Console.WriteLine(period.End.ToString("yyyy-MM-dd hh:mm:ss.ffff")); var events = FetchEvents(period, eventDataProvider); Console.WriteLine($"{events.Count} events fetched"); // 针对当前批次事件分组并处理 events .GroupBy(x => x.Channel) .Subscribe(group => { Console.WriteLine($"Channel - {group.Key}"); // 如需处理分组内的事件,可在此订阅group // group.Subscribe(@event => { /* 处理单个事件逻辑 */ }); }); });
效果说明
修改后,每次获取到时间段的事件数据时,都会对该批次的事件独立分组,所有存在的Channel都会触发一次Channel - xxx的输出,满足你持续运行时每次都打印分组的需求。
内容的提问来源于stack exchange,提问作者ironhide391
相关产品推荐
相关产品推荐

