Parallel.ForEach并行处理数据丢失问题排查
解决并行代码中的数据丢失问题
嘿,我一眼就揪出你并行代码里几个关键的坑,这些正是导致并发下数据丢失、结果不一致的核心原因:
1. 最致命的错误:多线程共享非线程安全的DataReader
你在Parallel.ForEach里遍历的是reader.ToDataTable().AsEnumerable()的行,但循环体里居然还在调用原始的reader.Field<T>()!DataReader(比如SqlDataReader)是绝对线程不安全的,它的读取指针是全局共享的,多个线程同时调用Field会让指针乱跳,导致读取到错误的字段值,甚至直接读到空值——这就是你看到某些字段变成null的直接原因!
你应该用循环变量row来获取当前行的字段值,而不是共享的reader。
2. ConcurrentDictionary的使用方式错误,存在竞态条件
你用FirstOrDefault去查询ConcurrentDictionary,然后判断是否存在再手动添加/更新——这整个过程不是原子操作。在你检查到键不存在到调用AddOrUpdate的间隙,其他线程可能已经添加了同一个键,导致你的更新被覆盖,或者数据混乱。
ConcurrentDictionary专门提供了GetOrAdd这种原子方法,应该用它来确保键的创建和获取是线程安全的。
3. SomeClass中的Dictionary不是线程安全的
即使你用了ConcurrentDictionary存储SomeClass实例,SomeClass内部的普通Dictionary依然不支持并发修改。多个线程同时给同一个SomeClass实例的Dictionary[b] = c赋值时,会导致数据丢失或者线程安全异常。
应该把内部的Dictionary换成ConcurrentDictionary。
修正后的完整代码
var result = GetDataTable(); // 用ConcurrentDictionary存储分组数据,内部字段也用线程安全的ConcurrentDictionary var concurrentCollection = new ConcurrentDictionary<string, ConcurrentDictionary<string, object>>(); Parallel.ForEach(reader.ToDataTable().AsEnumerable(), new ParallelOptions { MaxDegreeOfParallelism = 2 }, row => { // 从当前循环的row中获取值,彻底抛弃共享的reader! var a = row.Field<int>("a").ToString(); var b = row.Field<int>("b").ToString(); var c = row.Field<double>("c"); var d = row.Field<string>("d"); var e = row.Field<string>("e"); // 生成分组唯一键 var groupKey = $"{a}|{d}"; // 原子获取或创建分组的ConcurrentDictionary,初始化固定字段 var groupData = concurrentCollection.GetOrAdd(groupKey, _ => new ConcurrentDictionary<string, object> { ["A"] = a, ["D"] = d, ["E"] = e }); // 线程安全地更新动态字段值 groupData[b] = c; }); // 把并行处理后的结果转成DataTable foreach (var ins in concurrentCollection.OrderBy(x => x.Value["D"]).ThenBy(x => x.Value["A"])) { var datarow = result.NewRow(); foreach (var kvp in ins.Value) { // 处理null值,避免DataTable插入DBNull报错 datarow[kvp.Key] = kvp.Value ?? DBNull.Value; } result.Rows.Add(datarow); } concurrentCollection.Clear(); return result;
额外优化建议
- 如果
reader的数据量极大,reader.ToDataTable()会一次性把所有数据加载到内存,可能占用较多内存。如果你的数据库支持,也可以考虑先分批加载数据到内存,再并行处理,从根源避免DataReader的线程安全问题。 MaxDegreeOfParallelism不要盲目设置太大,一般设置为Environment.ProcessorCount或者稍高一点即可,过高的并发反而会因为线程切换开销导致性能下降。
内容的提问来源于stack exchange,提问作者developer
相关产品推荐
相关产品推荐

