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

如何在C#中高效实现并发处理数据项并按顺序写入结果

C# 高效实现并行读取、按序串行写入逻辑

核心约束

  • Process1(id):数据读取/预处理逻辑,支持多任务并行执行,是主要耗时环节
  • Process2(result):结果输出/写入逻辑,必须严格按照传入ID的原始顺序执行,不允许乱序
  • 原串行实现的性能瓶颈:逐ID等待Process1执行完成才进入下一步,并行能力完全闲置,IO等待阶段资源利用率极低

实现思路

对齐参考Python实现的核心逻辑,用异步队列做任务调度,固定数量的任务并行执行读取操作,读取到的结果暂存在线程安全字典中;写入端按ID顺序轮询检查,当前ID对应结果就绪就立刻执行写入,未就绪则等待新读取结果的通知,不会阻塞其他读取任务执行。
实现采用.NET原生Channel做异步生产消费队列,自带异步等待、背压能力,比手动实现的队列+事件组合更稳定高效。

完整实现代码

using System.Collections.Concurrent;
using System.Threading.Channels;

async Task ProcessAllIDs(List<int> ids, int parallelReadCount = 3)
{
    // 初始化待处理ID队列,设置容量上限做背压,避免读取过快占用过多内存
    var idQueue = Channel.CreateBounded<int>(new BoundedChannelOptions(capacity: parallelReadCount * 2)
    {
        FullMode = BoundedChannelFullMode.Wait
    });

    // 存储已经完成读取的结果,key为ID,value为Process1返回的处理结果
    var idToResult = new ConcurrentDictionary<int, ProcessResult>();
    // 异步信号:有新的读取结果就绪时触发
    var resultReadySignal = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
    
    // 启动固定数量的并行读取任务
    var readTasks = Enumerable.Range(0, parallelReadCount)
        .Select(_ => Task.Run(async () =>
        {
            await foreach (var id in idQueue.Reader.ReadAllAsync())
            {
                // 执行可并行的读取/处理逻辑
                var result = await Process1(id);
                idToResult[id] = result;
                // 触发信号通知写入端有新结果
                lock (resultReadySignal)
                {
                    if (!resultReadySignal.Task.IsCompleted)
                        resultReadySignal.SetResult();
                }
            }
        }))
        .ToArray();

    // 把所有ID写入队列供读取任务消费
    foreach (var id in ids)
    {
        await idQueue.Writer.WriteAsync(id);
    }
    idQueue.Writer.Complete(); // 所有ID入队后标记队列完成

    // 严格按ID顺序执行写入逻辑
    foreach (var targetId in ids)
    {
        while (true)
        {
            // 检查当前需要写入的ID结果是否就绪
            if (idToResult.TryRemove(targetId, out var result))
            {
                // 就绪则执行串行写入逻辑
                await Process2(result);
                break;
            }
            // 未就绪则等待新结果到来,重置信号后继续检查
            await resultReadySignal.Task;
            lock (resultReadySignal)
            {
                if (!idToResult.ContainsKey(targetId))
                    resultReadySignal = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
            }
        }
    }

    // 等待所有读取任务正常结束
    await Task.WhenAll(readTasks);
}

// --------------- 以下为示例逻辑,替换为实际业务代码即可 ---------------
async Task<ProcessResult> Process1(int id)
{
    // 模拟异步读取/轻量处理逻辑,如调用远程接口、读取磁盘文件、解析数据
    await Task.Delay(Random.Shared.Next(100, 500));
    return new ProcessResult { Id = id, Data = $"data_{id}" };
}

async Task Process2(ProcessResult result)
{
    // 模拟必须顺序执行的写入逻辑,如写入合并文件、提交顺序日志、按序落库
    await Task.Delay(50);
    Console.WriteLine($"ID={result.Id} 写入完成");
}

public class ProcessResult
{
    public int Id { get; set; }
    public string Data { get; set; }
}

调优说明

  • 并行度参数parallelReadCount按需调整:IO密集型场景(如HTTP接口调用、磁盘文件读取)可设置为5~20,CPU密集型场景建议设置为和CPU核心数一致
  • 队列容量默认设为并行度的2倍实现背压:如果读取速度远高于写入速度,不会无限制缓存结果占用内存,队列满时读取任务会自动等待
  • 信号操作加锁是为了避免多线程并发触发信号、重置信号时出现通知丢失,导致写入端永久等待
  • 全链路异步无阻塞,在Process1占主要耗时的场景下,吞吐可接近串行实现的parallelReadCount倍

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 20:33:39