生产者-消费者模式

  • ~3.05K 字
  1. 1. BlockingCollection
    1. 1.1. 常用方法
    2. 1.2. CancellationTokenSource
  2. 2. Channel
    1. 2.1. 常用方法
  3. 3. 背压问题

生产者-消费者模式是一种多线程协作模式:一个线程负责产生数据把数据加到队列里面,另一些线程负责处理数据将队列里的数据拿出来消耗掉。

BlockingCollection

System.Collections.Concurrent.BlockingCollection<T> 是.NET提供的一个线程安全的基于“同步阻塞”的生产者-消费者集合,它可以在多线程环境下安全地存放数据(普通集合存在数据竞争等问题需要加锁),并支持阻塞等待。

常用方法

  1. void Add() 生产者往集合添加数据,如果容量满,会阻塞

  2. bool TryAdd() 生产者尝试添加数据,不会阻塞

  3. void Take() 消费者从集合取数据

    1
    2
    var queue = new BlockingCollection<byte[]>();
    var item = queue.Take();//集合没有数据时会阻塞
  4. bool TryTake() 消费者尝试取数据,不会阻塞

  5. IEnumerable<T> GetConsumingEnumerable(CancellationToken cancellationToken) 返回一个阻塞式可枚举集合,消费者可以通过 foreach 持续读取数据,直到集合被标记为完成(CompleteAdding())或取消(CancellationTokenSource.Cancel())

  6. 构造函数 BlockingCollection<T>() 会构造出一个容量无限的集合,而 BlockingCollection<T>(int boundedCapacity) 构造出一个容量为boundedCapacity的有界集合

  7. void CompleteAdding() 此方法通知消费者:生产者已经不会再添加数据了

    为什么需要这个方法呢?因为在使用 GetConsumingEnumerable() 获取数据时,如果没有数据且生产者未调用 CompleteAdding() ,消费者不知道什么时候结束,它会一直阻塞下去,使用此方法告诉消费者:后面不会再有新数据了,可以在队列消费完后退出。

CancellationTokenSource

System.Threading.System.Threading.CancellationTokenSource 我们要知道一个线程不能被外部强行杀死,而是通过发送一个“取消信号”,让任务自己决定什么时候停止。而 new CancellationTokenSource() 会创建一个负责发出取消信号的对象,它是取消操作的控制器,通过 CancellationTokenSource.Token 拿到令牌传给需要的线程,如 GetConsumingEnumerable(CancellationToken cancellationToken) 这样我们在外部使用CancellationTokenSource.Cancel()时就会发送信号给接收 CancellationToken 的线程告诉它我想停止任务,这时线程接收到信号后就会自己停止。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
// CancellationTokenSource 简单示例
CancellationTokenSource cts = new();
CancellationToken token = cts.Token;
// 注册取消时自动执行的代码
token.Register(() =>
{
Console.WriteLine("任务被取消");
});

Task.Run(() =>
{
while(true)
{
// IsCancellationRequested true 已经请求取消
if(token.IsCancellationRequested)
{
Console.WriteLine("收到取消请求");
break;
}

Console.WriteLine("运行中...");

Thread.Sleep(100);
}

}, token);

Thread.Sleep(3000);
// 修改IsCancellationRequested=true
cts.Cancel();

Channel

System.Threading.Channels.Channel<T> 是.NET提供的一个线程安全的基于异步的生产者-消费者集合。适合在异步编程场景下传递数据。

常用方法

创建 System.Threading.Channels.Channel<T> 主要通过 System.Threading.Channels.Channel 类的静态工厂方法

  1. static Channel<T> CreateUnbounded<T>() 创建一个容量无上限的无界通道,生产者总是可以立即写入数据,适用于高吞吐、消费者能跟上生产者速度的场景,如果消费者处理速度跟不上,通道积压数据,会导致内存压力。
  2. static Channel<T> CreateBounded<T>(int capacity) 创建有固定容量 capacity 的有界通道
  3. static Channel<T> CreateBounded<T>(BoundedChannelOptions options) 创建有固定容量 options.Capacity 的有界通道,通道写满时,会触发“背压”机制,通 BoundedChannelFullMode options.FullMode 来决定如何处理新数据

生产者( ChannelWriter<T> )API,通过 Channel<T>.Writer 获取

  1. WriteAsync(T item): 异步写入,会等待空间。
  2. TryWrite(T item): 尝试立即写入。
  3. Complete(): 通知通道不再有新数据。
  4. TryComplete(): 尝试标记完成。
  5. WaitToWriteAsync(): 等待直到通道有空间可用。

消费者( ChannelReader<T> )API,通过 Channel<T>.Reader 获取

  1. ReadAllAsync(): 返回 IAsyncEnumerable,方便用 await foreach 遍历所有数据直到完成。
  2. ReadAsync(): 异步读取单个元素。
  3. TryRead(out T item): 尝试立即读取,成功返回 true。
  4. WaitToReadAsync(): 等待直到通道有数据可读。
  5. TryPeek(out T item): 尝试偷看下一个元素,但不移除它。

背压问题

背压(Backpressure)问题