如何获取Parallel.ForEachAsync各任务状态以处理部分失败场景?
解决Worker Service并行任务部分失败时的状态处理问题
问题核心
当前逻辑里,两个并行调用自研API的任务只要有一个成功就会更新数据库记录为完成(status=1),但如果另一个任务失败,实际只获取到部分游戏码,导致状态与实际情况完全不符。要解决这个问题,必须跟踪每个并行任务的执行状态,等所有任务跑完后统一判断整体结果,再更新数据库状态。
解决方案思路
- 定义任务结果模型:存每个任务的请求数量、执行状态(成功/失败)、实际获取到的游戏码数量(可选)
- 线程安全收集结果:用
ConcurrentBag<T>来并行收集任务结果,避免多线程冲突 - 统一处理结果:所有任务执行完后,统计成功任务数、总获取的游戏码数量,再根据统计结果更新数据库:
- 全部任务成功:标记为完成(status=1)
- 部分任务成功:标记为部分完成(比如status=2),同时记录已获取的游戏码数量
- 全部任务失败:标记为失败(status=3)
修改后的代码示例
// 定义任务结果模型 public class TaskResult { public int RequestQuantity { get; set; } public bool IsSuccess { get; set; } public int ReceivedCouponCount { get; set; } // 实际拿到的游戏码数量 } // 业务逻辑部分 var num = 20; var firstNum = 10; var secondNum = 10; if (num < 20) { firstNum = (num + 1) / 2; secondNum = num - firstNum; } var quantities = new List<int> { firstNum, secondNum }; var taskResults = new ConcurrentBag<TaskResult>(); // 线程安全的结果集合 var cts = new CancellationTokenSource(); ParallelOptions parallelOptions = new() { MaxDegreeOfParallelism = 2, CancellationToken = cts.Token }; try { await Parallel.ForEachAsync(quantities, parallelOptions, async (quantity, ct) => { var result = new TaskResult { RequestQuantity = quantity }; var content = new FormUrlEncodedContent(new[] { new KeyValuePair<string, string>("productCode", productCode), new KeyValuePair<string, string>("quantity", quantity.ToString()), new KeyValuePair<string, string>("clientTrxRef", bulkId.ToString()) }); try { using var response = await httpClient.PostAsync(_configuration["Razer:Production"], content, ct); if (response.IsSuccessStatusCode) { var coupon = await response.Content.ReadFromJsonAsync<Root>(cancellationToken: ct); // 假设Root模型里有游戏码列表的Count属性 result.IsSuccess = true; result.ReceivedCouponCount = coupon?.Coupons?.Count ?? 0; } else { result.IsSuccess = false; // 记录日志:请求失败,状态码:response.StatusCode } } catch (HttpRequestException ex) { result.IsSuccess = false; // 记录日志:请求异常,信息:ex.Message } finally { taskResults.Add(result); } }); // 所有任务完成后统一处理结果 var totalSuccess = taskResults.Count(r => r.IsSuccess); var totalReceived = taskResults.Sum(r => r.ReceivedCouponCount); var totalTasks = taskResults.Count; if (totalSuccess == totalTasks) { // 全部成功,标记为完成 await UpdateBulkStatus(id, status: 1, receivedCount: totalReceived); } else if (totalSuccess > 0) { // 部分成功,标记为部分完成并记录已获取数量 await UpdateBulkStatus(id, status: 2, receivedCount: totalReceived); } else { // 全部失败,标记为失败 await UpdateBulkStatus(id, status: 3, receivedCount: 0); } } catch (OperationCanceledException ex) { // 记录日志:任务被取消,信息:ex.Message // 标记任务为取消状态 await UpdateBulkStatus(id, status: 4, receivedCount: 0); } // 统一更新状态的方法 async Task UpdateBulkStatus(int id, int status, int receivedCount) { // 执行数据库更新逻辑:修改status和received_count字段 }
关键改进点
- 不再在单个任务里更新数据库,而是统一收集所有结果后再做状态更新,避免部分成功导致的状态错误
- 用
ConcurrentBag<TaskResult>安全收集并行任务结果,保证多线程环境下的数据一致性 - 新增了部分完成、失败、取消等状态的处理逻辑,让数据库状态能准确反映实际执行情况
内容的提问来源于stack exchange,提问作者raysefo
相关产品推荐
相关产品推荐

