具有Reactive Extensions的组的缓冲区组,嵌套订阅

Ron*_*erg 4 .net c# system.reactive reactive

我有一个事件源,生成属于某些组的事件.我想缓冲这些组并将组(分批)发送到存储.到目前为止我有这个:

eventSource
    .GroupBy(event => event.GroupingKey)
    .Select(group => new { group.Key, Events = group })
    .Subscribe(group => group.Events
                            .Buffer(TimeSpan.FromSeconds(60), 100)
                            .Subscribe(list => SendToStorage(list)));
Run Code Online (Sandbox Code Playgroud)

所以有一个嵌套订阅组中的事件.不知怎的,我认为有更好的方法,但我还没有弄清楚.

Shl*_*omo 9

这是解决方案:

eventSource
    .GroupBy(e => e.GroupingKey)
    .SelectMany(group => group.Buffer(TimeSpan.FromSeconds(60), 100))
    .Subscribe(list => SendToStorage(list));
Run Code Online (Sandbox Code Playgroud)

这里有一些可以帮助你'减少'的一般规则:

1)嵌套订阅通常使用Select嵌套订阅之后的所有内容进行修复,然后是a Merge,然后是嵌套订阅.所以应用它,你得到这个:

eventSource
    .GroupBy(e => e.GroupingKey)
    .Select(group => new { group.Key, Events = group })
    .Select(group => group.Events.Buffer(TimeSpan.FromSeconds(60), 100)) //outer subscription selector
    .Merge()
    .Subscribe(list => SendToStorage(list));
Run Code Online (Sandbox Code Playgroud)

2)你显然可以组合两个连续的选择(因为你没有对匿名对象做任何事情,可以删除它):

eventSource
    .GroupBy(e => e.GroupingKey)
    .Select(group => group.Buffer(TimeSpan.FromSeconds(60), 100)) 
    .Merge()
    .Subscribe(list => SendToStorage(list));
Run Code Online (Sandbox Code Playgroud)

3)最后,a Select后面的a Merge可以简化为SelectMany:

eventSource
    .GroupBy(e => e.GroupingKey)
    .SelectMany(group => group.Buffer(TimeSpan.FromSeconds(60), 100))
    .Subscribe(list => SendToStorage(list));
Run Code Online (Sandbox Code Playgroud)