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

.NET嵌套异步任务未正确等待问题排查

邮件特征扩展脚本的并行任务等待问题排查与修复

问题背景

我正在用C#编写脚本,调用外部API为邮件数据集扩展特征,具体场景:

  • 存在N个邮件(emails数组),每个邮件包含mailID和MailData属性,需调用3个异步方法:ComputeAttachmentsFeatures、ComputeVirusTotalFeatures、ComputeHeaderFeatures为邮件对象补充数据;
  • 每个mail对象的MailURLs属性包含0个或多个URL,每个URL需调用3个异步方法:ComputeDNSFeatures、ComputeWhoIsFeatures、ComputePageRankFeatures;

预期每个邮件迭代中,上述3个邮件任务与每个URL的3个任务并行执行,最多同时运行6个任务。但当前代码未等所有任务完成就进入下一个邮件处理流程,已尝试Task.WaitAll(Task[])方法,请求排查错误。

当前代码

foreach ((long mailID, MailData mail) in emails) 
{
    Logger.Info("Mail " + mailID);
    Task [] mailTasks =
    {
        Task.Run(async () => await mail.ComputeAttachmentsFeatures(client_Attachments, mailID)),
        Task.Run(async () => await mail.ComputeVirusTotalFeatures(client_VT, mailID)),
        Task.Run(async () => await mail.ComputeHeaderFeatures(client_Headers, mailID))
    };
    //For each URL in email, call the APIs
    foreach (var url in mail.MailURLs)
    {
        if (IsValidURL(url.FullHostName))
        {
            Task [] urlTasks =
            {
                Task.Run(async () => await url.ComputeDNSFeatures(mailID)), // Task.Factory.StartNew
                Task.Run(async () => await url.ComputeWhoIsFeatures(mailID)),
                Task.Run(async () => await url.ComputePageRankFeatures(client_PageRank, mailID))
            };
            Task.WaitAll(urlTasks);
        }
    }
    Task.WaitAll(mailTasks);
    Logger.Info("Mail completed " + mailID);
}

示例异步函数

public async Task<int> ComputeHeaderFeatures(HttpClient client, long mailId)
{
    //  Blacklists check of the traversed mailservers  -n_smtp_servers_blacklist-
    n_smtp_servers_blacklist = 0;
    foreach (string mail_server in ServersInReceivedHeaders)
    {
        if (!string.IsNullOrEmpty(mail_server) && Program.IsValidURL(mail_server))  // Only analyze it if it's a valid URL or IP
        {
            // API call to check the mail_server against more than 100 blacklists
            BlacklistURL alreadyAnalyzedURL = (BlacklistURL)Program.BlacklistedURLs.Find(mail_server);    // Checks if the IP has already been analyzed
            if (alreadyAnalyzedURL == null)
            {
                BlacklistURL blacklistsResult = new BlacklistURL(mail_server);
                while (true)
                {
                    try
                    {
                        int statusCode = await BlacklistURL_API.PerformAPICall(blacklistsResult, client);

                        if (statusCode == 429) // Rate Limit Hit
                        {
                            Program.Logger.Debug($"BlackListChecker - Rate limit hit for mail {mailId} - {mail_server}, will retry after {CapDelayBlackListChecking} s");
                            Thread.Sleep(CapDelayBlackListChecking * 1000);
                            continue;
                        }
                        if (statusCode == 503)  // Service temporarily unavailable
                        {
                            Program.Logger.Debug($"BlackListChecker - Service temporarily Unavailable (503 error), will retry after {DefaultDelay503} s");
                            Thread.Sleep(DefaultDelay503 * 1000);
                            continue;  // try again later
                        }
                        Program.BlacklistedURLs.Add(blacklistsResult);  // Adds the server and its result to the list of already analyzed servers
                        if (blacklistsResult.GetFeature() > 0) { n_smtp_servers_blacklist++; }  // If the server appears in at least 1 blacklist, we count it as malicious

                        Program.Logger.Debug($"BlackListChecker call for mail {mailId} - {mail_server} responded with status code {statusCode}");
                    }
                    catch (Exception ex)
                    {
                        Program.Logger.Error($"BlackListChecker - An exception occurred for mail {mailId} - {mail_server}:\n{ex}");
                        throw;
                    }
                    finally
                    {
                        Program.Logger.Debug($"BlackListChecker - Wait {DefaultGapBetweenCallsBlackListChecking} s before next call...");
                        Thread.Sleep(DefaultGapBetweenCallsBlackListChecking * 1000);
                    }
                    // Break the inner loop to proceed with the next argument
                    break;
                }
            }
            else  // The mailserver has already been analyzed, so we take the available result
            {
                if (alreadyAnalyzedURL.NBlacklists > 0) { n_smtp_servers_blacklist++; }
            }
        }
    }
    return 1;
}

错误原因分析

  1. 并行逻辑完全错误:当前代码启动邮件任务后,立刻串行处理每个URL的任务(等一个URL的3个任务完成再处理下一个),最后才等待邮件任务结束。这导致邮件任务在后台运行时,URL任务串行阻塞,若邮件任务先完成,会直接进入下一个邮件循环,完全不符合“邮件任务与所有URL任务并行”的预期。
  2. 冗余的Task.Run包装:异步方法本身返回Task,无需额外用Task.Run(async () => await ...)包裹,这种写法会多创建一层任务,徒增开销。
  3. 同步阻塞的Task.WaitAll:在异步代码中使用WaitAll会阻塞线程池线程,结合示例中异步方法里的Thread.Sleep,会严重降低并行效率,甚至引发线程饥饿。
  4. 全局集合线程不安全:Program.BlacklistedURLs是普通集合,多并行任务同时读写会导致数据不一致或异常。

修复方案

1. 重构并行逻辑,统一收集所有任务

将邮件任务和所有URL任务收集到同一个任务列表中,统一等待,实现真正的并行执行。

2. 用await Task.WhenAll替代Task.WaitAll

异步等待不会阻塞线程,能让线程池更高效地利用资源。

3. 移除冗余的Task.Run

直接使用异步方法返回的Task加入任务列表即可。

4. 修复全局集合线程安全问题

改用ConcurrentDictionary<TKey, TValue>存储已分析的黑名单URL,确保多线程读写安全。

5. 替换Thread.Sleep为await Task.Delay

Task.Delay会释放当前线程,避免线程阻塞,提升并行效率。

修复后的代码示例

// 替换原全局集合为线程安全版本
private static ConcurrentDictionary<string, BlacklistURL> BlacklistedURLs = new ConcurrentDictionary<string, BlacklistURL>();

// 单个邮件处理的异步方法
private async Task ProcessMailAsync((long mailID, MailData mail) mailItem)
{
    var (mailID, mail) = mailItem;
    Logger.Info($"Mail {mailID}");

    // 收集所有邮件任务
    var allTasks = new List<Task>
    {
        mail.ComputeAttachmentsFeatures(client_Attachments, mailID),
        mail.ComputeVirusTotalFeatures(client_VT, mailID),
        mail.ComputeHeaderFeatures(client_Headers, mailID)
    };

    // 收集所有URL任务
    foreach (var url in mail.MailURLs)
    {
        if (IsValidURL(url.FullHostName))
        {
            allTasks.AddRange(new[]
            {
                url.ComputeDNSFeatures(mailID),
                url.ComputeWhoIsFeatures(mailID),
                url.ComputePageRankFeatures(client_PageRank, mailID)
            });
        }
    }

    // 异步等待所有任务完成
    await Task.WhenAll(allTasks);

    Logger.Info($"Mail completed {mailID}");
}

// 主循环调用
foreach (var mailItem in emails)
{
    // 若需限制全局并行邮件数,可添加SemaphoreSlim控制
    await ProcessMailAsync(mailItem);
}

异步方法中的Sleep替换示例

// 替换原Thread.Sleep为Task.Delay
if (statusCode == 429) // Rate Limit Hit
{
    Program.Logger.Debug($"BlackListChecker - Rate limit hit for mail {mailId} - {mail_server}, will retry after {CapDelayBlackListChecking} s");
    await Task.Delay(CapDelayBlackListChecking * 1000);
    continue;
}

// 其他Thread.Sleep同理替换

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 18:15:54