[ blog / entry ]

CS149 Lecture 5:Work Distribution and Scheduling

学习笔记

[ contents ] 4 sections

Programming for high performance

优化并行程序的表现是一个循环的过程, 需要不断的重复: 分解, 分配, 组织等工作 重要的是做到下面的目标:

  • 在可用的计算资源中平衡工作负载
  • 降低通信成本
  • 降低额外的成本

TIP #1 总是先做最简单的优化, 再评估是否需要优化

Balancing workload 让处理器尽可能一直处于工作状态, 一般来说, 我们可以观察到不同的处理器处理的任务复杂度是不同的, 这导致整个任务的完成时间会受到分配到工作量最大的处理器拖长

等价地, 换一个观点, 如果有一段时间只有一个处理器在工作, 实际上这部分工作可以等价视作计算中的串行部分, 增加计算中的串行部分显然会限制我们的加速比

优化的方式是选择合适的任务分配来平衡负载

Static assignment

一种方式是静态分配工作, 即当计算任务给定时, 分配的方式是确定的 我们在 ass1 计算 Mandelbrot Set 的时候就采用了这种分配方式: 将图像的各区域分配给不同的 thread 这种方式的优点是: 简单, 并且几乎没有分配带来的多余负载 当工作的消耗以及任务量可以预计时, 这种分配方式是有效的, 我们可以估计每个处理器分配到的任务量

”Semi - static” assignment

这种半静态的分配方式是指我们能预计未来任务量, 即使用过去的任务分配作为分配策略的估计 我们在一段计算之中重新分配工作

一些例子: 粒子动力学模拟, 我们需要在粒子密度更大的地方进行更多的采样, 并且粒子在一段时间内局域, 于是在一段时间内这样的采样方式是有效的 自适应网格: 在我们更关心的, 更需要精度的地方加细网格

Dynamic assignment

程序可以动态地分配任务, 来保证负载能够平衡

下面是一个动态分配的例子

int N = 1024;

// assume all allocations are only executed by 1 thread
int* x = new int[N];
bool* is_prime = new bool[N];

// assume elements of x are initialized here

LOCK counter_lock;
int counter = 0;    // shared variable

while (1) {
    int i;

    lock(counter_lock);
    i = counter++;
    unlock(counter_lock);

    if (i >= N)
        break;

    is_prime[i] = test_primality(x[i]);
}

让每个 thread 去抢 lock , 抢到 lock 的 thread 就可以执行对应的计算

更一般的抽象, 就是把这种分配方式看作是一个任务队列, 让不同的处理器去接任务

很多待处理任务
┌────┐ ┌──┐ ┌──────┐ ┌───┐
│ A  │ │B │ │  C   │ │ D │
└────┘ └──┘ └──────┘ └───┘

        shared work queue
        ┌───────────┐
        │ task A    │
        │ task B    │
        │ task C    │
        │ task D    │
        │ ...       │
        └───────────┘
         ↓  ↓  ↓  ↓
        T1 T2 T3 T4

但一个任务的组成是什么? 从前面由 lock 实现的动态方式我们看到, thread 会处于抢 lock 和完成计算这两个状态, 这种分配方式如果有效, 就要求计算消耗的时间远大于抢 lock 的时间

事实上, 我们看到, 抢 lock 这个部分也是程序的串行部分

可以看到, 实际上我们面临着 lock 分配时间和完成计算的 trade off, 一种方式当然是增加完成计算的时间, 也就是一次分配给 thread 一批任务

下一个问题自然是: 如何选择任务的大小? 如果任务量分配过大, 那就又和 static assignment 没有区别了!

  • 任务量应该比处理器更多
  • 但任务划分也不能太细, 以减少分配任务的负载

一般来说, 一个合适的划分依赖于很多因素, 也和硬件有关系

我们来看另一种更聪明的任务分配方式: 如果将一系列任务分配给不同的处理器, 如果有处理器最后拿到了一个 workload 很重的任务, 这就会增加计算时间

当然我们会想到: 把这个重任务继续分解再分配就可以, 但这也许会增加同步的负载(任务内部有依赖, 分配给各处理器有通信成本), 如果是纯串行任务, 这种分解也不一定能做到!

另一种做法是: 把负载重的任务先分配出去, 多余的任务再匀给别的处理器, 尽可能让他们同时完成计算 这种策略需要我们对每个任务的大概负载有预计

对于降低同步负载也有一种应对, 我们提到, 对于单任务队列, 给处理器分配任务是纯串行的, 所以一种方式是把任务分配也并行化: 让每个处理器都有自己的任务队列, 但是不同处理器可以从别的任务队列中偷任务

        多个 local queues
      /      |      |      \
    Q1      Q2     Q3      Q4
    ↓       ↓      ↓       ↓
    T1      T2     T3      T4
     \
      \____ steal ____► 空闲 worker

并且这种任务队列中的工作并不要求工作之间独立, 只要可以分解就可以(当然会带来一些同步和通信成本)

Challenge: 实现一个好的负载平衡

  • 让所有处理器都在工作
  • 需要降低这种 solution 带来的成本(同步,通信以及需要的多余计算)

Static v.s. Dynamic

  • continuum spectrum
  • 尽可能利用知道的所有信息
               workload 越不可预测


 static     semi-static      dynamic queue      work stealing
   │             │                │                  │
   ▼             ▼                ▼                  ▼

overhead 低 --------------------------------> overhead 高

适应变化弱 ---------------------------------> 适应变化强

load balance
依赖预测   ---------------------------------> runtime 自动平衡

fork - join 并行化

从数据并行化的角度来看: 一个并行程序就是对许多数据做相同的操作 在许多并行方案上都有例子

// ISPC foreach
foreach (i = 0 ... N) {
    B[i] = foo(A[i]);
}

// ISPC bulk task launch
launch[numTasks] myFooTask(A, B);

// using higher-order function 'map'
map(foo, A, B);

// OpenMP parallel for
#pragma omp parallel for
for (int i = 0; i < N; i++) {
    B[i] = foo(A[i]);
}

// bulk CUDA thread launch
foo<<<numBlocks, threadsPerBlock>>>(A, B);

常见的并行化编程方案即为用 threads 显式管理并行

float* A;
float* B;

// initialize arrays A and B here

void myFunction(float* A, float* B) {
    ...
}

std::thread thread[NUM_HW_EXEC_CONTEXTS];

for (int i = 0; i < NUM_HW_EXEC_CONTEXTS; i++) {
    thread[i] = std::thread(myFunction, A, B);
}

for (int i = 0; i < NUM_HW_EXEC_CONTEXTS; i++) {
    thread[i].join();
}

一个典型的并行计算任务为 分治 divide - and - conquer 算法 例如: 快速排序, FFT 算法等

将一个任务分解为更小的独立任务, 然后继续分解, 可以让计算复杂度增加 log 因子

Fork - join 方案就是采用这种思想: 一个分解出来的任务可以不用立刻执行, 而是把它 fork 出来, 让别的处理器执行

我们以 Cpp 拓展 Cilk Plus 为例子, 它提供了两个基本的操作

cilk_spawn foo(args); 即 fork, 创建一个新的逻辑控制线程 Semantics: 调用 foo , 但调用者可以不用等待 foo 被执行就继续异步继续执行代码

cilk_sync; 即 join Semantics: 当通过 spawn 生成的所有函数都执行之后再返回

对于每个包含 cilk_spawn 的函数, 内部都隐式存在一个 sync 保证和这个函数关联的计算都结束

顾名思义, 每创建一个 cilk_spawn 调用, 实际上只是创建了一个需要执行的函数, 但不一定由当前执行程序的处理器执行

Abstraction v.s. implementation cilk 并没有指定调用的函数如何, 什么时候分配并执行, 这些调用可以被调用者, 以及调用者 spawn 出来的其他函数一起并行 事实上, cilk 的设计保证我们在代码中去掉这个 modifier , 整个代码仍是编译正确的 但 sync 确实会对调度施加约束

我们用 cilk 来优化 quick sort 的代码:

void quick_sort(int* begin, int* end) {

    if (begin >= end - PARALLEL_CUTOFF)
        std::sort(begin, end);

    else {
        int* middle = partition(begin, end);

        cilk_spawn quick_sort(begin, middle);
        quick_sort(middle + 1, end);
    }
}

注意, sync 没有出现是因为函数体内部在返回前隐式设置了一个 sync

流程抽象为

              quick_sort

             partition

          ┌───────┴───────┐
          │               │
 quick_sort(left)    quick_sort(right)
      spawn            当前线程继续
          │               │
          └───────┬───────┘

            implicit sync

                return

写一个 fork - join 的代码的方式就是把可独立计算的部分显式用 cilk 暴露出来即可

同样, 我们需要暴露足够多的独立工作来实现达到并行能力, 至少要比同时执行的硬件数要多 我们用 parallel slack 来评价这个指标, 即 独立的工作数 / 机器的并行能力 但不能暴露太多独立工作, 不然会隐式提高额外的成本(我们已经讨论过这个 trade off)

fork - join 程序的调度

我们只考虑两个简单的调度器:

  • pthread_create 来为每个 cilk_spawn 启动线程
  • cilk_sync 翻译为 pthread_join

潜在的性能问题

  • 启动线程额外的负载
  • 比核心数更多的线程数带来的: context 切换负担, 降低 cache locality 等

即我们现在来考虑 Cilk 的 implementation

Cilk Plus 的实现使用了叫 pool of worker threads 的概念, 实际上 Cilk 真正做的是维护一个线程池, 让这些线程接取任务, 而不是每次都新建线程

下一个问题是: 哪个线程去接任务, 当我们调用 cilk_spawn 的时候发生了什么

我们考虑一个简单的串行任务

foo();
bar();

对于正在执行的线程, 函数的调用顺序 continuation 隐式记录在了线程的调用栈 stack

于是我们可以用一个工作队列去储存: 每个线程做什么, 以及然后做什么 当执行到由 cilk_spawn modify 的函数时, 线程将它对应的 continuation 记录在工作队列中, 例如

cilk_spawn foo();
bar();

线程就执行 foo() , 然后将它的 continuation bar() 记录在工作队列中, 当别的线程空闲时, 它就会看到: 这个线程的工作队列中有一个 bar() 等待执行, 然后把它移动到自己的工作队列, 然后执行对应任务

这里有一个问题, 既然两个函数是独立的, 主线程应该执行哪个函数? 实际上两种方案都是可行的, 我们下面看不同实现的例子

考虑代码

for (int i=0; i<N; i++){

	cilk_spawn foo();
}
cilk_sync;

对于第一种: 先把 foo() 加入工作队列, 然后继续执行, 就会在自己的工作队列中不断添加待执行的 foo() . 这是一种 breadth - first 的想法. 但是问题是线程会一开始一直添加工作队列, 带来储存的负载, 并且在开始计算前就带来了 O(N)O(N) 的内存负担. 对于单线程, 与直接串行计算的消耗相比也很大

对于第二种: 线程把 continuation 的任务加入队列, 然后计算 foo(), 如果没有别的线程, 那么执行的操作和直接串行计算差别不大, 如果有别的线程, 就会拿走 continuation 的任务, 这是一种 depth - first 的想法

从这个例子中, 我们也能看到其他线程偷取, 这实际上就是队列结构的维护方式: 先进入队列的任务更可能是更重的任务, 被其他线程先偷取, 越靠近线程的任务是更小的任务, 同时在递归等类型的任务中也更好保持 locality 下一个问题是空闲的线程从哪个线程处偷取, 这里采用随机选择一个线程进行偷取

每个 Worker 都有自己的 deque

                   HEAD / TOP

               别人从这里偷

              ┌─────────────┐
              │   大 / 旧    │
              │             │
              │   小 / 新    │
              └─────────────┘

               自己从这里取

                 TAIL / BOTTOM

                    Worker


Worker 有活:
    → 一直做自己的
    → push/pop bottom

Worker 没活:
    → 随机找一个 victim
    → 从 victim 的 top 偷
    → 通常一次偷到较大的子任务

可以看到, 对于任务的分发等, 可能形成复杂的逻辑拓扑, 我们的 sync 也应该特殊设计

如果没有其他线程偷取任务, 那 sync 就没有实际效果(我们知道 sync 就是为了处理异步执行的数据同步问题)

对于存在其他线程偷取任务时, 需要维护一些数据, 来描述分发给各线程的哪些任务需要被哪个 sync 一同管理, 于是需要一个新的数据 Descriptor

Descriptor 记录了由某个任务块分发出的任务数量 spawn 以及被完成的任务数量 done , 每个线程完成对应的任务后增加 done 的计数

如何完成 done 的计数? 当某个线程上送入工作队列的任务被别的线程偷取后, 对应任务所在位置并不会清除, 而是更新为状态 STOLEN(id=A) 表明 block A 的任务被其他线程偷取, 该线程更新 spawn, 对应线程完成后更新 done 最后直到两者相等时, 释放 Descriptor , sync 完成, 进行下一步指令

当然也许会问中间的时候恰好两个数值相等怎么办, 但是注意到这种执行逻辑, 程序的 continuation 是送给了别的线程跑, 当某个线程真的拿到了 continuation 的 sync, 那么 spawn 其实已经达到最大值了, 此时一定有 done <= spawn , 这个逻辑流不会导致中间突然 done = spawn 时程序错误打破 barrier

Greedy join schedule

  • 所有线程在空闲时都尝试偷取其他线程的任务
  • 当系统中没有其他任务时, 线程才会空闲
  • 执行 spawn 的线程并不一定是执行 sync 后的逻辑的线程

并且当没有 stealing 时, 体系并没有额外的维护成本, 偷取任务通常也只在偷取较大任务时发生, 降低了通信成本, 大多数时候, 线程都在处理自己任务队列的任务

Fork - join parallelism: 一种自然并行化分治算法的方式 (除了 Cilk Plus, 也包括其他类似语法, 如 OpenMP)

Cilk Plus 对 spawn/sync 的实现, 以及 locality - aware work stealing scheduler 的抽象

  • 总是执行生成的任务(continuation stealing)
  • Greedy policy
← Back to blog