You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.23 23:13:30