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 的想法. 但是问题是线程会一开始一直添加工作队列, 带来储存的负载, 并且在开始计算前就带来了 的内存负担. 对于单线程, 与直接串行计算的消耗相比也很大
对于第二种: 线程把 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)表明blockA 的任务被其他线程偷取, 该线程更新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