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

C#如何实现保持顺序的并行映射迭代器?

并行处理Enumerator并保持输出顺序的实现方案

要解决并行处理高成本映射任务且严格保留输入顺序的问题,核心是跟踪每个元素的原始索引,并行处理后按索引顺序输出,同时控制内存缓冲只保留未到输出时机的结果。OrderablePartitioner确实是实现这个需求的合适工具,下面是具体的实现思路和代码示例:

核心思路

  1. 给输入枚举器的每个元素分配唯一的递增索引,标记其原始顺序。
  2. 使用OrderablePartitioner创建支持顺序跟踪的分区器,避免提前缓冲过多元素(减少内存占用)。
  3. 并行处理每个带索引的元素,将结果和对应索引存入线程安全的容器。
  4. 主线程按索引递增顺序提取结果,仅缓冲那些处理完成但还未到输出顺序的项,不会缓存整个序列。

代码实现

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 03:50:43