30万销售订单批量调用RestAPI:多线程与Semaphore实现方案咨询
解决方案:批量并发处理销售订单并收集API响应
针对你30万条销售订单的批量API调用需求,我整理了一个落地性很强的方案,结合批量拆分、并发控制、线程安全结果收集和Semaphore限流这些关键点,下面一步步给你拆解实现细节:
1. 核心思路梳理
- 复用你提供的
Batch扩展方法,把30万条订单拆分为每200/1000条的批次,减少API调用次数 - 用
Parallel.ForEach控制并发数(建议设为CPU核心数,比硬编码8更灵活),同时配合SemaphoreSlim额外限制API请求并发量(应对API限流场景) - 用线程安全的集合存储所有API响应,避免多线程下的列表操作冲突
2. 完整可运行代码实现
先定义基础类型,再实现核心处理逻辑:
using System.Collections.Concurrent; using System.Threading; using System.Text.Json; // 假设你的API返回响应类型,按需调整字段 public class ApiResponse { public int BatchItemCount { get; set; } public bool IsSuccess { get; set; } public string ResultMessage { get; set; } } // 你的销售订单类型示例 public class SalesOrder { public int OrderId { get; set; } public decimal Amount { get; set; } // 其他订单字段... } public class OrderProcessor { // 线程安全集合,专门用于并发场景下收集响应 private readonly ConcurrentBag<ApiResponse> _allApiResponses = new(); // SemaphoreSlim:限制同时发起的API请求数,避免触发API限流 private readonly SemaphoreSlim _apiRequestSemaphore; public OrderProcessor(int maxConcurrentApiCalls = 8) { _apiRequestSemaphore = new SemaphoreSlim(maxConcurrentApiCalls); } public List<ApiResponse> ProcessAllOrders(List<SalesOrder> totalOrders, int batchSize = 1000) { // 自动获取CPU核心数,作为并发上限的最优值 int maxParallelism = Environment.ProcessorCount; Parallel.ForEach(totalOrders.Batch(batchSize), new ParallelOptions { MaxDegreeOfParallelism = maxParallelism }, batch => { // 先获取Semaphore许可,再发起API调用 _apiRequestSemaphore.Wait(); try { // 调用实际的RestAPI,替换成你的业务逻辑 ApiResponse response = CallSalesOrderApi(batch); // 线程安全地添加响应结果 _allApiResponses.Add(response); } catch (Exception ex) { // 异常处理:记录错误并添加失败响应,保证流程不中断 _allApiResponses.Add(new ApiResponse { IsSuccess = false, ResultMessage = $"Batch failed: {ex.Message}", BatchItemCount = batch.Count() }); } finally { // 无论成功失败,都要释放Semaphore许可 _apiRequestSemaphore.Release(); } }); // 转换为普通List返回,方便后续处理 return _allApiResponses.ToList(); } // 模拟RestAPI调用,你需要替换成实际的HTTP请求逻辑 private ApiResponse CallSalesOrderApi(IEnumerable<SalesOrder> batch) { // 示例:用HttpClient发送POST请求(实际开发建议复用HttpClient实例) using var client = new HttpClient(); var batchJson = JsonSerializer.Serialize(batch); var content = new StringContent(batchJson, System.Text.Encoding.UTF8, "application/json"); var response = client.PostAsync("https://your-api-endpoint.com/process-batch", content).Result; response.EnsureSuccessStatusCode(); return JsonSerializer.Deserialize<ApiResponse>(response.Content.ReadAsStringAsync().Result); } } // 保留你提供的Batch扩展方法,优化了参数校验 public static class LinqExtensions { public static IEnumerable<IEnumerable<TSource>> Batch<TSource>(this IEnumerable<TSource> source, int size) { if (source == null) throw new ArgumentNullException(nameof(source)); if (size <= 0) throw new ArgumentOutOfRangeException(nameof(size), "Batch size must be greater than 0"); TSource[] bucket = null; var count = 0; foreach (var item in source) { bucket ??= new TSource[size]; bucket[count++] = item; if (count != size) continue; yield return bucket; bucket = null; count = 0; } // 返回最后一批不足指定大小的订单 if (bucket != null && count > 0) yield return bucket.Take(count); } }
3. 关键细节说明
- 线程安全的结果收集:用
ConcurrentBag<ApiResponse>替代普通List<T>,因为普通List不支持多线程同时写入,会导致数据丢失或异常;ConcurrentBag是.NET专为并发场景设计的集合,性能优于手动加lock的List。 - Semaphore的作用:
SemaphoreSlim用于额外限制API请求的并发数,比如即使Parallel.ForEach设置了8个并发,也可以通过Semaphore把API请求限制到5个,避免触发第三方API的限流规则。如果你的API没有严格并发限制,可以去掉Semaphore,直接用Parallel的并发控制。 - 并发数设置:用
Environment.ProcessorCount自动获取CPU核心数,比硬编码8更适配不同运行环境,比如在16核服务器上会自动调整为16。 - 异常处理:在API调用的
try-catch块中捕获异常,确保某一批次调用失败时,整个处理流程不会中断,同时记录错误信息便于后续排查。
4. 异步优化建议(.NET 6+)
如果你的API支持异步调用,用Parallel.ForEachAsync结合异步API会更高效——异步调用不会占用线程等待响应,能更充分利用系统资源:
public async Task<List<ApiResponse>> ProcessAllOrdersAsync(List<SalesOrder> totalOrders, int batchSize = 1000) { int maxParallelism = Environment.ProcessorCount; var allResponses = new ConcurrentBag<ApiResponse>(); var apiSemaphore = new SemaphoreSlim(maxParallelism); await Parallel.ForEachAsync(totalOrders.Batch(batchSize), new ParallelOptions { MaxDegreeOfParallelism = maxParallelism }, async (batch, cancellationToken) => { await apiSemaphore.WaitAsync(cancellationToken); try { ApiResponse response = await CallSalesOrderApiAsync(batch, cancellationToken); allResponses.Add(response); } catch (Exception ex) { allResponses.Add(new ApiResponse { IsSuccess = false, ResultMessage = $"Batch failed: {ex.Message}", BatchItemCount = batch.Count() }); } finally { apiSemaphore.Release(); } }); return allResponses.ToList(); } private async Task<ApiResponse> CallSalesOrderApiAsync(IEnumerable<SalesOrder> batch, CancellationToken cancellationToken) { // 复用HttpClient实例(建议在全局单例复用,不要每次调用都新建) using var client = new HttpClient(); var batchJson = JsonSerializer.Serialize(batch); var content = new StringContent(batchJson, System.Text.Encoding.UTF8, "application/json"); var response = await client.PostAsync("https://your-api-endpoint.com/process-batch", content, cancellationToken); response.EnsureSuccessStatusCode(); return await JsonSerializer.DeserializeAsync<ApiResponse>(await response.Content.ReadAsStreamAsync(cancellationToken), cancellationToken); }
这个异步版本特别适合IO密集型的API调用场景,能显著提升处理30万条订单的效率。
内容的提问来源于stack exchange,提问作者venkat
相关产品推荐
相关产品推荐

