如何在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
相关产品推荐
相关产品推荐

