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

Rx.Net实现:触发关闭序列时聚合消息并输出中间结果

Rx.NET 按Key聚合并按需刷新的实现方案

需求目标

实现对(int Key, int Value)类型的消息序列按Key求和聚合:

  • 当flush Observable发出标记时,输出当前所有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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 07:31:27