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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 09:58:10