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

