并行分发请求时接口响应越来越慢,现有实现存在什么问题?
现存问题
- 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
相关产品推荐
相关产品推荐

