在生产者/消费者场景中缓冲 IAsyncEnumerable

Raf*_*ael 4 c# linq producer-consumer async-await iasyncenumerable

我有一个场景,我正在从数据库中读取一些数据。该数据以 的形式返回IAsyncEnumerable<MyData>。读取数据后我想将其发送给消费者。该消费者是异步的。现在我的代码看起来像这样:

// C#
IAsyncEnumerable<MyData> enumerable = this.dataSource.Read(query);

await foreach (var data in enumerable) 
{
    await this.consumer.Write(data);
}
Run Code Online (Sandbox Code Playgroud)

我的问题是,当我枚举数据库时,我持有数据的锁。我不想持有这把锁的时间超过我需要的时间。

如果消费者消耗数据的速度比生产者生成数据的速度慢,有什么方法可以让我急切地从数据源中读取数据,而无需仅调用ToListor ToListAsync。我想避免一次将所有数据读入内存,如果现在生产者比消费者慢,这会导致相反的问题。如果数据库上的锁不是越短越好,我想要在内存中一次有多少数据与我们保持枚举运行多长时间之间进行可配置的权衡。

我的想法是,有某种方法可以使用队列或类似通道的数据结构来充当生产者和消费者之间的缓冲区。

在 Golang 中我会做这样的事情:

// go
queue := make(chan MyData, BUFFER_SIZE)
go dataSource.Read(query, queue)

// Read sends data on the channel, closes it when done

for data := range queue {
    consumer.Write(data)
}
Run Code Online (Sandbox Code Playgroud)

有没有办法在 C# 中获得类似的行为?

The*_*ias 5

这是Rafael答案ConsumeBuffered中扩展方法的更强大的实现。这个使用 a作为缓冲区,而不是. 优点是,源序列和缓冲序列这两个序列的枚举不会各阻塞一个线程。已注意完成源序列的枚举,以防缓冲序列的枚举被下游消费者过早地放弃。Channel<T>BlockingCollection<T>

public static async IAsyncEnumerable<T> ConsumeBuffered<T>(
    this IAsyncEnumerable<T> source, int capacity,
    [EnumeratorCancellation] CancellationToken cancellationToken = default)
{
    ArgumentNullException.ThrowIfNull(source);
    Channel<T> channel = Channel.CreateBounded<T>(new BoundedChannelOptions(capacity)
    {
        SingleWriter = true,
        SingleReader = true,
    });
    using CancellationTokenSource completionCts = new();

    Task producer = Task.Run(async () =>
    {
        try
        {
            await foreach (T item in source.WithCancellation(completionCts.Token)
                .ConfigureAwait(false))
            {
                await channel.Writer.WriteAsync(item).ConfigureAwait(false);
            }
        }
        catch (ChannelClosedException) { } // Ignore
        finally { channel.Writer.TryComplete(); }
    });

    try
    {
        await foreach (T item in channel.Reader.ReadAllAsync(cancellationToken)
            .ConfigureAwait(false))
        {
            yield return item;
            cancellationToken.ThrowIfCancellationRequested();
        }
        await producer.ConfigureAwait(false); // Propagate possible source error
    }
    finally
    {
        // Prevent fire-and-forget in case the enumeration is abandoned
        if (!producer.IsCompleted)
        {
            completionCts.Cancel();
            channel.Writer.TryComplete();
            await Task.WhenAny(producer).ConfigureAwait(false);
        }
    }
}
Run Code Online (Sandbox Code Playgroud)

设置有界通道的SingleWriter和SingleReader选项有点学术性,可以省略。目前 (.NET 6) ,无论提供什么选项,System.Threading.Channels库中都只有一种有界Channel<T>实现。此实现基于与 .NET 同步的(与 类似的内部 .NET 类型)。Deque<T>Queue<T>lock

通道在try/块内枚举,因为当枚举被放弃时,finallyC# 迭代器将finally块作为自动生成的Dispose/方法的一部分执行。DisposeAsyncIEnumerator<T>IAsyncEnumerator<T>

注意:如果外部CancellationToken被取消,取消将作为 传播OperationCanceledException,并且所有缓冲的项目都会丢失。在具有多个生产者和消费者的生产者-消费者场景中,这可能是一个问题。建议CancellationToken仅用于破坏整个处理管道,而不是部分处理管道。