Partition an observable stream of IEnumerable<T> with a delay in Reactive Extensions

Viewed 461

I have the following Rx extension method for partitioning an IEnumerable<T> and delaying the producing of each partitioned value. It uses an IEnumerable<T> extension to partition the data, which is also shown with a unit test.

Is there a better way to do the 'delay' than using the Observable.Timer().Wait() method call?

public static class RxExtensions
{
    public static IObservable<IEnumerable<T>> PartitionWithInterval<T>(
        this IObservable<IEnumerable<T>> source, int size, TimeSpan interval,
        IScheduler scheduler = null)
    {
        if (scheduler == null)
        {
            scheduler = TaskPoolScheduler.Default;
        }

        var intervalEnabled = false;
        return source.SelectMany(x => x.Partition(size).ToObservable())
            .Window(1)
            .SelectMany(x =>
            {
                if (!intervalEnabled)
                {
                    intervalEnabled = true;
                }
                else
                {
                    Observable.Timer(interval, TaskPoolScheduler.Default).Wait();
                }

                return x;
            })
            .ObserveOn(scheduler);
    } 
}

public static class EnumerableExtensions
{
    public static IEnumerable<IEnumerable<T>> Partition<T>(
        this IEnumerable<T> source, int size)
    {
        using (var enumerator = source.GetEnumerator())
        {
            var items = new List<T>();
            while (enumerator.MoveNext())
            {
                items.Add(enumerator.Current);
                if (items.Count == size)
                {
                    yield return items.ToArray();

                    items.Clear();
                }
            }
           
            if (items.Any())
            {
                yield return items.ToArray();
            }
        }
    }
}

Test for the Rx extension method is shown below:

static void Main(string[] args)
{
     try
     {
         var data = Enumerable.Range(0, 10);
         var interval = TimeSpan.FromSeconds(1);

         Observable.Return(data)
            .PartitionWithInterval(2, interval)
            .Timestamp()
            .Subscribe(x =>
                {
                   var message = $"{x.Timestamp} - count = {x.Value.Count()}" +
                       $", values - {x.Value.First()}, {x.Value.Last()}";
                   Console.WriteLine(message);
                });

           Console.ReadLine();
       }
       catch (Exception e)
       {
           Console.WriteLine(e);
       }
}
3 Answers

Here is an implementation of the PartitionWithInterval operator, that is optimized towards memory efficiency. The enumerables emitted by the IObservable<IEnumerable<T>> are enumerated lazily, just enough to produce the next one or two partitions. Then their enumeration is suspended until the next interval. To achieve this laziness, the implementation uses IAsyncEnumerables instead of IObservables, and makes use of operators from the packages System.Linq.Async and System.Interactive.Async.

public static IObservable<IList<T>> PartitionWithInterval<T>(
    this IObservable<IEnumerable<T>> source, int size,
    TimeSpan interval, IScheduler scheduler = null)
{
    scheduler ??= Scheduler.Default;
    return Observable.Defer(() =>
    {
        Task delayTask = Task.CompletedTask;
        return source
            .ToAsyncEnumerable()
            .SelectMany(x => x.ToAsyncEnumerable()).Buffer(size) /* Behavior A */
            //.SelectMany(x => x.ToAsyncEnumerable().Buffer(size)) /* Behavior B */
            .Do(async (_, cancellationToken) =>
            {
                await delayTask;
                var timer = Observable.Timer(interval, scheduler);
                delayTask = timer.ToTask(cancellationToken);
            })
            .ToObservable();
    });
}

Below is a marble diagram that shows the behavior of the PartitionWithInterval operator, configured with size: 2:

Source: +----[1,2,3,4,5]--------------------[6,7,8,9]---|
Output: +----[1,2]-------[3,4]--------------[5,6]-------[7,8]-------[9]|

As shown, an output partition may contain values from more than one enumerables (the partition [5,6] in the above diagram). In case this is undesirable, just comment the line "Behavior A" and uncomment the line "Behavior B". The marble diagram below shows the effect of this change:

Source: +----[1,2,3,4,5]--------------------[6,7,8,9]---|
Output: +----[1,2]-------[3,4]-------[5]-------[6,7]-------[8,9]|

Note: The above solution is not absolutely satisfactory regarding the intention to enumerate lazily the enumerables emitted by the source observable. The ideal would be to produce each partition exactly at the time it should be emitted. Instead, the above implementation gathers the elements of the next partition immediately after emitting the previous partition. The alternative would be to enforce a delay after emitting each partition, including the last one. This would postpone the completion of the resulting IObservable by a time span equal to interval, which is not ideal either (this behavior is implemented by the revision 3 of this answer). The ideal behavior could probably be achieved by re-implementing the operators ToAsyncEnumerable, SelectMany, Buffer and Do, so that they communicate the state IsLast of the currently emitted element. Even if this is possible, it would require a lot of effort for such a measly improvement.

Related