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); } }
异常原因分析
咱们来拆解一下问题核心:
- 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了。 - 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
相关产品推荐
相关产品推荐

