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

运行Parallel.ForEach时如何向其绑定的List动态新增元素并利用空闲线程

核心问题说明

直接使用List<T>配合Parallel.ForEach无法实现动态新增元素处理,原因有两个:

  1. Parallel.ForEach启动时会获取集合的快照枚举,后续新增的元素不会被纳入遍历范围
  2. List<T>不是线程安全的,多线程读写会出现元素丢失、索引异常等问题

可行解决方案

方案1:使用BlockingCollection<T> + Parallel.ForEach(改动最小)

BlockingCollection<T>是.NET内置的线程安全阻塞集合,天然适配生产者消费者场景,消费者可以阻塞等待新元素加入,直到标记添加完成。
代码实现如下:

// 初始化线程安全集合,底层默认是ConcurrentQueue,支持先进先出
// 可指定最大容量避免内存溢出,比如new BlockingCollection<Customer>(1000)
BlockingCollection<Customer> customerQueue = new BlockingCollection<Customer>();
// 先加载第一批未处理数据
foreach (var cust in GetNotProcessedCostumer())
{
    customerQueue.Add(cust);
}

// 启动后台生产者任务,定期拉取数据库新增数据
_ = Task.Run(async () =>
{
    // 可通过CancellationToken控制停止逻辑,比如程序关闭时终止拉取
    while (!stoppingToken.IsCancellationRequested)
    {
        var newCustomers = GetNotProcessedCostumer();
        foreach (var cust in newCustomers)
        {
            customerQueue.Add(cust, stoppingToken);
        }
        // 自定义拉取间隔,示例为1分钟
        await Task.Delay(TimeSpan.FromMinutes(1), stoppingToken);
        
        // 确认后续无新增数据时,调用CompleteAdding,Parallel处理完所有元素后会自动退出
        // if (noMoreData) customerQueue.CompleteAdding();
    }
}, stoppingToken);

// 并行处理逻辑,GetConsumingEnumerable会阻塞等待新元素,直到集合被标记为添加完成
Parallel.ForEach(customerQueue.GetConsumingEnumerable(stoppingToken),
    new ParallelOptions { MaxDegreeOfParallelism = 2, CancellationToken = stoppingToken },
    cust =>
    {
        ExecuteSomething(cust);
    });

方案2:使用TPL Dataflow的ActionBlock<T>(更推荐,拓展性更强)

TPL Dataflow是.NET官方提供的数据流处理库,内置并行度控制、线程调度,原生支持动态接收输入项,不需要手动维护集合和遍历逻辑。
首先需要通过NuGet安装System.Threading.Tasks.Dataflow包,代码实现如下:

// 初始化处理块,指定并行度和取消令牌
var processBlock = new ActionBlock<Customer>(
    cust => ExecuteSomething(cust),
    new ExecutionDataflowBlockOptions
    {
        MaxDegreeOfParallelism = 2,
        CancellationToken = stoppingToken
    });

// 加载第一批未处理数据
foreach (var cust in GetNotProcessedCostumer())
{
    processBlock.Post(cust);
}

// 启动后台拉取任务
_ = Task.Run(async () =>
{
    while (!stoppingToken.IsCancellationRequested)
    {
        var newCustomers = GetNotProcessedCostumer();
        foreach (var cust in newCustomers)
        {
            processBlock.Post(cust);
        }
        await Task.Delay(TimeSpan.FromMinutes(1), stoppingToken);
        
        // 确认无新增数据时调用Complete,等待所有处理完成
        // if (noMoreData) processBlock.Complete();
    }
}, stoppingToken);

// 若需要等待所有处理完成,可等待Completion任务
// await processBlock.Completion;

该方案优势:

  • 内置线程调度,空闲线程会自动处理新提交的元素,完全匹配需求
  • 可方便拓展重试、批量处理、限流等逻辑,适合复杂业务场景
  • 无需手动处理集合线程安全问题,Post方法原生线程安全

内容的提问来源于stack exchange,提问作者Gabriel Burato

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 03:45:06