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

Parallel.ForEach结合SQL每20条批量插入的线程安全问题

问题

在Parallel.ForEach中执行任务,每个线程完成后将结果存入List,当List元素达到20条时批量插入数据库,以此避免单条插入的连接开销。但当前代码因多线程同时修改List触发Collection was modified异常,不想用Lock以免损失并行优势,尝试暂停恢复Foreach也没成功,求最优实现方式。

相关代码

主逻辑代码

private ItemsDAL _itemsDal = new ItemsDAL();
public void Run()
{
    List<Item> items = new List<Item>();

    for (int i = 0; i < 1500; i++)
    {
        items.Add(new Item()
        {
            Details = "Test " + i
        });
    }

    Parallel.ForEach(items, (oneItem) =>
    {
        Console.WriteLine("Adding item with Details " + oneItem.Details);
        _itemsDal.InsertToDatabase(oneItem);
    });
}

DAL类代码

public class ItemsDAL
{
    object _lock = new object();
    private List<Item> _itemsToInsert = new List<Item>();

    public void InsertToDatabase(Item item)
    {
        _itemsToInsert.Add(item);
        if (_itemsToInsert.Count >= 20)
        {
            // 批量插入数据库(每20条一组)
            using (ForeachContext db = new ForeachContext())
            {
                foreach (var itemToInsert in _itemsToInsert)
                {
                    db.Items.Add(itemToInsert);
                }
                db.SaveChanges();
            }

            _itemsToInsert.Clear();
        }
    }
}
最优实现方案

方案2:分区并行处理+局部批量(推荐,完全无锁)

直接利用Parallel.ForEach的分区特性,让每个线程(分区)维护自己的局部List,达到20条阈值时独立批量提交,完全避免跨线程共享集合的冲突,并行效率拉满。

修改后的主逻辑代码:

public void Run()
{
    List<Item> items = new List<Item>();

    for (int i = 0; i < 1500; i++)
    {
        items.Add(new Item()
        {
            Details = "Test " + i
        });
    }

    Parallel.ForEach(items, 
        // 每个分区初始化一个专属的局部List,不与其他线程共享
        () => new List<Item>(), 
        (oneItem, state, localList) =>
        {
            Console.WriteLine("Adding item with Details " + oneItem.Details);
            localList.Add(oneItem);
            
            // 局部List攒够20条就批量提交
            if (localList.Count >= 20)
            {
                using (ForeachContext db = new ForeachContext())
                {
                    db.Items.AddRange(localList); // 用AddRange比循环Add效率更高
                    db.SaveChanges();
                }
                localList.Clear();
            }
            return localList;
        },
        // 每个分区结束后,处理剩余不足20条的数据
        (localList) => 
        {
            if (localList.Count > 0)
            {
                using (ForeachContext db = new ForeachContext())
                {
                    db.Items.AddRange(localList);
                    db.SaveChanges();
                }
            }
        });
}

这种方式不需要修改DAL类,也没有任何锁,每个线程独立处理自己的批次,完全不会出现集合修改异常,并行性能不受影响。

方案1:线程安全集合+轻量锁(适合集中管理场景)

如果需要集中控制批量提交的逻辑,可以用ConcurrentBag<Item>替代普通List,配合原子计数+仅在批量提交时加轻量锁,把锁的影响降到最低。

修改后的DAL类代码:

public class ItemsDAL
{
    // 线程安全集合,支持多线程添加/取出
    private readonly ConcurrentBag<Item> _itemsToInsert = new ConcurrentBag<Item>();
    // 原子计数,避免多线程计数冲突
    private int _count = 0;
    // 仅在批量提交时加锁,防止多个线程同时执行提交操作
    private readonly object _batchLock = new object(); 

    public void InsertToDatabase(Item item)
    {
        _itemsToInsert.Add(item);
        // 原子递增计数,线程安全
        int currentCount = Interlocked.Increment(ref _count);

        // 达到20条阈值时触发批量提交
        if (currentCount >= 20)
        {
            lock (_batchLock)
            {
                // 二次检查,防止多个线程同时进入提交逻辑
                if (_count >= 20)
                {
                    var batchItems = new List<Item>();
                    // 取出所有待提交的元素
                    while (_itemsToInsert.TryTake(out var itemToInsert))
                    {
                        batchItems.Add(itemToInsert);
                    }

                    using (ForeachContext db = new ForeachContext())
                    {
                        db.Items.AddRange(batchItems);
                        db.SaveChanges();
                    }

                    // 重置计数
                    Interlocked.Exchange(ref _count, 0);
                }
            }
        }
    }

    // 必须在Parallel.ForEach结束后调用,处理最后一批不足20条的数据
    public void FlushRemaining()
    {
        lock (_batchLock)
        {
            var batchItems = new List<Item>();
            while (_itemsToInsert.TryTake(out var itemToInsert))
            {
                batchItems.Add(itemToInsert);
            }

            if (batchItems.Count > 0)
            {
                using (ForeachContext db = new ForeachContext())
                {
                    db.Items.AddRange(batchItems);
                    db.SaveChanges();
                }
                Interlocked.Exchange(ref _count, 0);
            }
        }
    }
}

主逻辑需要补充调用FlushRemaining:

public void Run()
{
    List<Item> items = new List<Item>();

    for (int i = 0; i < 1500; i++)
    {
        items.Add(new Item()
        {
            Details = "Test " + i
        });
    }

    Parallel.ForEach(items, (oneItem) =>
    {
        Console.WriteLine("Adding item with Details " + oneItem.Details);
        _itemsDal.InsertToDatabase(oneItem);
    });

    // 处理最后一批不足20条的数据
    _itemsDal.FlushRemaining();
}

方案对比

  • 方案2是最优选择,完全无锁,并行效率最高,代码简洁,不需要额外维护线程安全集合。
  • 方案1适合需要集中管理批量提交规则的场景,仅在批量提交时加锁,对并行性能影响极小。

内容的提问来源于stack exchange,提问作者Rafa Ayadi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 10:43:10