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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 23:32:31