ASP.NET Core 6 Web API实现同ID串行、异ID并行处理
实现同一ID串行、不同ID并行的处理方案
针对你的需求,这里提供几种实用的实现方案,适配ASP.NET Core 6环境:
方案一:分组后异步任务调度(一次性处理列表场景)
如果是一次性处理已有的数据列表,先按Id分组,再为每个分组启动独立的异步任务,任务内部串行处理组内数据,组间自动并行:
// 按Id分组 var groupedData = l.GroupBy(d => d.Id); // 为每个分组创建处理任务 var processingTasks = groupedData.Select(group => Task.Run(async () => { foreach (var data in group) { // 替换为你的实际处理逻辑 await ProcessData(data); Console.WriteLine($"完成Id={data.Id}的处理,内容:{data.o}"); } }, cancellationToken)); // 等待所有分组处理完成 await Task.WhenAll(processingTasks);
如果需要控制并行度,可用Parallel.ForEachAsync(.NET 6+支持):
await Parallel.ForEachAsync(groupedData, new ParallelOptions { MaxDegreeOfParallelism = Environment.ProcessorCount }, async (group, ct) => { foreach (var data in group) { await ProcessData(data); Console.WriteLine($"完成Id={data.Id}的处理,内容:{data.o}"); } });
方案二:多Channel+BackgroundService(持续接收请求场景)
如果Web API需要持续接收请求并实时处理,可给每个Id创建独立的Channel,每个Channel对应一个消费者任务,保证同一Id的请求串行处理,不同Id的请求并行处理:
1. 实现数据处理服务
创建单例服务管理各Id的Channel和消费者:
public class DataProcessingService { private readonly Dictionary<int, Channel<Data>> _channelMap = new(); private readonly SemaphoreSlim _semaphore = new(1, 1); // 入队待处理数据 public async Task EnqueueData(Data data, CancellationToken cancellationToken) { await _semaphore.WaitAsync(cancellationToken); try { // 不存在对应Id的Channel则创建,并启动消费者 if (!_channelMap.TryGetValue(data.Id, out var channel)) { channel = Channel.CreateUnbounded<Data>(); _channelMap[data.Id] = channel; _ = StartConsumer(channel, data.Id, cancellationToken); } } finally { _semaphore.Release(); } // 将数据写入对应Channel await _channelMap[data.Id].Writer.WriteAsync(data, cancellationToken); } // 启动对应Id的消费者任务 private async Task StartConsumer(Channel<Data> channel, int id, CancellationToken cancellationToken) { try { // 串行读取并处理当前Id的所有数据 await foreach (var data in channel.Reader.ReadAllAsync(cancellationToken)) { await ProcessData(data); Console.WriteLine($"完成Id={id}的处理,内容:{data.o}"); } } finally { // 消费者结束后清理Channel await _semaphore.WaitAsync(); try { _channelMap.Remove(id); } finally { _semaphore.Release(); } } } // 替换为你的实际处理逻辑 private async Task ProcessData(Data data) { // 模拟处理耗时 await Task.Delay(500); } }
2. 注册服务
在Program.cs中将服务注册为单例:
builder.Services.AddSingleton<DataProcessingService>();
3. 在Controller中使用
[ApiController] [Route("api/[controller]")] public class DataController : ControllerBase { private readonly DataProcessingService _processingService; public DataController(DataProcessingService processingService) { _processingService = processingService; } [HttpPost] public async Task<IActionResult> PostData([FromBody] Data data) { await _processingService.EnqueueData(data, HttpContext.RequestAborted); return Ok("数据已入队待处理"); } }
方案三:Parallel.ForEach同步处理(同步场景)
如果处理逻辑是同步的,可直接用Parallel.ForEach处理分组后的数据:
var groupedData = l.GroupBy(d => d.Id).ToList(); Parallel.ForEach(groupedData, new ParallelOptions { MaxDegreeOfParallelism = 4 }, group => { foreach (var data in group) { // 同步处理逻辑 ProcessDataSync(data); Console.WriteLine($"完成Id={data.Id}的处理,内容:{data.o}"); } }); // 同步处理方法 private void ProcessDataSync(Data data) { Thread.Sleep(1000); }
为什么之前的单个Channel方案串行?
你之前用单个Channel<Data>+BackgroundService,所有数据都进入同一个Channel,只有一个消费者线程读取处理,自然是全串行。而多Channel方案为每个Id分配独立的消费队列,不同Id的队列并行消费,同一Id的队列串行处理,正好匹配你的需求。
内容的提问来源于stack exchange,提问作者peterpie
相关产品推荐
相关产品推荐

