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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 09:30:24