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

ConcurrentBag迭代触发System.ObjectDisposedException异常的排查与解决

问题描述

我有一个Dictionary记录列表,将这些记录转换为对象后,通过CheckPropertyIsObjectorArrayOrDefault方法处理数据并添加至ConcurrentBag<JObject>类型的parsedObjectList中。后续使用while循环迭代该ConcurrentBag时抛出了System.ObjectDisposedException异常,异常信息如下:

[12:30:16] [Error] (2) Exception :: System.ObjectDisposedException: Cannot access a disposed object. Object name: 'The ThreadLocal object has been disposed.'. at System.Threading.ThreadLocal`1.GetValueSlow() at myApp.Client.LoadTenantConfigCache(JObject tableObject, String tenantId, String tenantName, String dbQuery) in C:\code\v2\myApp\myApp\Client.cs:line 1147, Properties:: {"MethodName": "LoadTenantConfigCache()", "Payload": null, "innerException": null, "MainException": "System.ObjectDisposedException: Cannot access a disposed object.\r\nObject name: 'The ThreadLocal object has been disposed.'.\r\n at System.Threading.ThreadLocal`1.GetValueSlow()\r\n at myApp.Client.LoadTenantConfigCache(JObject tableObject, String tenantId, String tenantName, String dbQuery) in C:\code\v2\myApp\myApp\Client.cs:line 1147"} {"Environment": "Production"} 
[12:30:16] [Error] (2) Exception :: System.ObjectDisposedException: Cannot access a disposed object. Object name: 'The ThreadLocal object has been disposed.'. at System.Threading.ThreadLocal`1.GetValueSlow() at myApp.Client.LoadTenantConfigCache(JObject tableObject, String tenantId, String tenantName, String dbQuery) in C:\code\v2\myApp\myApp\Client.cs:line 1147[TableName : Entities] [Count: 30] , Properties:: {"MethodName": "LoadTenantConfigCache()", "Payload": null, "innerException": null, "MainException": "System.ObjectDisposedException: Cannot access a disposed object.\r\nObject name: 'The ThreadLocal object has been disposed.'.\r\n at System.Threading.ThreadLocal`1.GetValueSlow()\r\n at myApp.Client.LoadTenantConfigCache(JObject tableObject, String tenantId, String tenantName, String dbQuery) in C:\code\v2\myApp\myApp\Client.cs:line 1147"} {"Environment": "Production"}

相关代码如下:

List<Dictionary<string, object>> objectList = new List<Dictionary<string, object>>();
//adding data into objectList logic ... ...
ConcurrentBag<JObject> parsedObjectList = new ConcurrentBag<JObject>();
var tasks = new ConcurrentBag<Task>();
foreach (var dictObject in objectList)
{
    var task = Task.Factory.StartNew(async (object entity) =>
    {
        JObject obj = new JObject();
        var res = "";
        var key = "";
        var column = "";
        Dictionary<string, object> data = (Dictionary<string, object>)entity;
        foreach (var readerObject in data)
        {
            key = readerObject.Key;
            column = readerObject.Value.ToString();
            res = await _common.CheckPropertyIsObjectorArrayOrDefault(column, key, tableName, tempInstance, obj);
        }
        if (res == "Success")
            parsedObjectList.Add(obj);
        else
            Console.WriteLine("Missing Records Key:" + key + " Value: " + column);
    }, state: dictObject);
    tasks.Add(task);
}
await Task.WhenAll(tasks);

using (var ldr = Ignite.GetDataStreamer<string, ICustomCacheStore>(cacheName))
{
    ldr.AutoFlushFrequency = 10;
    JObject item = null;
    while (parsedObjectList.TryTake(out item))//from here throwing error
    {
        var tInstance = (ICustomCacheStore)Activator.CreateInstance(tempType);
        JObject keyObj = new JObject();
        foreach (var keyName in keyArray)
        {
            var pk = keyName.ToString();
            if (item.ContainsKey(pk))
                keyObj[pk] = item[pk];
        }
        keyObj["TableName"] = tableName;
        //Create Instance For PrimarykeyModel
        var priKeyInstance = (IKeyModel)Activator.CreateInstance(Type.GetType(primaryKeyModel + ", Models"));
        var serializerSettings = new JsonSerializerSettings { NullValueHandling = NullValueHandling.Ignore };
        //populate Object for Type Key
        JsonConvert.PopulateObject(keyObj.ToString(), priKeyInstance, serializerSettings);
        JsonConvert.PopulateObject(item.ToString(), tInstance, serializerSettings);
        string json = JsonConvert.SerializeObject(priKeyInstance, Formatting.None);
        string base64EncodedKey = Convert.ToBase64String(Encoding.UTF8.GetBytes(json));
        await ldr.AddData(base64EncodedKey, tInstance);
    }
}

异常原因分析

咱们来拆解一下问题核心:

  1. Task.Factory.StartNew与async lambda的兼容问题:你用Task.Factory.StartNew传入了一个async委托,这个方法会返回Task<Task>——外层Task只是启动任务的容器,内层Task才是async方法实际执行的逻辑。但你直接把这个Task<Task>加到了ConcurrentBag<Task>里,调用await Task.WhenAll(tasks)时,其实只等待了外层Task完成,内层的async任务可能还在后台运行。这就导致后续代码(比如using块里的逻辑)可能在部分async任务还没处理完时就开始执行,甚至如果tempInstance或者_common内部持有ThreadLocal实例,这个实例可能在async任务还在使用时就被提前Dispose了。
  2. ThreadLocal生命周期管理失误:异常明确指向ThreadLocal对象已被释放,说明在调用CheckPropertyIsObjectorArrayOrDefault方法时,方法内部依赖的ThreadLocal实例已经被Dispose了。结合上面的任务等待问题,很大概率是因为外层Task完成后,某些持有ThreadLocal的对象(比如tempInstance)被提前回收或释放,而内层async任务还在尝试访问它。

解决建议

针对这两个核心问题,给出具体的修复方案:

1. 替换Task.Factory.StartNew为Task.Run

Task.Run专门为async委托设计,会自动展开Task<Task>返回正确的Task,确保Task.WhenAll能等待所有async任务完成。修改代码如下:

foreach (var dictObject in objectList)
{
    var task = Task.Run(async () =>
    {
        JObject obj = new JObject();
        var res = "";
        var key = "";
        var column = "";
        Dictionary<string, object> data = dictObject;
        foreach (var readerObject in data)
        {
            key = readerObject.Key;
            column = readerObject.Value.ToString();
            res = await _common.CheckPropertyIsObjectorArrayOrDefault(column, key, tableName, tempInstance, obj);
        }
        if (res == "Success")
            parsedObjectList.Add(obj);
        else
            Console.WriteLine("Missing Records Key:" + key + " Value: " + column);
    });
    tasks.Add(task);
}
await Task.WhenAll(tasks);

这里直接用Task.Run,不需要传递state参数,直接捕获dictObject即可(C# 5+已经解决了foreach循环变量捕获的问题,不用担心线程安全问题)。

2. 确保ThreadLocal实例的生命周期覆盖所有任务

检查tempInstance或者_common内部的ThreadLocal对象:

  • 如果ThreadLocal是在LoadTenantConfigCache方法内创建的,确保它的Dispose调用在所有async任务完成之后。
  • 如果tempInstance会在方法内被Dispose,需要确保它的生命周期延续到await Task.WhenAll(tasks)之后再释放。
  • 避免在using块中持有ThreadLocal相关对象,除非能确保所有依赖它的任务都已完成。

3. 可选:优化ConcurrentBag的使用(非必须,但更规范)

虽然ConcurrentBag是线程安全的,但如果所有任务都完成后再处理数据,可以考虑用List<Task<JObject>>来收集结果,避免使用ConcurrentBag,这样代码更直观:

var tasks = new List<Task<JObject>>();
foreach (var dictObject in objectList)
{
    var task = Task.Run(async () =>
    {
        JObject obj = new JObject();
        var res = "";
        var key = "";
        var column = "";
        Dictionary<string, object> data = dictObject;
        foreach (var readerObject in data)
        {
            key = readerObject.Key;
            column = readerObject.Value.ToString();
            res = await _common.CheckPropertyIsObjectorArrayOrDefault(column, key, tableName, tempInstance, obj);
        }
        if (res != "Success")
        {
            Console.WriteLine("Missing Records Key:" + key + " Value: " + column);
            return null;
        }
        return obj;
    });
    tasks.Add(task);
}
var results = await Task.WhenAll(tasks);
var parsedObjectList = results.Where(x => x != null).ToList();

这样后续处理parsedObjectList时就不需要用ConcurrentBag的TryTake逻辑,直接遍历列表即可,代码更清晰易维护。

内容的提问来源于stack exchange,提问作者cj devin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 09:04:57