【问题标题】:Partitioning observables in C#在 C# 中对 observables 进行分区
【发布时间】:2017-09-30 17:11:49
【问题描述】:

我正在寻找某种方法将可观察序列拆分为单独的序列,我可以根据给定的谓词独立处理这些序列。像这样的东西是理想的:

var (evens, odds) = observable.Partition(x => x % 2 == 0);
var strings = evens.Select(x => x.ToString());
var floats = odds.Select(x => x / 2.0);

我能想到的最接近的方法是做两个 where 过滤器,但这需要评估条件并处理源序列两次,我对此并不感兴趣。

observable = observable.Publish().RefCount();
var strings = observable.Where(x => x % 2 == 0).Select(x => x.ToString());
var floats = observable.Where(x => x % 2 != 0).Select(x => x / 2.0);

F# 似乎通过 Observable.partition<'T>Observable.split<'T,'U1,'U2> 对此提供了很好的支持,但我找不到与 C# 等效的任何东西。

【问题讨论】:

  • 您可以随时拉入 F# 库并从 C# 中使用它
  • 查看 F# 源代码,看起来它实际上只是对源流应用了两个过滤器,所以它与我的提议基本相同,有两个 wheres
  • 两个链接都死了,找不到合适的替换链接。

标签: c# system.reactive


【解决方案1】:

GroupBy 可能会删除“观察两次”限制,但您仍然会得到 Where 子句:

public static class X
{
    public static (IObservable<T> trues, IObservable<T> falsies) Partition<T>(this IObservable<T> source, Func<T, bool> partitioner)
    {
        var x = source.GroupBy(partitioner).Publish().RefCount();
        var trues = x.Where(g => g.Key == true).Merge();
        var falsies = x.Where(g => g.Key == false).Merge();
        return (trues, falsies);
    }
}

【讨论】:

  • 不错。我想我也许可以用GroupBy 做点什么,但Merge 是我缺少的步骤。谢谢。
  • 哦,实际上,当您这样做时,看起来源序列确实被评估了两次。该死的
  • 糟糕,我太懒了。更新的代码:Var x 应该是 Published+Refcounted。优点仍然是,如果您的过滤器功能很昂贵,则只运行一次。
  • 这个解决方案的问题是订阅trues 会自动加热source observable,这使得IGroupedObservable&lt;false, T&gt; 可能在订阅falsies observable 之前被发射。在这种情况下,falsies 永远不会发出任何元素。您可以通过将.Publish() 替换为.Replay() 来缓解此问题,但即便如此,您仍可能会丢失一些falsies,它们可能会在订阅trues 期间同步发出。
【解决方案2】:

这样的事情怎么样

var (odds,evens) = (collection.Where(a=> a % 2 == 1), collection.Where(a=> a % 2 == 0));?

或者如果你想根据一个条件进行分区

Func<int,bool> predicate = a => a%2==0;

var (odds,evens) = (collection.Where(a=> !predicate(a)), collection.Where(a=> predicate(a)));

我认为没有办法解决这样一个事实,即您以这种方式迭代项目两次,其他可以做的就是拥有一个接受谓词并传入 2 个单独集合并在一次迭代中填充它们的方法一个 foreach 或 for。

类似这样的:

var collection = new[] { 1, 2, 3, 4, 5, 6, 7, 8, 9};

Func<int,bool> predicate = a => a%2==0;
var odds = new List<int>();
var evens = new List<int>();

Action<List<int>, List<int>, Func<int, bool>> partition = (collection1, collection2, pred) =>
{
    foreach (int element in collection)
    {
        if (pred(element))
        {
            collection1.Add(element);
        }
        else 
        {
            collection2.Add(element);
        }
    }
};

partition(evens, odds, predicate);

扩展最后一个想法,你在寻找这样的东西吗?

public static (ObservableCollection<T>, ObservableCollection<T>) Partition<T>(this ObservableCollection<T> collection, Func<T, bool> predicate)
{
    var collection1 = new ObservableCollection<T>();
    var collection2 = new ObservableCollection<T>();

    foreach (T element in collection)
    {
        if (predicate(element))
        {
            collection1.Add(element);
        }
        else
        {
            collection2.Add(element);
        }
    }

    return (collection1, collection2);
}

【讨论】:

  • 感谢您的想法。您对ObservableCollection 的建议是正确的,但我正在处理可观察序列(IObservable)而不是集合,因此它需要立即返回新序列,因为它们可能是无限序列。我做了一些东西,创建了一对Subjects,但它有点大,我希望有一些预先存在的东西。
  • 嗯,现在没有想到会在元组中实时“产生”,也许不是在方法中创建它们,而是传递 IObservables 并在这些序列发生变化时拦截调用。
【解决方案3】:

使用RefCount 操作符加热源序列不是一个好主意,因为源序列可能在派生序列的所有订阅到位之前就开始发射元素。在这种情况下,一些发射的元素可能会丢失。一种更安全的方法是推迟预热源序列,直到所有观察者都已订阅。这是一个如何做的例子:

var published = observable.Publish(); // Make sure not to warm it too early

var strings = published.Where(x => x % 2 == 0).Select(x => x.ToString());
var floats = published.Where(x => x % 2 != 0).Select(x => x / 2.0);

strings.Subscribe(x => Console.WriteLine(x));
floats.Subscribe(x => Console.WriteLine(x));

published.Connect(); // Now that all subscriptions are in place, it's time to warm it

await published; // Wait for the completion of the source sequence

您可以通过使用LookupObservable&lt;TSource, TKey&gt; 类(包含在an answer of a relevant question 中)来减少上述代码的重复性。实现这个类是因为创建多个Where 子序列可能非常低效,以防子序列的总数很大(因为将检查源发出的每个元素的许多条件)。在您的情况下,您只有两个子序列,一个用于键 true,另一个用于键 false,因此使用 LookupObservable 类不太引人注目。无论如何,这是一个用法示例:

var published = observable.Publish(); // Make sure not to warm it too early

var lookup = new LookupObservable<int, bool>(published, x => x % 2 == 0);

var strings = lookup[true].Select(x => x.ToString());
var floats = lookup[false].Select(x => x / 2.0);

strings.Subscribe(x => Console.WriteLine(x));
floats.Subscribe(x => Console.WriteLine(x));

published.Connect(); // Now that all subscriptions are in place, it's time to warm it

await published; // Wait for the completion of the source sequence

【讨论】:

    猜你喜欢
    • 2011-12-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-02-16
    • 1970-01-01
    • 2018-03-17
    • 1970-01-01
    相关资源
    最近更新 更多