ASP.NET6中百万级Task引发System.OutOfMemoryException问题及优化咨询
大规模.sig文件处理的内存溢出问题与多线程优化方案
问题背景
我有600万个平均大小约15字节的.sig文件需要读取并处理,此前在ASP.NET Core 2.1中用Task.Factory实现,耗时约20小时且运行正常。迁移到ASP.NET 6后,测试服务器启动文件操作后Web应用无响应,日志出现System.OutOfMemoryException错误。当前实现方式不够理想,寻求该任务的多线程优化方案。
相关代码
控制器中的ImportSignatures方法
[HttpPost("ImportSignatures")] public JsonResult ImportSignatures() { try { return Json(SignatureImportService.ImportSigningCerts()); } catch (Exception e) { LogsHelper.WriteLog("api/Settings/ImportSignatures", e); return Json(new ImportSigningCertsResult(e.Message, SignatureImportService.WasCancelled)); } }
ImportSigningCerts方法
public static ImportSigningCertsResult ImportSigningCerts() { LogsHelper.WriteEventLog("Launching SignatureImportService"); WasCancelled = false; IsWorking = true; ResultStr = ""; totalSignatures = 0; processedSignatures = 0; var cancelMsg = "Certificate import was interrupted. \n"; var endMsg = "Certificate import completed successfully. \n"; var toDelete = new List<string>(); try { var configuration = SignatureImportConfiguration.FromCfg(); using (s_tokenSource = new CancellationTokenSource()) { List<string> signatures = Directory.EnumerateFiles(configuration.Path, "*.sig").ToList(); totalSignatures = signatures.Count; Store mainStore = StoreMan.GetStore("Main"); var importStats = new ImportStats(); var tasks = new List<Task>(); int saveIndex = 1; const int proccessedForSave = 100000; // 处理多少个签名后执行中间存储和文件删除 CancellationToken token = s_tokenSource.Token; int minWorkerThreads, minCompletionPortThreads, maxWorkerThreads, maxCompletionPortThreads; ThreadPool.GetMinThreads(out minWorkerThreads, out minCompletionPortThreads); ThreadPool.GetMaxThreads(out maxWorkerThreads, out maxCompletionPortThreads); ThreadPool.SetMaxThreads(minWorkerThreads * 2, maxCompletionPortThreads); signatures.ForEach(path => { tasks.Add(Task.Factory.StartNew(() => { token.ThrowIfCancellationRequested(); // 读取当前文件并将必要证书上传到存储 if (UploadSigningCerts(mainStore, path, importStats)) { if (configuration.NeedCleaning) { lock (s_toDeleteListLockObj) toDelete.Add(path); } } // 中间存储保存和已处理文件删除 lock (s_intermediateSaveLockObj) { if (++processedSignatures > proccessedForSave * saveIndex) { LogsHelper.WriteEventLog("Intermediate saving of the certificate store..."); mainStore.WriteIfChanged(); StartRemovingSignatures(toDelete); saveIndex++; } } }, token)); }); try { Task.WaitAll(tasks.ToArray()); } catch (AggregateException ae) { foreach (Exception e in ae.InnerExceptions) { if (e is not TaskCanceledException) LogsHelper.WriteLog("SignatureImportService/ImportSigningCerts", e); } } mainStore.WriteIfChanged(); StartRemovingSignatures(toDelete); ResultStr = (WasCancelled ? cancelMsg : endMsg) + $"Certificates found: {importStats.all}. Was imported: {importStats.imported}." + (importStats.parsingFailed > 0 ? $" Unrecognized files: {importStats.parsingFailed}" : ""); } LogsHelper.WriteEventLog(ResultStr); return s_tokenSource == null ? new ImportSigningCertsResult(ResultStr) : new ImportSigningCertsResult(ResultStr, WasCancelled); } catch (Exception) { throw; } finally { IsWorking = false; } }
UploadSigningCerts方法
private static bool UploadSigningCerts(Store store, string path, ImportStats importStats) { bool toBeDeleted = true; CryptoClient client = CryptoServiceContext.DefaultInstance.CryptoClient; try { List<CertInfo> certs = client.GetSignCmsInfo(File.ReadAllBytes(path)).Certs.ToList(); Interlocked.Add(ref importStats.all, certs.Count); for (int i = 0; i < certs.Count; i++) { lock (s_importLockObj) { // 验证文件中的每个证书,决定是否导入到存储... } } return toBeDeleted; } catch (Exception e) { LogsHelper.WriteLog("SignatureImportService/UploadSigningCerts", e); LogsHelper.WriteEventLog($"Error importing certificate from signature: {Path.GetFileName(path)};"); Interlocked.Increment(ref importStats.errors); return false; } }
StartRemovingSignatures方法
private static void StartRemovingSignatures(List<string> toDelete) { if (toDelete.Count > 0) { List<string> tempToDelete; lock (s_toDeleteListLockObj) { tempToDelete = new List<string>(toDelete); toDelete.Clear(); } LogsHelper.WriteEventLog("Deleting successfully processed signature files..."); Task.Factory.StartNew(() => { tempToDelete.ForEach(path => { try { File.Delete(path); } catch (Exception e) { LogsHelper.WriteLog("ImportResult/DeleteSignatures", e); } }); }); } }
错误日志
20.08.2023 11:58:01 api/Settings/ImportSignatures Exception of type 'System.OutOfMemoryException' was thrown. at System.Threading.Tasks.Task.EnsureContingentPropertiesInitializedUnsafe() at System.Threading.Tasks.Task.AssignCancellationToken(CancellationToken cancellationToken, Task antecedent, TaskContinuation continuation) at System.Threading.Tasks.Task.TaskConstructorCore(Delegate action, Object state, CancellationToken cancellationToken, TaskCreationOptions creationOptions, InternalTaskOptions internalOptions, TaskScheduler scheduler) at Store.Services.SignatureImportService.<>c__DisplayClass20_0.<ImportSigningCerts>b__0(String path) at System.Collections.Generic.List`1.ForEach(Action`1 action) at Store.Services.SignatureImportService.ImportSigningCerts() at Store.Controllers.SettingsController.ImportSignatures()
优化方案
1. 限制并发任务数量,避免内存暴涨
当前为600万个文件创建600万个Task,每个Task对象的内存开销直接导致溢出。改用Parallel.ForEach或TPL数据流控制并发数:
// 使用Parallel.ForEach限制并发数,建议设置为CPU核心数的2-4倍 Parallel.ForEach(signatures, new ParallelOptions { MaxDegreeOfParallelism = Environment.ProcessorCount * 2, CancellationToken = token }, path => { // 原Task内的业务逻辑 });
2. 异步化IO操作,释放线程资源
替换同步文件读取为异步方法,减少线程阻塞,提升线程池利用率:
// 修改UploadSigningCerts为异步方法 private static async Task<bool> UploadSigningCertsAsync(Store store, string path, ImportStats importStats) { bool toBeDeleted = true; CryptoClient client = CryptoServiceContext.DefaultInstance.CryptoClient; try { byte[] fileBytes = await File.ReadAllBytesAsync(path); List<CertInfo> certs = client.GetSignCmsInfo(fileBytes).Certs.ToList(); Interlocked.Add(ref importStats.all, certs.Count); for (int i = 0; i < certs.Count; i++) { lock (s_importLockObj) { // 证书验证与导入逻辑 } } return toBeDeleted; } catch (Exception e) { LogsHelper.WriteLog("SignatureImportService/UploadSigningCerts", e); LogsHelper.WriteEventLog($"Error importing certificate from signature: {Path.GetFileName(path)};"); Interlocked.Increment(ref importStats.errors); return false; } }
3. 分批加载文件路径,避免一次性占用内存
Directory.EnumerateFiles是延迟加载的,不要直接转成List,改为分批处理:
var fileEnumerator = Directory.EnumerateFiles(configuration.Path, "*.sig").GetEnumerator(); while (fileEnumerator.MoveNext()) { var batch = new List<string>(); batch.Add(fileEnumerator.Current); // 每次处理1000个文件的批次 for (int i = 0; i < 999 && fileEnumerator.MoveNext(); i++) { batch.Add(fileEnumerator.Current); } // 异步处理当前批次 await Parallel.ForEachAsync(batch, new ParallelOptions { MaxDegreeOfParallelism = Environment.ProcessorCount * 2 }, async (path, token) => { token.ThrowIfCancellationRequested(); if (await UploadSigningCertsAsync(mainStore, path, importStats)) { if (configuration.NeedCleaning) { lock (s_toDeleteListLockObj) toDelete.Add(path); } } // 中间保存逻辑 lock (s_intermediateSaveLockObj) { if (++processedSignatures > proccessedForSave * saveIndex) { LogsHelper.WriteEventLog("Intermediate saving of the certificate store..."); mainStore.WriteIfChanged(); StartRemovingSignatures(toDelete); saveIndex++; } } }); }
4. 异步化整个调用链,避免阻塞Web请求
控制器和服务方法改为异步,替换Task.WaitAll为await Task.WhenAll,防止Web请求线程被阻塞:
// 控制器异步方法 [HttpPost("ImportSignatures")] public async Task<JsonResult> ImportSignatures() { try { return Json(await SignatureImportService.ImportSigningCertsAsync()); } catch (Exception e) { LogsHelper.WriteLog("api/Settings/ImportSignatures", e); return Json(new ImportSigningCertsResult(e.Message, SignatureImportService.WasCancelled)); } }
5. 优化线程池设置
ASP.NET 6中线程池会自动调整,无需手动强制修改MaxThreads。如果必须调整,应根据服务器CPU核心数合理设置,避免线程过多导致上下文切换开销激增。
6. 批量优化文件删除逻辑
每次删除文件时避免创建新Task,可复用异步批量删除逻辑,减少Task创建开销,同时控制删除速率避免IO压力过大。
内容的提问来源于stack exchange,提问作者Dany
相关产品推荐
相关产品推荐

