C#如何实现保持顺序的并行映射迭代器?
并行处理Enumerator并保持输出顺序的实现方案
要解决并行处理高成本映射任务且严格保留输入顺序的问题,核心是跟踪每个元素的原始索引,并行处理后按索引顺序输出,同时控制内存缓冲只保留未到输出时机的结果。OrderablePartitioner确实是实现这个需求的合适工具,下面是具体的实现思路和代码示例:
核心思路
- 给输入枚举器的每个元素分配唯一的递增索引,标记其原始顺序。
- 使用
OrderablePartitioner创建支持顺序跟踪的分区器,避免提前缓冲过多元素(减少内存占用)。 - 并行处理每个带索引的元素,将结果和对应索引存入线程安全的容器。
- 主线程按索引递增顺序提取结果,仅缓冲那些处理完成但还未到输出顺序的项,不会缓存整个序列。
代码实现
using System; using System.Collections; using System.Collections.Generic; using System.Collections.Concurrent; using System.Threading; using System.Threading.Tasks; public static class OrderedParallelExtensions { public static IEnumerator<S> ParallelMap<T, S>(this IEnumerator<T> input, Func<T, S> mapFunc, int maxDegreeOfParallelism = -1) { if (input == null) throw new ArgumentNullException(nameof(input)); if (mapFunc == null) throw new ArgumentNullException(nameof(mapFunc)); // 将输入枚举器转换为带索引的序列,同时释放原枚举器 var indexedItems = new List<(int Index, T Item)>(); using (input) { int index = 0; while (input.MoveNext()) { indexedItems.Add((index++, input.Current)); } } if (indexedItems.Count == 0) yield break; // 创建不缓冲的OrderablePartitioner,减少内存占用 var partitioner = OrderablePartitioner.Create( indexedItems, EnumerablePartitionerOptions.NoBuffering); var resultQueue = new ConcurrentQueue<(int Index, S Result)>(); var processingCompleted = new ManualResetEventSlim(false); // 后台并行处理任务 _ = Task.Run(() => { try { Parallel.ForEach(partitioner, new ParallelOptions { MaxDegreeOfParallelism = maxDegreeOfParallelism }, item => { // 执行高成本映射操作 var result = mapFunc(item.Item); resultQueue.Enqueue((item.Index, result)); }); } finally { processingCompleted.Set(); } }); var pendingResults = new Dictionary<int, S>(); int currentOutputIndex = 0; while (true) { // 把队列中的结果转移到待处理字典 while (resultQueue.TryDequeue(out var indexedResult)) { pendingResults[indexedResult.Index] = indexedResult.Result; } // 输出连续的已完成结果 while (pendingResults.TryGetValue(currentOutputIndex, out var readyResult)) { yield return readyResult; pendingResults.Remove(currentOutputIndex); currentOutputIndex++; } // 检查是否所有任务完成且无待处理结果 if (processingCompleted.IsSet && resultQueue.IsEmpty && pendingResults.Count == 0) { break; } // 短暂等待避免空轮询,可根据需求调整时长 Thread.Sleep(10); } } }
关键细节说明
- OrderablePartitioner的作用:它会确保每个元素的原始索引被正确传递给并行处理任务,不会打乱顺序标记。
NoBuffering选项让分区器不会提前读取过多元素到内存,符合你不想缓存整个序列的需求。 - 内存控制:只有那些处理完成但还没到输出顺序的结果会被存在
pendingResults字典中,内存占用取决于并发处理的任务数,而非整个输入序列的大小。 - 线程安全:用
ConcurrentQueue收集并行结果,避免锁竞争;主线程的结果输出逻辑通过字典保证顺序,无需额外锁(因为主线程是唯一的写入和读取者)。
使用示例
// 假设有一个输入枚举器 var inputEnumerator = GetInputEnumerator(); // 并行映射,保留顺序 var resultEnumerator = inputEnumerator.ParallelMap(item => ExpensiveMapFunction(item), maxDegreeOfParallelism: 4); while (resultEnumerator.MoveNext()) { Console.WriteLine(resultEnumerator.Current); }
内容的提问来源于stack exchange,提问作者Michael Kay
相关产品推荐
相关产品推荐

