【问题标题】:Producer-Consumer using OpenMP-Tasks使用 OpenMP 任务的生产者-消费者
【发布时间】:2012-10-23 16:05:42
【问题描述】:

我正在尝试使用 OpenMP 中的任务实现并行算法。 并行编程模式是基于生产者-消费者的思想,但是 由于消费者进程比生产者慢,我想使用一些 生产者和几个消费者。 主要思想是创建与生产者一样多的操作系统线程,然后每个 这些将创建要并行完成的任务(由消费者)。每一个 生产者将与一定数量的消费者相关联(即 numCheckers/numSeekers)。 我在英特尔双芯片服务器上运行算法,每个芯片有 6 个内核。 问题是当我只使用一个生产者(搜索者)并且数量越来越多时 消费者(跳棋)的性能下降得非常快,因为数量 消费者增长(见下表),即使正确的核心数量是 以 100% 的速度工作。 另一方面,如果我增加生产者的数量,平均时间 减少或至少保持稳定,即使有成比例的数量 消费者。 在我看来,所有的改进都是通过输入的划分来实现的 在生产者之间,任务只是窃听。但同样,我没有任何 解释一个生产者的行为。我是否遗漏了什么 OpenMP 任务逻辑?我是不是做错了什么?

-------------------------------------------------------------------------
|   producers   |   consumers   |   time        |
-------------------------------------------------------------------------
|       1       |       1       |   0.642935    |
|       1       |       2       |   3.004023    |
|       1       |       3       |   5.332524    |
|       1       |       4       |   7.222009    |
|       1       |       5       |   9.472093    |
|       1       |       6       |   10.372389   |
|       1       |       7       |   12.671839   |
|       1       |       8       |   14.631013   |
|       1       |       9       |   14.500603   |
|       1       |      10       |   18.034931   |
|       1       |      11       |   17.835978   |
-------------------------------------------------------------------------
|       2       |       2       |   0.357881    |
|       2       |       4       |   0.361383    |
|       2       |       6       |   0.362556    |
|       2       |       8       |   0.359722    |
|       2       |      10       |   0.358816    |
-------------------------------------------------------------------------

我的代码的主要部分是休闲:

int main( int argc, char** argv) {

  // ... process the input (read from file, etc...)

  const char *buffer_start[numSeekers];
  int buffer_len[numSeekers];

  //populate these arrays dividing the input
  //I need to do this because I need to overlap the buffers for
  //correctness, so I simple parallel-for it's not enough 

  //Here is where I create the producers
  int num = 0;
  #pragma omp parallel for num_threads(numSeekers) reduction(+:num)
  for (int i = 0; i < numSeekers; i++) {
      num += seek(buffer_start[i], buffer_len[i]);
  }

  return (int*)num;
}

int seek(const char* buffer, int n){

  int num = 0;

  //asign the same number of consumers for each producer 
  #pragma omp parallel num_threads(numCheckers/numSeekers) shared(num)
  {
    //only one time for every producer
    #pragma omp single
    {
      for(int pos = 0; pos < n; pos += STEP){
    if (condition(buffer[pos])){
      #pragma omp task shared(num)
      {
        //check() is a sequential function
        num += check(buffer[pos]);
      }
    }
      }
      #pragma omp taskwait
    }
  return num;
}

【问题讨论】:

  • 您启用了嵌套并行,不是吗?请注意,标准规定任务执行可能会延迟到达到调度点(例如taskwait)。任务中还有num 的数据竞争。您应该使用atomic 构造来保护累积。如果您的系统由两个 AMD64 芯片或 Nehalem 和后 Nehlaem Intel 芯片组成,请记住 NUMA 位置。

标签: c multithreading task openmp


【解决方案1】:

观察到的行为是由于您没有启用嵌套的parallel 区域。发生的情况是,在第一种情况下,您实际上正在经历 OpenMP 任务的巨大开销。这很可能是由于 check() 与 OpenMP 运行时引入的开销相比没有做足够的工作。为什么它对 1 个和 2 个生产者的行为如此?

仅使用一个生产者运行时,外部parallel 区域仅使用一个线程执行。根据 OpenMP API 规范,此类parallel 区域不活动,它们只是串行执行内部代码(唯一的开销是额外的函数调用和通过指针访问共享变量)。在这种情况下,内部parallel 区域虽然在嵌套并行性被禁用的情况下被嵌套,但会变为活动并激发大量任务。任务引入了相对较高的开销,并且这种开销随着线程数的增加而增加。对于 1 个消费者,内部 parallel 区域也不活动,因此可以连续运行而没有任务开销。

当与两个生产者一起运行时,外部parallel 区域活动,因此内部parallel 区域呈现不活动(请记住 - 没有启用嵌套并行) 因此,根本没有创建任何任务 - seek() 只是串行运行。没有任务开销,代码运行速度几乎是 1 个生产者/1 个消费者案例的两倍。运行时间不依赖于消费者的数量,因为无论指定多少线程,内部 parallel 区域始终不活动

任务分配和对共享变量的协同访问会带来多大的开销?我创建了一个执行以下代码的简单综合基准测试:

for (int i = 0; i < 10000000; i++) {
   ssum += sin(i*0.001);
}

在默认优化级别为 GCC 4.7.2 的 Westmere CPU 上连续执行不到一秒。然后我介绍了使用简单的atomic 构造的任务来保护对共享变量ssum 的访问:

#pragma omp parallel
{
   #pragma omp single
   for (int i = 0; i < 10000000; i++) {
      #pragma omp task
      {
         #pragma omp atomic
         ssum += sin(i*0.001);
      }
   }
}

(这里不需要taskwait,因为parallel区域末尾的隐式屏障有一个调度点)

我还创建了一个更复杂的变体,它执行归约的方式与 Massimiliano 提出的相同:

#define STRIDE 8

#pragma omp parallel
{
   #pragma omp single
   for (int i = 0; i < 10000000; i++) {
      #pragma omp task
      {
         const int idx = omp_get_thread_num();
         ssumt[idx*STRIDE] += sin(i*0.001);
      }
   }
   #pragma omp taskwait

   const int idx = omp_get_thread_num();
   #pragma omp atomic
   ssum += ssumt[idx*STRIDE];
}

代码是用 GCC 4.7.2 编译的,例如:

g++ -fopenmp -o test.exe test.cc

在双插槽 Westmere 系统(总共 12 个内核)上以批处理模式运行它(因此没有其他进程可以干预),具有不同的线程数和插槽上的不同线程位置,可以获得以下运行时间两个代码:

Configuration   ATOMIC       Reduction    ATOMIC slowdown
2 + 0            2,79 ±0,15   2,74 ±0,19   1,8%
1 + 1            2,72 ±0,21   2,51 ±0,22   8,4% <-----
6 + 0           10,14 ±0,01  10,12 ±0,01   0,2%
3 + 3           22,60 ±0,24  22,69 ±0,33  -0,4%
6 + 6           37,85 ±0,67  38,90 ±0,89  -2,7%

(运行时间以秒为单位给出,由omp_get_wtime() 测量,10 次运行/std 的平均值。还显示了偏差/;Configuration 列中的x + y 表示第一个套接字上的x 线程和@987654344 @第二个套接字上的线程)

如您所见,任务的开销是巨大的。它比使用 atomic 而不是对线程私有累加器应用归约的开销要高得多。此外,atomic+= 的赋值部分编译为锁定的比较和交换指令 (LOCK CMPXCHG) - 不会比每次调用 omp_get_thread_num() 高多少。

还应注意,双插槽 Westmere 系统是 NUMA,因为每个 CPU 都有自己的内存,并且对另一个 CPU 内存的访问通过 QPI 链路进行,因此延迟增加(并且可能带宽降低)。由于ssum 变量在atomic 情况下是共享的,所以在第二个处理器上运行的线程实际上是在发出远程请求。尽管如此,两种代码之间的差异仍然可以忽略不计(除了标记的双线程情况 - 我必须调查原因),当线程数量增加时,atomic 代码甚至开始优于减少的代码。

在多尺度 NUMA 系统上,atomic 方法中的同步可能会成为更大的负担,因为它会为已经很慢的远程访问增加锁定开销。一个这样的系统是我们的 BCS 耦合节点之一。 BCS (Bull Coherence Switch) 是 Bull 的专有解决方案,它使用 XQPI (eXternal QPI) 将多个 Nehalem-EX 板连接到一个系统中,引入了三个级别的 NUMA(本地内存;同一板上的远程内存) ; 远程板上的远程内存)。当在一个这样的系统上运行时,由 4 个带有 4 个八核 Nehalem-EX CPU(总共 128 个内核)的板组成,atomic 可执行文件运行 1036 秒(!!),而缩减方法运行 1047 秒,即两者仍然执行大约相同的时间(我之前声明atomic 方法慢 21.5% 是由于测量期间操作系统服务抖动)。这两个数字都来自单次运行,因此不太具有代表性。请注意,在此系统上,XQPI 链接为板间 QPI 消息引入了非常高的延迟,因此锁定非常昂贵,但不会那么昂贵。使用归约可以消除部分开销,但必须正确实施。首先,减少变量的本地副本也应该是线程执行的 NUMA 节点的本地副本,其次,应该找到一种不调用 omp_get_thread_num() 的方法。这两个可以通过多种不同的方式实现,但最简单的一种就是使用threadprivate 变量:

static double ssumt;
#pragma omp threadprivate(ssumt)

#pragma omp parallel
{
   ssumt = 0.0;

   #pragma omp single
   for (int i = 0; i < 10000000; i++) {
      #pragma omp task
      {
         ssumt += sin(i*0.001);
      }
   }
   #pragma omp taskwait

   #pragma omp atomic
   ssum += ssumt;
}

访问ssumt 不需要保护,因为两个任务很少在同一个线程中同时执行(必须进一步调查这是否符合 OpenMP 规范)。此版本的代码执行时间为 972 秒。再一次,这与 1036 秒相差不远,仅来自一次测量(即它可能只是一个统计波动),但理论上它应该更快。

带回家的教训:

  • 阅读有关嵌套parallel 区域的OpenMP 规范。通过将环境变量OMP_NESTED 设置为true 或调用omp_set_nested(1); 来启用嵌套并行。如果启用,活动嵌套的级别可以由 OMP_MAX_ACTIVE_LEVELS 控制,正如 Massimiliano 所指出的那样。
  • 注意数据争用并尝试使用最简单的方法来防止它们。并非每次使用更复杂的方法都能为您带来更好的性能。
  • 特殊系统通常需要特殊编程。
  • 如有疑问,请使用线程检查工具(如果有)。英特尔有一个(商业),Oracle 的 Solaris Studio(以前称为 Sun Studio)也有一个(免费;尽管产品名称有 Linux 版本)。
  • 注意开销!尝试将作业分成足够大的块,这样创建数百万个任务的开销不会抵消获得的并行增益。

【讨论】:

  • 感谢 Hristo 的回答。你是对的,我没有使用嵌套并行。一旦我激活了嵌套并行性,内部并行块就被激活了,现在多个生产者的时间与我之前只为一个生产者展示的时间相似。真正让我吃惊的是 OpenMP 中的任务是多么的低效。我已经使用 Task 类在 .net 中实现了几个算法,它似乎对性能没有太大影响,或者平台本身可能会用自己的开销隐藏部分开销。
  • 你建议我使用任务还是我应该尝试使用队列或类似的方式自己安排和拆分工作?我知道如果每个任务要完成的工作量太小,任务是一个坏主意,但在我的情况下,这是一种简单的方法,可以避免实现保存检查器使用的工作的数据结构并制作它并行安全,因为这非常接近任务的工作方式。
  • 一种可能的解决方案是使用parallel for schedule(dynamic,10) reduction(+:num) 而不是single + task(调整schedule 中的块大小,直到获得最佳性能)。另一种可能的解决方案是重新实现任务 :) 在single 块中构建与condition() 匹配的所有指针(或buffer 的任何元素)的列表。然后只需使用 for 使用静态计划(或动态,如果 check() 可能需要不同的时间来执行)运行它。
  • 我也可以添加这个 - 明智地使用并行结构。您的代码已经单线程运行了 0.642935 秒。并行构造会增加开销(或 很多 开销;取决于您如何使用它们)。除非您有 很多 个项目要处理,否则并行处理是没有意义的。您可能还想查看if 指令,该指令允许有选择地停用parallel 区域,例如当工作项的数量太少而并行处理无法生效时。
【解决方案2】:

正如 Hristo 在评论中所建议的,您应该启用嵌套并行。这是设置环境变量完成的:

  • OMP_NESTED(启用或禁用嵌套并行)
  • OMP_MAX_ACTIVE_LEVELS(控制嵌套活动并行区域的最大数量)

另一方面,与其使用atomic 构造来保护积累,我建议采用以下策略:

...
// Create a local buffer to accumulate partial results
const int nthreads = numCheckers/numSeekers;
const int stride   = 32; // Choose a value that avoids false sharing
int numt[stride*nthreads];
// Initialize to zero as we are reducing on + operator
for (int ii = 0; ii < stride*nthreads; ii++)    
  numt[ii] = 0;

#pragma omp parallel num_threads(numCheckers/numSeekers)
{

  //only one time for every producer
  #pragma omp single
  {
    for(int pos = 0; pos < n; pos += STEP){
      if (condition(buffer[pos])){
      #pragma omp task
      {
        //check() is a sequential function
        const int idx = omp_get_thread_num();
        numt[idx*stride] += check(buffer[pos]);
      }
    }
  }
  #pragma omp taskwait

  // Accumulate partial results
  const int idx = omp_get_thread_num();
  #pragma atomic
  num += numt[stride*idx];
}

这应该可以防止由于同时请求写入同一内​​存位置而导致的潜在减速。

注意以前版本的答案,建议在最里面的并行区域使用reduction是错误的:

出现在最里面的归约子句中的列表项 封闭的工作共享或并行结构可能无法在 明确的任务

OpenMP 3.1 规范的 §2.9.3.6 不允许。

【讨论】:

  • 不幸的是,您的示例不符合要求。 OpenMP API 规范对可以从何处访问reduction 变量施加了限制,并且显式任务(由task 构造创建)不在允许的位置(OpenMP API 规范 v3.1,§2.9.3.6)。这就是为什么我推荐使用atomic。 OpenMP API 规范 v4.0 中将减少任务。
  • @HristoIliev 该死的,你是对的! :-) 非常感谢您指出 §2.9.3.6,因为我不知道。我会删除答案,因为它明显有缺陷。
  • @HristoIliev 现在应该没问题了。如果您仍然发现任何问题,请告诉我。
  • 这与过早的优化接壤。我做了一个简单的测试,运行了 1000 万个任务:sum += sin(i*0.01); 都使用atomic 构造和你的减少方法。在双插槽 Westmere 系统上运行 2、6 和 12 线程时,两者之间没有明显差异。 atomic 仅在两个线程在不同套接字上运行的情况下(但不是在 3+3 线程情况下)慢 8%,而在 12 个线程上运行时你的速度较慢。不过,在 128 核 3 级 NUMA 系统上,您的方法要快 20%(但它需要特别注意 OpenMP)。
猜你喜欢
  • 2012-01-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-09-25
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多