生产者-消费者模式是一种多线程协作模式:一个线程负责产生数据把数据加到队列里面,另一些线程负责处理数据将队列里的数据拿出来消耗掉。
BlockingCollection
System.Collections.Concurrent.BlockingCollection<T> 是.NET提供的一个线程安全的基于“同步阻塞”的生产者-消费者集合,它可以在多线程环境下安全地存放数据(普通集合存在数据竞争等问题需要加锁),并支持阻塞等待。
常用方法
void Add() 生产者往集合添加数据,如果容量满,会阻塞
bool TryAdd() 生产者尝试添加数据,不会阻塞
void Take() 消费者从集合取数据
1
2var queue = new BlockingCollection<byte[]>();
var item = queue.Take();//集合没有数据时会阻塞bool TryTake() 消费者尝试取数据,不会阻塞
IEnumerable<T> GetConsumingEnumerable(CancellationToken cancellationToken) 返回一个阻塞式可枚举集合,消费者可以通过 foreach 持续读取数据,直到集合被标记为完成(CompleteAdding())或取消(CancellationTokenSource.Cancel())
构造函数 BlockingCollection<T>() 会构造出一个容量无限的集合,而 BlockingCollection<T>(int boundedCapacity) 构造出一个容量为boundedCapacity的有界集合
void CompleteAdding() 此方法通知消费者:生产者已经不会再添加数据了
为什么需要这个方法呢?因为在使用 GetConsumingEnumerable() 获取数据时,如果没有数据且生产者未调用 CompleteAdding() ,消费者不知道什么时候结束,它会一直阻塞下去,使用此方法告诉消费者:后面不会再有新数据了,可以在队列消费完后退出。
CancellationTokenSource
System.Threading.System.Threading.CancellationTokenSource 我们要知道一个线程不能被外部强行杀死,而是通过发送一个“取消信号”,让任务自己决定什么时候停止。而 new CancellationTokenSource() 会创建一个负责发出取消信号的对象,它是取消操作的控制器,通过 CancellationTokenSource.Token 拿到令牌传给需要的线程,如 GetConsumingEnumerable(CancellationToken cancellationToken) 这样我们在外部使用CancellationTokenSource.Cancel()时就会发送信号给接收 CancellationToken 的线程告诉它我想停止任务,这时线程接收到信号后就会自己停止。
1 | // CancellationTokenSource 简单示例 |
Channel
System.Threading.Channels.Channel<T> 是.NET提供的一个线程安全的基于异步的生产者-消费者集合。适合在异步编程场景下传递数据。
常用方法
创建 System.Threading.Channels.Channel<T> 主要通过 System.Threading.Channels.Channel 类的静态工厂方法
- static Channel<T> CreateUnbounded<T>() 创建一个容量无上限的无界通道,生产者总是可以立即写入数据,适用于高吞吐、消费者能跟上生产者速度的场景,如果消费者处理速度跟不上,通道积压数据,会导致内存压力。
- static Channel<T> CreateBounded<T>(int capacity) 创建有固定容量 capacity 的有界通道
- static Channel<T> CreateBounded<T>(BoundedChannelOptions options) 创建有固定容量 options.Capacity 的有界通道,通道写满时,会触发“背压”机制,通 BoundedChannelFullMode options.FullMode 来决定如何处理新数据
生产者( ChannelWriter<T> )API,通过 Channel<T>.Writer 获取
- WriteAsync(T item): 异步写入,会等待空间。
- TryWrite(T item): 尝试立即写入。
- Complete(): 通知通道不再有新数据。
- TryComplete(): 尝试标记完成。
- WaitToWriteAsync(): 等待直到通道有空间可用。
消费者( ChannelReader<T> )API,通过 Channel<T>.Reader 获取
- ReadAllAsync(): 返回 IAsyncEnumerable
,方便用 await foreach 遍历所有数据直到完成。 - ReadAsync(): 异步读取单个元素。
- TryRead(out T item): 尝试立即读取,成功返回 true。
- WaitToReadAsync(): 等待直到通道有数据可读。
- TryPeek(out T item): 尝试偷看下一个元素,但不移除它。
背压问题
背压(Backpressure)问题