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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 20:03:21