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

并行分发请求时接口响应越来越慢,现有实现存在什么问题?

现存问题

  • C# 事件默认是同步执行的:你在RequestList.Add方法中触发RequestReceived事件时,会在当前控制器请求线程上,同步执行所有订阅者的回调方法(也就是你写的RequestList_RequestReceived,包括日志打印、Task创建逻辑),所有回调执行完才会返回Add方法,直接拖慢接口响应速度。
  • List<T>不是线程安全的:多请求并发调用Add、Remove方法时,会出现内部数组扩容冲突、索引覆盖、线程自旋等待等问题,并发量越高,等待时间越长,正好符合你说的「请求越多响应越慢」的现象。
  • Task.Factory.StartNew使用不当:如果你的Distribute方法存在同步阻塞逻辑(比如用同步HTTP请求而非异步),会大量占用线程池线程,高并发下线程池耗尽后,新任务需要排队等待线程释放,不仅拖累分发逻辑,还会抢占ASP.NET Core本身的请求处理线程。

最优实现方案

推荐用.NET 原生的Channel<T>实现高性能生产者消费者模式,完全适配你的需求:控制器作为生产者只负责写队列,写完立即返回;后台托管服务作为消费者异步处理分发逻辑,全程无阻塞,线程安全,性能远高于自定义List+事件的实现。

具体实现步骤

1. 注册服务(Program.cs)

// 注册请求队列Channel,可根据需求调整为有界队列避免内存溢出
builder.Services.AddSingleton(Channel.CreateUnbounded<DistributionRequest>(new UnboundedChannelOptions
{
    SingleReader = true, // 单消费者场景开启可提升性能
    SingleWriter = false // 控制器请求为多生产者
}));
// 注册线程安全的结果存储
builder.Services.AddSingleton<ConcurrentDictionary<string, DistributionResult>>();
// 注册后台分发服务
builder.Services.AddHostedService<DistributionBackgroundService>();
// 注册HttpClientFactory,避免手动创建HttpClient的资源泄漏问题
builder.Services.AddHttpClient();

2. 控制器代码(仅做入队操作,毫秒级返回)

[ApiController]
[Route("api/distribute")]
public class DistributionController : ControllerBase
{
    private readonly ChannelWriter<DistributionRequest> _channelWriter;
    private readonly ConcurrentDictionary<string, DistributionResult> _resultStore;

    public DistributionController(ChannelWriter<DistributionRequest> channelWriter, 
        ConcurrentDictionary<string, DistributionResult> resultStore)
    {
        _channelWriter = channelWriter;
        _resultStore = resultStore;
    }

    [HttpPost]
    public IActionResult Post([FromForm] DistributionRequest request)
    {
        // 生成唯一请求ID
        request.RequestId = Guid.NewGuid().ToString();
        // 初始化结果为处理中
        _resultStore.TryAdd(request.RequestId, new DistributionResult { Status = ProcessStatus.Processing });
        // 写入队列,非阻塞操作
        _channelWriter.TryWrite(request);
        // 直接返回响应
        return Ok(new DistributionResponse { Succeeded = true, RequestId = request.RequestId });
    }

    // 补充查询结果的接口示例
    [HttpGet("result/{requestId}")]
    public IActionResult GetResult(string requestId)
    {
        if (_resultStore.TryGetValue(requestId, out var result))
        {
            return Ok(result);
        }
        return NotFound("请求ID不存在");
    }
}

3. 后台分发服务代码(异步处理逻辑,不阻塞请求线程)

public class DistributionBackgroundService : BackgroundService
{
    private readonly ChannelReader<DistributionRequest> _channelReader;
    private readonly ConcurrentDictionary<string, DistributionResult> _resultStore;
    private readonly IHttpClientFactory _httpClientFactory;
    private readonly ILogger<DistributionBackgroundService> _logger;

    public DistributionBackgroundService(ChannelReader<DistributionRequest> channelReader, 
        ConcurrentDictionary<string, DistributionResult> resultStore,
        IHttpClientFactory httpClientFactory,
        ILogger<DistributionBackgroundService> logger)
    {
        _channelReader = channelReader;
        _resultStore = resultStore;
        _httpClientFactory = httpClientFactory;
        _logger = logger;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        // 异步监听队列,无请求时挂起不占用系统资源
        await foreach (var request in _channelReader.ReadAllAsync(stoppingToken))
        {
            // 丢到异步任务处理,不阻塞队列消费,可通过SemaphoreSlim控制并发度避免打挂下游
            _ = ProcessRequestAsync(request, stoppingToken);
        }
    }

    private async Task ProcessRequestAsync(DistributionRequest request, CancellationToken stoppingToken)
    {
        try
        {
            _logger.LogInformation("开始处理请求 {RequestId}", request.RequestId);
            var client = _httpClientFactory.CreateClient();
            // 这里替换为你实际的分发逻辑,所有IO操作都用异步方法,不要加同步阻塞代码
            var response = await client.PostAsJsonAsync("你的下游endpoint地址", request, stoppingToken);
            // 更新处理结果
            _resultStore.TryUpdate(request.RequestId, new DistributionResult
            {
                Status = ProcessStatus.Completed,
                Content = await response.Content.ReadAsStringAsync(stoppingToken)
            }, _resultStore[request.RequestId]);
            _logger.LogInformation("请求 {RequestId} 处理完成", request.RequestId);
        }
        catch (Exception ex)
        {
            _logger.LogError(ex, "请求 {RequestId} 处理失败", request.RequestId);
            // 更新失败状态
            _resultStore.TryUpdate(request.RequestId, new DistributionResult
            {
                Status = ProcessStatus.Failed,
                ErrorMessage = ex.Message
            }, _resultStore[request.RequestId]);
        }
    }
}

额外优化建议

  • 如果担心内存溢出,可以将Channel改为有界队列,设置最大容量,超过容量时可选择拒绝请求或者持久化到Redis/数据库。
  • 结果存储可以添加过期清理逻辑,比如用定时任务删除24小时以上的历史结果,避免内存持续上涨。
  • 如果需要服务重启不丢失排队请求,可以在写入Channel前先将请求持久化到存储,消费完成后再删除,只会增加极少量的接口响应耗时,远优于你当前的实现。
  • 分发逻辑如果需要控制并发量,可以用SemaphoreSlim限制同时处理的请求数,避免下游接口被压垮。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 20:45:02