【问题标题】:AsParallel() and internal buffer sizeAsParallel() 和内部缓冲区大小
【发布时间】:2011-12-07 15:53:37
【问题描述】:

如何限制 AsParallel() 预先读取并放入其内部缓冲区的项目数量?

这是一个例子:

int returnedCounter;

IEnumerable<int> Enum()
{
    while (true)
        yield return Interlocked.Increment(ref returnedCounter);
}

[TestMethod]
public void TestMethod1()
{
    foreach (var i in Enum().AsParallel().Select(a => a))
    {
        Thread.Sleep(3000);
        break;
    }
    Console.WriteLine(returnedCounter);
}

我消耗了 1 个项目,睡眠,停止枚举。它在我的机器上打印 526400。在我的真实项目中,每个项目分配数千字节。 AsParallel() 会预先读取大量项目,这会导致非常糟糕的内存消耗和 CPU 浪费。

Put WithMergeOptions(ParallelMergeOptions.NotBuffered) 有点帮助。它打印 4544。但对我来说仍然太多了。

在 Enum() 中等待会冻结主线程中的循环。

【问题讨论】:

    标签: .net multithreading plinq


    【解决方案1】:

    关于Partitioners的另一个问题!

    在您的情况下,您将必须找到/编写一个一次只需要一个项目的分区器。

    这是一篇关于Custom Partitioners的文章


    更新:

    我只记得我在哪里看到了 SingleItemPartitioner 实现:它在 ParallelExtensionsExtras 项目中:Samples for Parallel Programming with the .NET Framework

    我也刚刚阅读了您的测试代码。我可能应该第一次这样做!

    这段代码:

    Enum().AsParallel().Select(a => a)
    

    意思是:取Enum()并尽可能快地并行枚举它,并返回一个新的IEnumerable&lt;int&gt;

    因此,您的 foreach 不是从 Enum() 中提取项目 - 它是从 linq 语句创建的新 IEnumerable&lt;int&gt; 中提取项目。

    另外,您的foreach 在主线程上运行,因此每个项目的工作都是单线程的。

    如果您想并行运行,但只在需要时产生一个项目,请尝试:

    Parallel.ForEach( SingleItemPartitioner.Create( Enum() ), ( i, state ) =>
        {
            Thread.Sleep( 3000 );
            state.Break();
        }
    

    【讨论】:

    • 我采用了来自msdn.microsoft.com/en-us/library/dd997416.aspx 的每个分区一个项目的分区器。还从您指出的示例中尝试了 SingleItemPartitioner。代码类似于 `foreach (var i in SingleItemPartitioner.Create(Enum()).AsParallel().Select(a => a))` 相同的输出 - 超过 500000 个项目被预先读取并由 plinq 缓存。
    • > 所以每个项目的工作都是单线程的。正确的。但也有一部分工作是并行完成的。让我澄清一下我的问题。我的代码模拟了常见的 plinq 使用场景 `foreach (var i in Enum().AsParallel().Select(a => DoExpensiveWorkInParallel(a))) DoRestOfWorkSynchronously(i);'
    • 问题是同步部分可能会等待一些事件,而不是消耗由 ParallelQuery.Select() 生成的项目。我希望 plinq 停止从源 IEnumerable(由 Enum() 返回)获取项目并暂停并行处理。 plinq 实际上是这样做的,但首先它会用超过 500000 个项目填充一些内部队列。我的问题是如何让 plinq 预先处理更少的项目。
    • PLinq 旨在尽可能快地处理工作,而您要求它做的只是生成一个 IEnumerable。在您问题的示例中,DoExpensiveWorkInParallel 很快,DoRestOfWorkSynchronously 很慢。
    • PLinq 创建传送带或管道。 DoRestOfWorkSynchronously 确实很慢 - 它使传送带暂停 3 秒。输送机的平行部分不断生产和堆叠产品到内部队列。队列的大小似乎根本无法配置:(
    【解决方案2】:

    找到了解决方法。

    首先,让我澄清一下最初的问题。我需要一个可在无限序列上工作的可暂停管道。管道是:

    1. 同步读取序列:Enum()
    2. 并行处理项目:AsParallel().Select(a =&gt; a)
    3. 继续同步处理:foreachbody

    第 3 步可能会暂停流水线。这是由Sleep() 模拟的。问题是当管道暂停时,第 2 步提前获取了太多元素。 Plinq 必须有一些内部队列。队列大小不能显式配置。不过,大小取决于ParallelMergeOptionsParallelMergeOptions.NotBuffered 降低了队列大小,但对我来说还是太大了。

    我的解决方法是知道有多少项目正在处理,达到限制时停止并行处理,当管道再次启动时重新启动并行处理。

    int sourceCounter;
    
    IEnumerable<int> SourceEnum() // infinite input sequence
    {
        while (true)
            yield return Interlocked.Increment(ref sourceCounter);
    }
    
    [TestMethod]
    public void PlainPLinq_PausedConsumtionTest()
    {
        sourceCounter = 0;
        foreach (var i in SourceEnum().AsParallel().WithMergeOptions(ParallelMergeOptions.NotBuffered).Select(a => a))
        {
            Thread.Sleep(3000);
            break;
        }
        Console.WriteLine("fetched from source sequence: {0}", sourceCounter); // prints 4544 on my machine
    }
    
    [TestMethod]
    public void MyParallelSelect_NormalConsumtionTest()
    {
        sourceCounter = 0;
        foreach (var i in MyParallelSelect(SourceEnum(), 64, a => a))
        {
            if (sourceCounter > 1000000)
                break;
        }
        Console.WriteLine("fetched from source sequence: {0}", sourceCounter);
    }
    
    [TestMethod]
    public void MyParallelSelect_PausedConsumtionTest()
    {
        sourceCounter = 0;
        foreach (var i in MyParallelSelect(SourceEnum(), 64, a => a))
        {
            Thread.Sleep(3000);
            break;
        }
        Console.WriteLine("fetched from source sequence: {0}", sourceCounter);
    }
    
    class DataHolder<D> // reference type to store class or struct D
    {
        public D Data;
    }
    
    static IEnumerable<DataHolder<T>> FetchSourceItems<T>(IEnumerator<T> sourceEnumerator, DataHolder<int> itemsBeingProcessed, int queueSize)
    {
        for (; ; )
        {
            var holder = new DataHolder<T>();
            if (Interlocked.Increment(ref itemsBeingProcessed.Data) > queueSize)
            {
                // many enought items are already being processed - stop feeding parallel processing
                Interlocked.Decrement(ref itemsBeingProcessed.Data);
                yield break;
            }
            if (sourceEnumerator.MoveNext())
            {
                holder.Data = sourceEnumerator.Current;
                yield return holder;
            }
            else
            {
                yield return null; // return null DataHolder to indicate EOF
                yield break;
            }
        }
    }
    
    IEnumerable<OutT> MyParallelSelect<T, OutT>(IEnumerable<T> source, int queueSize, Func<T, OutT> selector)
    {
        var itemsBeingProcessed = new DataHolder<int>();
        using (var sourceEnumerator = source.GetEnumerator())
        {
            for (;;) // restart parallel processing
            {
                foreach (var outData in FetchSourceItems(sourceEnumerator, itemsBeingProcessed, queueSize).AsParallel().WithMergeOptions(ParallelMergeOptions.NotBuffered).Select(
                    inData => inData != null ? new DataHolder<OutT> { Data = selector(inData.Data) } : null))
                {
                    Interlocked.Decrement(ref itemsBeingProcessed.Data);
                    if (outData == null)
                        yield break; // EOF reached
                    yield return outData.Data;
                }
            }
        }
    }
    

    【讨论】:

      猜你喜欢
      • 2015-02-16
      • 2014-09-20
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-09-05
      相关资源
      最近更新 更多