【问题标题】:How to enforce a sequence of ordered execution in parallel.for?如何在 parallel.for 中强制执行一系列有序执行?
【发布时间】:2021-01-18 14:06:16
【问题描述】:

我有一个简单的并行循环在做一些事情,然后我将结果保存到一个文件中。

object[] items; // array with all items
object[] resultArray = new object[numItems];
Parallel.For(0, numItems, (i) => 
{ 
    object res = doStuff(items[i], i);
    resultArray[i] = res;
});

foreach (object res in resultArray)
{
    sequentiallySaveResult(res);
}

为了节省,我需要按正确的顺序写入结果。通过将结果放入resultArray,结果的顺序又是正确的。

但是,由于结果非常大并且占用大量内存。 我想按顺序处理项目,例如四个线程启动并处理项目 1-4,下一个空闲线程处理项目 5,依此类推。

这样,我可以启动另一个线程,监视数组中接下来需要写入的项目(或者每个线程可以在项目完成时发出一个事件),所以我已经可以开始编写第一个结果,而稍后的项目仍在处理中,然后释放内存。

Parallel.For 是否可以按给定顺序处理项目?我当然可以使用concurentQueue,将所有索引按正确的顺序放在那里,然后手动启动线程。

但如果可能的话,我想保留 'Parallel.For' 实现中关于使用多少线程等的所有自动化。

免责声明:我无法切换到ForEach,我需要i。

编辑 #1:
目前,执行顺序是完全随机的,举个例子:

Processing item 1/255
Processing item 63/255
Processing item 32/255
Processing item 125/255
Processing item 94/255
Processing item 156/255
Processing item 187/255
Processing item 249/255
...

编辑 #2:
有关已完成工作的更多详细信息:

我处理一个灰度图像,需要为每个“层”(上例中的项目)提取信息,所以我从 0 到 255(对于 8 位)并在图像上执行任务。

我有一个类可以同时访问像素值:

 unsafe class UnsafeBitmap : IDisposable
    {

        private BitmapData bitmapData;
        private Bitmap gray;
        private int bytesPerPixel;
        private int heightInPixels;
        private int widthInBytes;
        private byte* ptrFirstPixel;

        public void PrepareGrayscaleBitmap(Bitmap bitmap, bool invert)
        {
            gray = MakeGrayscale(bitmap, invert);

            bitmapData = gray.LockBits(new Rectangle(0, 0, gray.Width, gray.Height), ImageLockMode.ReadOnly, gray.PixelFormat);
            bytesPerPixel = System.Drawing.Bitmap.GetPixelFormatSize(gray.PixelFormat) / 8;
            heightInPixels = bitmapData.Height;
            widthInBytes = bitmapData.Width * bytesPerPixel;
            ptrFirstPixel = (byte*)bitmapData.Scan0;
        }

        public byte GetPixelValue(int x, int y)
        {
            return (ptrFirstPixel + ((heightInPixels - y - 1) * bitmapData.Stride))[x * bytesPerPixel];
        }

        public void Dispose()
        {
            gray.UnlockBits(bitmapData);
        }
    }

循环是

UnsafeBitmap ubmp; // initialized, has the correct bitmap
int numLayers = 255;
int bitmapWidthPx = 10000;
int bitmapHeightPx = 10000;
object[] resultArray = new object[numLayer];
Parallel.For(0, numLayers, (i) => 
{ 
        for (int x = 0; x < bitmapWidthPx ; x++)
    {
        inLine = false;
        for (int y = 0; y < bitmapHeightPx ; y++)
        {
            byte pixel_value = ubmp.GetPixelValue(x, y);
            
            if (i <= pixel_value && !inLine)
            {
                result.AddStart(x,y);
                inLine = true;
            }
            else if ((i > pixel_value || y == Height - 1) && inLine)
            {
                result.AddEnd(x, y-1);
                inLine = false;
            }
        }
    }
    result_array[i] = result;
});

foreach (object res in resultArray)
{
    sequentiallySaveResult(res);
}

我还想启动一个线程进行保存,检查下一个需要写入的项目是否可用,写入它,从内存中丢弃。为此,最好按顺序开始处理,以便结果大致按顺序到达。如果第 5 层的结果倒数第二个到达,我必须等待写入第 5 层(以及所有后续)直到最后。

如果有 4 个线程启动,则开始处理第 1-4 层,当一个线程完成后,开始处理第 5 层,下一个第 6 层,依此类推,结果将或多或少以相同的顺序出现,我可以开始将结果写入文件并从内存中丢弃。

【问题讨论】:

  • 您可以切换到Parallel.ForEach 并使用提供索引的重载。此外,如果您已经拥有所有适当的项目,则不需要索引。排序是额外的工作,这就是为什么 Parallel.For、Foreach 或 PLINQ 都不会产生有序的结果。您可以通过将AsOrdered 添加到 PLINQ 查询来请求排序结果
  • 你不需要也不应该不尝试使用你自己的线程。对于初学者,PLINQ 和 Parallel 使用所有可用的内核。他们还处理数据分区、负载平衡、批处理等,这比尝试自己做同样的事情要高效得多
  • 正如所写,我什至不需要得到完美排序的结果。它应该只是按给定的顺序开始处理,所以结果以正确的顺序到达,所以我可以开始保存它们了。如果第 2 项的结果先到,第 1 项紧随其后,没关系,我可以在 writer 中处理,它会等待第一项,然后按正确的顺序写入。
  • 不会的。当您使用并行处理时,结果将以 CPU 生成它们的任何顺序出现。如果其中一项工作任务比其他任务工作得更快,也许是因为它的工作更容易,它会比其他任务产生更多的结果。这就是为什么必须再次订购结果。你确定你需要 parallelism 吗?并行意味着处理大量内存中数据。如果您需要加载文件或执行多个步骤,您可能正在寻找管道处理
  • 每项任务都需要相当精确的时间。但即使在任务上工作得更快,也没关系。所以让我们说 4 个线程开始,开始处理项目 1-4。线程 2 首先完成,然后开始处理项目 5,依此类推。所以结果将大致以相同的正确顺序出现,所以我可以开始按顺序编写它们。但是,如果第 5 项仅是倒数第二项,我必须等到一切都完成后才能编写第 5 项和以下所有内容。

标签: c# multithreading .net-core seq concurrent-processing


【解决方案1】:

Parallel 类知道如何并行化工作负载,但不知道如何合并处理后的结果。所以我建议改用PLINQ。您要求以原始顺序保存结果并与处理同时进行,这比平时有点棘手,但它仍然是完全可行的:

IEnumerable<object> results = Partitioner
    .Create(items, EnumerablePartitionerOptions.NoBuffering)
    .AsParallel()
    .AsOrdered()
    .WithMergeOptions(ParallelMergeOptions.NotBuffered)
    .Select((item, index) => DoStuff(item, index))
    .AsEnumerable();

foreach (object result in results)
{
    SequentiallySaveResult(result);
}

解释:

  1. AsOrdered 运算符是按原始顺序检索结果所必需的。
  2. WithMergeOptions 运算符是防止结果缓冲所必需的,以便在结果可用时立即保存。
  3. Partitioner.Create 是必需的,因为数据源是一个数组,而 PLINQ 默认情况下会静态地对数组进行分区。这意味着数组被分成多个范围,并分配一个线程来处理每个范围。一般来说,这是一个很好的性能优化,但在这种情况下,它违背了及时有序地检索结果的目的。所以需要一个动态分区器,从头到尾依次枚举源。
  4. EnumerablePartitionerOptions.NoBuffering 配置可防止 PLINQ 使用的工作线程一次抓取多个项目(这是默认的 PLINQ 分区技巧,称为“块分区”)。
  5. AsEnumerable 并不是真正需要的。它只是为了表示并行处理的结束。后面的foreach 将ParallelQuery&lt;object&gt; 视为IEnumerable&lt;object&gt;。

由于需要所有这些技巧,并且由于此解决方案不够灵活,以防您稍后需要在处理管道中添加更多并发异构步骤,因此我建议您牢记升级到TPL Dataflow 图书馆。它是一个库,可在并行处理领域解锁许多强大的选项。

【讨论】:

  • here 介绍了使用 PLINQ 创建处理管道的更高级方法。
【解决方案2】:

如果您想对线程操作进行排序,线程同步 101 会教我们使用条件变量,并在 C# 任务中实现这些条件变量,您可以使用提供异步等待功能的SemaphoreSlim SemaphoreSlim.WaitAsync。再加上计数器检查,您将获得所需的结果。

但是我不相信它是必要的,因为如果我理解正确并且您只想按顺序保存它们以避免将它们存储在内存中,您可以使用内存映射文件来:

  1. 如果结果大小相同,只需将缓冲区写入位置index * size。

  2. 如果结果的大小不同,请在获得结果时写入临时映射文件,并让另一个线程在它们出现时复制正确的顺序输出文件。这是一个 IO 绑定操作,所以不要为此使用任务池。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2022-01-04
    • 2015-03-01
    • 2010-12-23
    • 2021-01-09
    • 2010-09-10
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多