Rx.Net实现:触发关闭序列时聚合消息并输出中间结果
Rx.NET 按Key聚合并按需刷新的实现方案
需求目标
实现对(int Key, int Value)类型的消息序列按Key求和聚合:
- 当
flushObservable发出标记时,输出当前所有Key的聚合结果 - 当原消息序列完成时,输出最后一批聚合结果
示例输入序列:(1,1) → (2,3) → (1,3) → [Flush] → (2,1) → (2,5) → [Complete]
期望输出:
- 触发Flush时:
[(1,4), (2,3)] - 序列完成时:
[(2,6)]
已尝试方案及问题
方案1:GroupBy + Aggregate + Buffer(flush)
这是直观的组合逻辑,但无法输出中间聚合结果,仅在序列最终完成时一次性输出所有值。
代码实现:
using System.Reactive; using System.Reactive.Linq; using System.Reactive.Subjects; using System.Reactive.Threading.Tasks; var source = new Subject<(int Key, int Value)>(); var flush = new Subject<Unit>(); var completion = StartProcessing(source, flush); source.OnNext((1, 1)); source.OnNext((2, 3)); source.OnNext((1, 3)); flush.OnNext(Unit.Default); // 期望输出 [(1,4), (2,3)] source.OnNext((2, 1)); source.OnNext((2, 5)); source.OnCompleted(); // 期望输出 [(2,6)] await completion; return; static Task StartProcessing( IObservable<(int Key, int Value)> source, IObservable<Unit> flush) { return source .GroupBy(message => message.Key) .SelectMany(group => group .Aggregate( seed: (group.Key, Value: 0), accumulator: (output, item) => (output.Key, output.Value + item.Value) ) ) .Buffer(flush) .Select(buffer => Observable.FromAsync(() => Flush(buffer))) .Merge() .ToTask(); } static async Task Flush(IEnumerable<(int Key, int Value)> data) { Console.WriteLine($"Flushing [{string.Join(", ", data)}]"); // 模拟处理耗时 await Task.Delay(TimeSpan.FromSeconds(1)); }
实际输出:
Flushing [] Flushing [(1, 4), (2, 9)]
问题原因:Aggregate仅在分组序列完成时才推送最终聚合值,无法在flush触发时捕获中间聚合状态,导致Buffer无法正确收集批次结果。
方案2:GroupByUntil
用GroupByUntil替代GroupBy和Buffer,可以在flush触发时关闭分组并输出单个Key的聚合结果,但每个Key的结果会单独输出,无法批量打包刷新。
实际输出:
Flushing (1, 4) Flushing (2, 3) Flushing (2, 6)
可行解决方案
通过Window运算符实现按批次分割序列并批量聚合,代码如下:
static Task StartProcessing( IObservable<(int Key, int Value)> source, IObservable<Unit> flush) { return source .Window(flush) .SelectMany(group => group .Aggregate( seed: new Dictionary<int, int>(), accumulator: (output, item) => { output[item.Key] = output.TryGetValue(item.Key, out var value) ? value + item.Value : item.Value; return output; } ) .SelectMany(buffer => Observable .FromAsync(() => Flush(buffer.Select(x => (x.Key, x.Value))))) ) .ToTask(); }
方案解释
- Window:将原序列分割为多个子序列,每个子序列在
flush触发或原序列完成时结束,支持多批次的序列分割。 - Aggregate:对每个子序列内的消息按Key求和,子序列结束时输出包含所有Key聚合结果的字典。
- Observable.FromAsync:将异步刷新处理包装为Observable,确保处理完成后才标记该批次流程结束。
- SelectMany:将多个子序列的处理Observable扁平化为单一序列,替代
Select+Merge的组合逻辑。 - ToTask:返回一个Task,当所有刷新处理完成时结束,代表整个聚合流程完成。
内容的提问来源于stack exchange,提问作者Ivan Zaruba
相关产品推荐
相关产品推荐

