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
相关产品推荐
相关产品推荐

