【问题标题】:Parallelization of CPU bound task continuing with IO boundCPU 绑定任务的并行化继续与 IO 绑定
【发布时间】:2012-07-26 19:39:44
【问题描述】:

我正在尝试找出一种并行化处理大数据集的代码的好方法,然后将生成的数据导入 RavenDb。

数据处理受 CPU 限制和数据库导入 IO 限制。

我正在寻找一种在 Environment.ProcessorCount 线程数上并行处理的解决方案。然后将生成的数据导入到 RavenDb 上的 x(比如说 10)个池线程上,与上述过程并行。

这里的主要内容是我希望在导入完成的数据时继续处理,以便在等待导入完成时继续处理下一个数据集。

另一个问题是每个批次的内存在成功导入后需要丢弃,因为私有工作内存很容易达到 >5GB。

下面的代码是我到目前为止所得到的。请注意,它不能满足上述并行化要求。

datasupplier.GetDataItems()
    .Partition(batchSize)
    .AsParallel()
    .WithDegreeOfParallelism(Environment.ProcessorCount)
    .ForAll(batch =>
    {
        Task.Run(() =>
        {
            ...
        }
    }

GetDataItem 产生可枚举的数据项,这些数据项被划分为批处理数据集。 GetDataItem 将产生约 2,000,000 个项目,每个项目平均处理大约 0.3 毫秒。

该项目在 x64 平台上的最新 .NET 4.5 RC 上运行。

更新。

我当前的代码(如上所示)将获取项目并分批对其进行分区。每个批次在八个线程上并行处理(i7 上的 Environment.ProcessorCount)。处理速度很慢,CPU 密集型和内存密集型。 当单个批次的处理完成时,将启动一个任务以将结果数据异步导入 RavenDb。批量导入作业本身是同步的,如下所示:

using (var session = Store.OpenSession())
{
    foreach (var data in batch)
    {
        session.Store(data);
    }
    session.SaveChanges();
}

这种方法存在一些问题:

  1. 每完成一个批次,就会启动一个任务来运行导入作业。我想限制并行运行的任务数量(例如,最多 10 个)。此外,即使启动了许多任务,它们似乎也永远不会并行运行。

  2. 内存分配是个大问题。处理/导入批次后,它似乎仍保留在内存中。

我正在寻找解决上述问题的方法。理想情况下我想要:

  • 每个逻辑处理器一个线程执行繁重的批量数据处理。
  • 十个左右的并行线程等待完成的批次导入 RavenDb。
  • 将内存分配保持在最低限度,这意味着在导入任务完成后取消分配批次。
  • 不要在其中一个线程上运行导入作业以进行批处理。已完成批次的导入应与正在处理的下一批并行运行。

解决方案

var batchSize = 10000;
var bc = new BlockingCollection<List<Data>>();
var importTask = Task.Run(() =>
{
    bc.GetConsumingEnumerable()
        .AsParallel()
        .WithExecutionMode(ParallelExecutionMode.ForceParallelism)
        .WithMergeOptions(ParallelMergeOptions.NotBuffered)
        .ForAll(batch =>
        {
            using (var session = Store.OpenSession())
            {
                foreach (var i in batch) session.Store(i);
                session.SaveChanges();
            }
        });
});
var processTask = Task.Run(() =>
{
    datasupplier.GetDataItems()
        .Partition(batchSize)
        .AsParallel()
        .WithDegreeOfParallelism(Environment.ProcessorCount)
        .ForAll(batch =>
        {
            bc.Add(batch.Select(i => new Data()
            {
                ...
            }).ToList());
        });
});

processTask.Wait();
bc.CompleteAdding();
importTask.Wait();

【问题讨论】:

  • 您的代码不起作用吗?使用太多内存?是不是太慢了?
  • 您能否详细说明并具体输入您的问题?理想情况下,您当前的代码会发生什么与您期望它应该做什么。
  • 我对 RavenDb 的 API 不是很熟悉。它是否具有非阻塞 I/O 或仅同步的 BeginWrite/EndWrite 模式?
  • 看起来很像 map/reduce。您可能想研究该模式,或使用支持该模式的框架。我相信 RavenDB 在某种程度上支持 map/reduce。
  • 您可能想查看 WorkList 模式 here,只需添加您想要的块大小的工作项,然后分配一个 num_of_workers 以在进入 ConcurrentQueue 时对其进行处理。

标签: c# multithreading asynchronous parallel-processing task


【解决方案1】:

您的任务总体上听起来像是一个生产者-消费者工作流程。您的批处理器是生产者,而您的 RavenDB 数据“导入”是生产者输出的消费者。

考虑使用BlockingCollection&lt;T&gt; 作为批处理处理器和数据库导入器之间的连接。一旦批处理器将完成的批处理推送到阻塞集合中,数据库导入器将立即唤醒,并在它们“赶上”并清空集合时重新进入睡眠状态。

批处理器生产者可以全速运行,并且始终与处理先前完成的批处理的数据库导入器任务并行运行。如果您担心批处理器可能比数据库导入器领先太多(b/c 数据库导入比处理每个批处理花费的时间要长得多),您可以在阻塞集合上设置一个上限,以便生产者在添加时会阻塞超过这个限制,让消费者有机会赶上。

不过,您的一些 cmets 令人担忧。启动一个 Task 实例以异步执行 db 导入到批处理并没有什么特别的错误。任务!=线程。创建新任务实例不会产生与创建新线程相同的巨大开销。

不要过于精确地控制线程。即使您指定您想要的存储桶数量与您拥有的核心数量完全相同,您也无法独占使用这些核心。来自其他进程的数百个其他线程仍将安排在您的时间片之间。使用任务指定逻辑工作单元并让 TPL 管理线程池。让自己免于因错误的控制感而感到沮丧。 ;>

在您的 cmets 中,您指出您的任务似乎没有彼此异步运行(您如何确定这一点?)并且在每批完成后似乎没有释放内存。我建议放弃一切,直到您首先弄清楚这两个问题的原因。您是否忘记在某处调用 Dispose() ?你是否持有一个让整个对象树不必要地存活的引用?你在测量正确的东西吗?并行任务是否被阻塞的数据库或网络 I/O 序列化?在解决这两个问题之前,您的并行计划是什么并不重要。

【讨论】:

  • 感谢这些建议。我将研究 BlockingCollection。我想避免的任务问题是无法控制正在使用的线程数。我可以轻松地优化活动的并发 io-bound 任务数量,而调度程序几乎不可能做到这一点。在逻辑内核数量上并行化批处理的原因是批处理作业不是多线程的。即使 TPL 将管理线程,您也需要至少与逻辑内核相同数量的线程才能在最佳情况下利用整个 CPU。
  • 由于调试输出表明一个任务总是在另一个(第一行命中)开始之前完成(最后一行命中),任务似乎没有彼此异步运行。的确,它可以是任意数量的东西,不一定要与任务代码相关。
【解决方案2】:

对于每个批次,您都在开始一项任务。这意味着您的循环会很快完成。它留下了(批次数)不是您想要的任务。您想要(CPU 数量)。

解决方案:不要为每个批次开始一个新任务。 for 循环已经是并行的了。

针对您的评论,这里有一个改进的版本:

//this runs in parallel
var processedBatches = datasupplier.GetDataItems()
    .Partition(batchSize)
    .AsParallel()
    .WithDegreeOfParallelism(Environment.ProcessorCount)
    .Select(x => ProcessCpuBound(x));

foreach (var batch in processedBatches) {
 PerformIOIntensiveWorkSingleThreadedly(batch); //this runs sequentially
}

【讨论】:

  • 我不想阻塞批处理线程与导入所述批处理的 io 绑定任务。导入任务必须与正在处理的下一批并行运行。因此,我的有缺陷的代码启动了一项新任务,将每个批次导入 RavenDb。这允许批处理不间断地继续进行,但会引发另一组问题。
  • 我添加了一个可能的解决方案。
  • 我希望能够同时运行多个 PerformIOIntensiveWorkSingleThreadedly(batch)。另一个问题是我想避免存储对所有批次的引用(消除 foreach),因为数据集很大。根本没有足够的内存来保存它。
  • a) 您可以使 foreach 循环成为具有(可能不同)并行度的 Parallel.ForEach。 b) PLINQ 查询处于流模式,这意味着它不会占用所有内存。它将在请求时生成项目(带有一个小的有界缓冲区)。您可以尝试添加合并选项来控制缓冲 (msdn.microsoft.com/en-us/library/dd997424.aspx)。
【解决方案3】:

我最近构建了类似的东西,我使用了 Queue 类 vs List 和 Parallel.Foreach。我发现过多的线程实际上会减慢速度,这是一个甜蜜点。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-04-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-06-11
    • 2021-08-30
    相关资源
    最近更新 更多