【问题标题】:How do I use PLINQ to process chunks of an array如何使用 PLINQ 处理数组块
【发布时间】:2012-07-13 09:42:49
【问题描述】:

我需要获取大量的双精度数组,并使用处理器密集型函数将其分块处理。

我的原始数组非常大,大约 200MB 的双信号数据。

我需要将它分成 5000 个双精度块,使用返回单个双精度的函数处理具有一些处理器密集型数学的那些。需要这些函数中的每一个结果来创建稍后使用的有序数组。

认为这对于使用 PLINQ 进行并行处理来说是最佳的,但我不太确定如何去做。

我写的幼稚实现是这样的:

        var processedList = new List<double>();

        var chunk = new List<double>;
        foreach (var rawSample in drop.RawSamples)
        {
            chunk.Add(rawSample);

            if (chunk.Count == 5000)
            {
                // Do long processing here
                processedList.Add(LongProcessingFunction(chunk));

                chunk.Clear();
            }
        }

        // Do something later with the list of processed values.....

那么,我应该从哪里开始使用 PLINQ?我需要能够使用处理器的所有内核来执行长时间、密集的功能。

我看到 IEnumerable 有一个 Take(n) 函数.....我可以使用它吗?

我可以在这里使用 AsParallel 吗?

谢谢!

【问题讨论】:

  • 也许你可以重复使用Take(5000),把这个chunk的处理放到一个新的TaskWaitAll任务中。
  • 我想摆脱管理自己的线程池。我认为创建几千个任务然后等待它们都不会很有效。我认为应该有一些 plinq 语法可以为我做到这一点?
  • 如果您的日程安排不是太复杂,我认为使用TaskFactory.StartNew 提供的机制就可以了,也许可以使用TaskCreationOptions.LongRunning

标签: .net linq c#-4.0 parallel-processing plinq


【解决方案1】:

首先,如果您要处理如此大量的数据,则应尽可能避免逐个元素地处理它。在您的代码中,您可以通过以 5000 为增量迭代整数并使用 Array.Copy() 之类的东西来做到这一点。

或者,更好的是,根本不进行任何复制,让LongProcessingFunction 接受一个数组(或IList&lt;T&gt;,或IReadOnlyList&lt;T&gt;,如果您使用的是.Net 4.5;但使用接口确实有一些开销)和该数组的偏移量。

如果你想使你的代码并行,你可以使用ParallelEnumerable.Range()AsOrdered()(这是使结果按正确顺序所必需的)和Select()

double[] result = ParallelEnumerable.Range(0, drop.RawSamples.Length / chunkSize)
    .AsOrdered()
    .Select(i => LongProcessingFunction(drop.RawSamples, i * chunkSize))
    .ToArray();

【讨论】:

    【解决方案2】:

    我建议在您的索引源上实现一个分区器,它将您的源分成 5000 个元素的块。然后可以使用AsParallel 并行处理它们中的每一个。

    class Program
    {
        static void Main(string[] args)
        {
    
            IList<double> rawData =  [Your raw data here];
    
            IList<double> result =
                rawData
                    .Partition(5000)
                    .AsParallel()
                    .AsOrdered()
                    .Select(chunk => LongProcessingFunction(chunk))
                    .ToList();
        }
    
        private static double LongProcessingFunction(IList<double> chunk)
        {
            throw new NotImplementedException();
        }
    }
    
    public static class MyExtensions
    {
        public static IEnumerable<List<T>> Partition<T>(this IList<T> source, Int32 size)
        {
            for (int i = 0; i < Math.Ceiling(source.Count / (Double)size); i++)
            {
                yield return new List<T>(source.Skip(size*i).Take(size));
            }
        }
    }
    

    【讨论】:

    • 自定义分区器可以提高效率,但并不意味着实现基本逻辑。
    • 在您的代码中,您实际上并没有实现partitioner,所以我认为您的方法的另一个名称会更好。如果我正确理解了这个问题,您还需要在查询中添加AsOrdered()
    猜你喜欢
    • 1970-01-01
    • 2011-06-06
    • 2014-02-23
    • 2019-02-21
    • 1970-01-01
    • 1970-01-01
    • 2012-06-29
    • 2020-08-25
    • 1970-01-01
    相关资源
    最近更新 更多