C#生产者/消费者场景下IAsyncEnumerable的缓冲实现方案咨询
C# 实现带缓冲的异步生产者-消费者模式(类似Go通道)
要解决数据库锁持有时间过长的问题,同时避免一次性加载所有数据到内存,C# 可以通过异步通道/队列实现类似Go带缓冲通道的生产者-消费者模式,核心是让生产者提前读取数据到缓冲队列,消费者从队列异步处理,平衡内存占用与锁持有时长。
推荐实现:使用官方 Channel<T>(.NET Core 3.0+)
System.Threading.Channels.Channel<T> 是.NET官方提供的异步通道组件,专门用于高效的生产者-消费者场景,支持配置缓冲大小,自动处理缓冲满时的等待逻辑,比手动实现队列更简洁可靠。
完整代码示例
using System.Threading.Channels; // 配置缓冲大小,根据内存和锁需求调整 var bufferCapacity = 10; var channel = Channel.CreateBounded<MyData>(new BoundedChannelOptions(bufferCapacity) { // 缓冲满时,生产者等待消费者取走数据(和Go带缓冲通道行为一致) FullMode = BoundedChannelFullMode.Wait }); // 生产者任务:异步读取数据库并写入通道 var producer = Task.Run(async () => { try { await foreach (var data in this.dataSource.Read(query)) { // 写入通道,缓冲满时自动暂停等待 await channel.Writer.WriteAsync(data); } } finally { // 标记通道写入完成,消费者会知道没有更多数据 channel.Writer.Complete(); } }); // 消费者任务:从通道读取数据并异步处理 var consumer = Task.Run(async () => { await foreach (var data in channel.Reader.ReadAllAsync()) { await this.consumer.Write(data); } }); // 等待生产和消费任务全部完成,处理异常 try { await Task.WhenAll(producer, consumer); } catch (Exception ex) { // 按需处理异常(数据库读取失败、消费者处理错误等) Console.WriteLine($"执行出错: {ex.Message}"); }
关键细节说明
- 缓冲控制:通过
BoundedChannelOptions设置缓冲容量,FullMode可选不同策略:Wait(生产者等待)、DropWrite(丢弃当前写入)、DropOldest(丢弃最旧数据)等,按需选择。 - 锁持有优化:数据库读取的
await foreach仅在通道有剩余空间时继续读取,不会因为消费者处理慢而一直持有数据库锁,锁的持有时间仅为单条数据的读取耗时。 - 自动结束通知:生产者完成后调用
Writer.Complete(),消费者通过ReadAllAsync()的枚举结束自动感知,无需手动处理结束信号。
手动实现:基于队列与信号量
如果无法使用Channel<T>,也可以手动用ConcurrentQueue加SemaphoreSlim实现类似逻辑,但需要自行处理更多细节:
using System.Collections.Concurrent; var bufferSize = 10; var dataQueue = new ConcurrentQueue<MyData>(); var semaphore = new SemaphoreSlim(bufferSize, bufferSize); bool producerCompleted = false; // 生产者 var producerTask = Task.Run(async () => { try { await foreach (var data in this.dataSource.Read(query)) { await semaphore.WaitAsync(); dataQueue.Enqueue(data); } } finally { producerCompleted = true; // 释放一个信号量,让消费者退出循环 semaphore.Release(); } }); // 消费者 var consumerTask = Task.Run(async () => { while (!producerCompleted || dataQueue.TryDequeue(out var data)) { if (data != null) { await this.consumer.Write(data); semaphore.Release(); } else { // 队列为空时等待生产者写入 await semaphore.WaitAsync(); } } }); await Task.WhenAll(producerTask, consumerTask);
这种方式需要手动维护信号量和完成标记,代码复杂度更高,优先推荐使用官方Channel<T>。
内容的提问来源于stack exchange,提问作者Rafael
相关产品推荐
相关产品推荐

