You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现父对象填充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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 09:16:40