[API Proposal]: Rx-like operators for IAsyncEnumerable<T>

Open
#114,269 3 comments 2 reactions 0 assignees View on GitHub

Nobody has claimed this yet.

Assessment

Difficulty
5/5
Estimated time
Over a week
Newbie friendliness
25/100
Issue type
Feature
Clarity
Mostly clear
Activity status
Stale
Tech stack
csharp

Research direction

Start by reviewing the current LINQ operators for IAsyncEnumerable and the corresponding dotnet-reactive operators referenced in the proposal. Evaluate the proposed Switch, Merge, and Concat APIs along with the open questions about Unit, GroupByLazy, Throttle, and Debounce. Done means reaching an agreed API scope and design, with implementation and tests defined.

Written by the indexing model from the issue text.

Description

api-suggestion area-System.Linq
Background and motivation

Current LINQ operators for IAsyncEnumerable<T> are the async transposition of the Synchronous LINQ over IEnumerable<T> and that is great. But IAsyncEnumerable<T> has a lot of similarities with IObservable<T> also, since its not blocking nature, and it has full language support (with await foreach).

I think its a great opportunity to pick some brilliant ideas (operators) from dotnet-reactive that are more oriented to event processing.

API Proposal

For example:


public static partial class AsyncEnumerable
{
    /// <summary>
    /// Flattens an asynchronous sequence of asynchronous sequences into a single sequence.
    /// Always switches to the latest inner sequence, canceling previous ones.
    /// </summary>
    /// <typeparam name="T">The type of elements in the sequence.</typeparam>
    /// <param name="sources">The source sequence of sequences.</param>
    /// <returns>A single flattened sequence of elements.</returns>
    public static IAsyncEnumerable<T> Switch<T>(this IAsyncEnumerable<IAsyncEnumerable<T>> sources);

    /// <summary>
    /// Merges multiple asynchronous sequences into a single sequence.
    /// Unlike Concat, elements are interleaved as they arrive.
    /// </summary>
    /// <typeparam name="T">The type of elements in the sequence.</typeparam>
    /// <param name="sources">The source sequences to merge.</param>
    /// <returns>A merged asynchronous sequence.</returns>
    public static IAsyncEnumerable<T> Merge<T>(this IAsyncEnumerable<IAsyncEnumerable<T>> sources);

    /// <summary>
    /// Concatenates multiple asynchronous sequences into a single sequence.
    /// Processes them sequentially: each sequence starts after the previous one completes.
    /// </summary>
    /// <typeparam name="T">The type of elements in the sequence.</typeparam>
    /// <param name="sources">The source sequences to concatenate.</param>
    /// <returns>A concatenated asynchronous sequence.</returns>
    public static IAsyncEnumerable<T> Concat<T>(this IAsyncEnumerable<IAsyncEnumerable<T>> sources);
}

  • It could be useful to process some IAsyncEnumerable<void> sequences. Rx solves this problem using an empty Unit struct. I dont know if this could be an acceptable solution.
  • Maybe could be useful a GroupByLazy?
  • Throttle and Debounce?
API Usage
async IAsyncEnumerable<IAsyncEnumerable<int>> GetDataSources()
{
    yield return GenerateNumbers(1); 
    await Task.Delay(3000); 
    yield return GenerateNumbers(100); 
}

async IAsyncEnumerable<int> GenerateNumbers(int start)
{
    for (int i = 0; i < 5; i++)
    {
        yield return start + i;
        await Task.Delay(1000); 
    }
}

await foreach (var number in GetDataSources().Switch())
{
    Console.WriteLine(number);
}

await foreach (var number in GetDataSources().Merge())
{
    Console.WriteLine(number);
}

await foreach (var number in GetDataSources().Concat())
{
    Console.WriteLine(number);
}

Alternative Designs

No response

Risks

No response

Dominant language
C#
Stars
18.3k
Forks
5.6k
PR merge metrics
PR metrics pending

Contributor guide

Open the contributing guide

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

More from dotnet/runtime

All issues in dotnet/runtime

Similar issues

More C# issues

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.