Partitioning observables in C#

Viewed 267

I'm looking for some way of splitting an observable sequence into separate sequences that I can process independently based on a given predicate. Something like this would be ideal:

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

The closest I've been able to come up with is doing two where filters, but that requires evaluating the condition and processing the source sequence twice, which I'm not wild about.

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# seems to have good support for this with Observable.partition<'T> and Observable.split<'T,'U1,'U2>, but I've not been able to find anything equivalent for C#.

3 Answers

A GroupBy may remove the "observe twice" restriction, though you'll still end up with Where clauses:

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);
    }
}

How about something like

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

or if you want to partition based on one condition

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

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

I think there is no working around the fact that you iterate the items twice this way, what else could be done would be to have a method that accepts a predicate and pass in the 2 sepatate collections and populate them in one iteration in a foreach or for.

Something like this:

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);

Expanding on the last idea, are you looking for something like this?

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);
}

Warming the source sequence by using the RefCount operator is not a good idea, because the source sequence may start emitting elements before all subscriptions to derived sequences are in place. In that case some of the emitted elements could be lost. A safer approach is to postpone warming the source sequence until all observers have been subscribed. Here is an example of how to do it:

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

You could make the above code a bit less repetitive by using the LookupObservable<TSource, TKey> class, that is included in an answer of a relevant question. This class was implemented because creating multiple Where subsequences can be quite inefficient, in case the total number of subsequences is large (because each element emitted by the source will be checked for numerous conditions). In your case you have only two subsequences, one for the key true and one for the key false, so using the LookupObservable class is less compelling. In any case, here is a usage example:

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
Related