如何实现父对象填充List、子对象异步消费List的模型
这个场景其实就是典型的生产者-消费者模式,刚好适合用.NET里的线程安全集合和异步处理来实现,我给你捋捋具体的实现思路和代码示例:
核心思路:生产者-消费者模式
你的需求完全匹配生产者-消费者模型:父对象是生产者,负责往队列里添加元素;子对象是消费者,持续从队列里取元素处理,空队列时进入等待状态,不用空转浪费资源。
关键技术点
要实现这个异步、线程安全的模型,需要用到这几个核心组件:
- 线程安全集合:比如
ConcurrentQueue<T>,替代普通List<T>,避免多线程操作时的并发冲突(普通List不是线程安全的,直接在多线程下用会出问题)。 - 信号量(SemaphoreSlim):用来通知消费者有新元素到来,避免消费者一直循环检查队列(空转),大幅降低CPU占用。
- 异步循环:消费者用
async/await实现持续运行的处理逻辑,非阻塞,不影响其他操作。
完整代码示例
下面是一个可直接运行的C#示例,我会逐部分解释:
1. 子对象(消费者)的实现
using System.Collections.Concurrent; using System.Threading; using System.Threading.Tasks; public class ElementProcessor<T> { // 线程安全的队列,存储待处理元素 private readonly ConcurrentQueue<T> _elementQueue = new ConcurrentQueue<T>(); // 信号量:用来通知有新元素需要处理,初始为0(没有元素时等待) private readonly SemaphoreSlim _newElementSignal = new SemaphoreSlim(0); // 取消令牌,用于优雅停止处理循环(可选,如果你需要停止子对象的话) private readonly CancellationTokenSource _cts = new CancellationTokenSource(); // 存储处理循环的任务,方便外部跟踪(可选) public Task ProcessingTask { get; } public ElementProcessor() { // 启动异步处理循环 ProcessingTask = RunProcessingLoopAsync(); } // 父对象调用这个方法添加元素 public void AddElement(T element) { _elementQueue.Enqueue(element); // 释放信号量,通知处理循环有新元素 _newElementSignal.Release(); } // 优雅停止处理循环的方法(可选) public async Task StopProcessingAsync() { _cts.Cancel(); // 释放一次信号量,让等待中的循环退出 _newElementSignal.Release(); await ProcessingTask; _newElementSignal.Dispose(); _cts.Dispose(); } // 核心处理循环 private async Task RunProcessingLoopAsync() { try { while (!_cts.Token.IsCancellationRequested) { // 等待信号量:如果队列空,这里会阻塞,直到有新元素或者取消 await _newElementSignal.WaitAsync(_cts.Token); // 循环取出所有可用元素处理(避免每次只处理一个,提高效率) while (_elementQueue.TryDequeue(out var element)) { // 这里是你的元素处理逻辑,换成你实际需要的操作 await ProcessElementAsync(element); } } } catch (OperationCanceledException) { // 取消时正常退出,不需要处理异常 } } // 模拟异步处理元素的方法,你可以替换成自己的业务逻辑 private async Task ProcessElementAsync(T element) { // 模拟处理耗时,比如调用API、读写文件等 await Task.Delay(1000); Console.WriteLine($"处理完成元素: {element}"); } }
2. 父对象(生产者)的使用示例
class Program { static async Task Main(string[] args) { // 创建子对象(消费者) var processor = new ElementProcessor<string>(); // 父对象在不同时间添加元素,不用等待处理完成 Console.WriteLine("父对象添加元素: 元素1"); processor.AddElement("元素1"); await Task.Delay(2000); // 模拟2秒后添加下一个元素 Console.WriteLine("父对象添加元素: 元素2"); processor.AddElement("元素2"); await Task.Delay(1500); // 模拟1.5秒后添加下一个元素 Console.WriteLine("父对象添加元素: 元素3"); processor.AddElement("元素3"); // 等待所有元素处理完成,然后停止处理器 await Task.Delay(5000); await processor.StopProcessingAsync(); Console.WriteLine("处理循环已停止"); } }
代码解释
- 线程安全队列:
ConcurrentQueue<T>自动处理多线程下的入队/出队操作,不需要手动加锁,比lock包裹普通List更高效。 - 信号量的作用:
SemaphoreSlim(0)初始没有可用信号,当父对象添加元素后调用Release(),处理循环的WaitAsync()就会被唤醒,开始处理元素。处理完所有元素后,循环又会回到等待状态,不会空转。 - 异步处理:
ProcessElementAsync是异步方法,处理元素时不会阻塞整个循环,即使某个元素处理耗时较长,后续元素也能在它完成后继续处理(如果需要并行处理,可以调整逻辑,但这里是按顺序处理)。 - 优雅停止:
StopProcessingAsync方法通过取消令牌和释放信号量,让处理循环安全退出,避免线程泄漏。
扩展建议
- 如果需要并行处理元素:可以在
ProcessElementAsync里用Task.WhenAll或者限制并发数(比如用SemaphoreSlim控制最大并发数)。 - 如果需要批量处理:可以在队列积累到一定数量后再一次性处理,调整
RunProcessingLoopAsync里的逻辑即可。 - 如果是.NET Core 3.0+:可以用
Channel<T>替代ConcurrentQueue+SemaphoreSlim,Channel是专门为生产者-消费者模式设计的,API更简洁。
内容的提问来源于stack exchange,提问作者AdvApp
相关产品推荐
相关产品推荐

