Skip to content

CPU、SIMD 与任务调度 ​

现代 CPU 同时使用指令级并行、向量执行、多核和硬件多线程。它们增加吞吐量的方式不同:多个核执行不同线程;SIMD 用一条指令处理多条数据 lane;硬件线程共享部分核心资源,帮助隐藏等待。不能把“4 核、每核 2 线程、8 lane”简单相乘后当成所有程序的加速保证。

SIMD 与掩码 ​

考虑 8 个独立元素做相同计算,8-lane SIMD 有机会用一组向量指令完成。若各元素循环次数为 [1,1,1,1,1,1,1,8],在朴素掩码执行模型中,最后一个 lane 仍需 8 轮,其余 lane 在结束后被屏蔽。总有效轮次为 15,最大槽位为 64,利用率只有约 23.4%。

把相似工作量的数据放在一起可以提高控制流一致性,但重排需要额外成本,还可能影响原始顺序。宽向量的潜力依赖可并行数据、规则访存和控制流一致性,不是宽度增加就按比例变快。

尾部也需要正确掩码。N=10、宽度=8 时,第二组只有两个有效元素;未屏蔽的载入会越界。用能被向量宽度整除的 N 测试,容易遗漏这个错误。官方 Assignment 1 的向量练习适合观察这些机制。

ISPC 用 SPMD 语义表达多个程序实例,uniform 表示一组实例共享的值,varying 表示每实例不同的值;编译器将适当执行映射到向量硬件。它的 task 抽象还可进一步使用多个核,单核 SIMD 与多核任务要分别分析。ISPC 用户手册

连续区间与动态领取 ​

本站 benchmark 的静态策略给线程 id 分配区间:

cpp
begin = N * id / P;
end   = N * (id + 1) / P;

整数除法自然处理 N 不能整除 P 的尾部,也允许 P>N 时部分线程得到空区间。各线程写不同元素,因此无需为结果数组加一把总锁;主线程 join 后再检查完整数组。

连续区间具有低调度开销和良好局部性,但面对聚集的重任务会失衡。动态策略维护一个原子 next,每次领取 grain 个项:

cpp
begin = next.fetch_add(grain, std::memory_order_relaxed);
if (begin >= N) break;
end = std::min(begin + grain, N);

原子读改写确保不同线程拿到不重叠区间。这里 relaxed 只承担任务编号唯一分配,不用于发布其他线程刚写入的数据;所有结果最终通过 join 建立可见性。内存序分析见下一章。

cpp
// Original scheduling benchmark: uneven independent jobs, exact comparison.
#include <algorithm>
#include <atomic>
#include <chrono>
#include <cstdint>
#include <cstdlib>
#include <iomanip>
#include <iostream>
#include <stdexcept>
#include <string>
#include <thread>
#include <vector>

using Clock = std::chrono::steady_clock;
using Values = std::vector<std::uint64_t>;

static std::size_t argument(const char *text, std::size_t maximum) {
    std::string value(text);
    if (value.empty() || value.find_first_not_of("0123456789") != std::string::npos)
        throw std::invalid_argument("arguments must be positive decimal integers");
    std::size_t used = 0;
    auto n = std::stoull(value, &used);
    if (used != value.size() || n == 0 || n > maximum)
        throw std::invalid_argument("argument outside supported range");
    return static_cast<std::size_t>(n);
}

static std::uint64_t work(std::size_t i, std::size_t n) {
    // The first quarter is deliberately heavier: static contiguous chunks
    // give one worker much more work even though item counts are equal.
    unsigned steps = i < n / 4 ? 512 : 16;
    std::uint64_t x = static_cast<std::uint64_t>(i) + UINT64_C(0x9e3779b97f4a7c15);
    for (unsigned k = 0; k < steps; ++k) {
        x ^= x >> 12;
        x ^= x << 25;
        x ^= x >> 27;
        x *= UINT64_C(2685821657736338717);
    }
    return x;
}

static void range(Values &output, std::size_t begin, std::size_t end) {
    for (std::size_t i = begin; i < end; ++i) output[i] = work(i, output.size());
}

static void execute(Values &output, const std::string &policy,
                    std::size_t workers, std::size_t grain) {
    if (policy == "serial") { range(output, 0, output.size()); return; }
    std::atomic<std::size_t> next{0};
    std::vector<std::thread> threads;
    threads.reserve(workers);
    try {
        for (std::size_t id = 0; id < workers; ++id) {
            threads.emplace_back([&, id] {
                if (policy == "static") {
                    range(output, output.size() * id / workers,
                          output.size() * (id + 1) / workers);
                } else {
                    for (;;) {
                        std::size_t begin = next.fetch_add(grain, std::memory_order_relaxed);
                        if (begin >= output.size()) break;
                        range(output, begin, std::min(begin + grain, output.size()));
                    }
                }
            });
        }
    } catch (...) {
        for (auto &thread : threads) thread.join();
        throw;
    }
    for (auto &thread : threads) thread.join();
}

static std::uint64_t checksum(const Values &values) {
    std::uint64_t hash = UINT64_C(1469598103934665603);
    for (std::uint64_t value : values) { hash ^= value; hash *= UINT64_C(1099511628211); }
    return hash;
}

static void measure(const Values &expected, const std::string &policy,
                    std::size_t workers, std::size_t grain, std::size_t repeats) {
    Values output(expected.size());
    execute(output, policy, workers, grain); // One untimed warmup.
    if (output != expected) throw std::runtime_error("warmup mismatch: " + policy);
    std::vector<double> times;
    for (std::size_t trial = 0; trial < repeats; ++trial) {
        std::fill(output.begin(), output.end(), 0); // Not timed.
        const auto begin = Clock::now();
        execute(output, policy, workers, grain); // Includes thread creation/join.
        const auto end = Clock::now();
        if (output != expected) throw std::runtime_error("result mismatch: " + policy);
        times.push_back(std::chrono::duration<double, std::milli>(end - begin).count());
    }
    std::sort(times.begin(), times.end());
    double median = times[times.size()/2];
    if (times.size() % 2 == 0) median = (times[times.size()/2-1] + median) / 2;
    std::cout << policy << ',' << workers << ',' << grain << ',' << expected.size() << ','
              << repeats << ',' << times.front() << ',' << median << ',' << times.back()
              << ',' << checksum(output) << '\n';
}

int main(int argc, char **argv) {
    try {
        if (argc > 5) throw std::invalid_argument("usage: benchmark [items [threads [grain [repeats]]]]");
        std::size_t count = argc > 1 ? argument(argv[1], 5000000) : 200000;
        std::size_t workers = argc > 2 ? argument(argv[2], 128) : 4;
        std::size_t grain = argc > 3 ? argument(argv[3], 5000000) : 64;
        std::size_t repeats = argc > 4 ? argument(argv[4], 31) : 5;
        Values expected(count);
        range(expected, 0, count);
        std::cerr << "compiler=" << __VERSION__ << " hardware_concurrency_hint="
                  << std::thread::hardware_concurrency() << '\n';
        std::cout << "policy,threads,grain,items,repeats,min_ms,median_ms,max_ms,checksum\n";
        std::cout << std::fixed << std::setprecision(6);
        measure(expected, "serial", 1, count, repeats);
        measure(expected, "static", workers, grain, repeats);
        measure(expected, "dynamic", workers, grain, repeats);
        std::cerr << "all output elements matched serial reference\n";
    } catch (const std::exception &error) {
        std::cerr << error.what() << '\n';
        return 1;
    }
}

粒度的折中 ​

grain 太小,线程频繁争用 next,调度成本可能超过有效计算;grain 太大,最后剩余大块无法再细分,部分线程先空闲。一个有效实验是固定 N 和 P,扫描 grain=1、16、64、256、4096,并观察 min/median/max。不要在只测一个输入后宣称某个 grain 普遍最优。

动态集中队列也不是唯一方案。循环分配 i=id, id+P, ... 无需共享领取,对有规律的负载可能有效,但访存更分散。工作窃取则让线程优先处理本地队列,空闲后从其他线程取任务,试图同时获得局部性与负载均衡;队列并发、任务依赖和结束检测会增加实现复杂度。

任务图与线程池 ​

任务描述“要做的一段工作”,线程是执行载体。为 10000 个小任务各创建一个线程,可能耗费大量启动与调度成本;线程池复用少量工作线程来运行许多任务。存在依赖时,任务维护未完成前驱计数,前驱完成后递减,变成零才进入 ready 队列。

flowchart LR
  A[任务 A] --> C[任务 C]
  B[任务 B] --> C
  C --> D[任务 D]
  Q[ready 队列] --> W1[worker 1]
  Q --> W2[worker 2]
查看流程图文本
flowchart LR
  A[任务 A] --> C[任务 C]
  B[任务 B] --> C
  C --> D[任务 D]
  Q[ready 队列] --> W1[worker 1]
  Q --> W2[worker 2]

注意任务全部“已提交”不代表全部“已完成”。运行时需要区分未提交、等待依赖、就绪、执行中、完成;线程退出前还要确保不会再有新 ready 任务。锁内更新状态,锁外执行用户函数,避免一个长任务阻塞整个调度器。官方 Assignment 2 可用于进一步练习完整任务系统。

练习与自测 ​

先运行 ./labs/parallel/benchmark 1003 4 17 3,确认所有结果逐元素一致。然后把重任务分布改为每四项一个重任务,比较同样静态策略;最后增大 N,看线程启动成本的相对占比怎样变化。当前 benchmark 每次创建线程,没有池化,报告时应明确这一点。

自测:多个线程写同一个 vector 的不同元素是否合法?

在容器大小不变、对象不重叠且没有其他冲突访问的前提下,普通元素可以独立写。不要一边 push_back 导致重分配一边写;vector<bool> 的位压缩也有特殊风险。本例使用预分配的 vector<uint64_t>。

自测:动态调度总能胜过静态调度吗?

不能。工作均匀时静态调度低开销且局部性好;任务很小时原子领取可能主导。动态调度的收益来自负载适应能力,需抵消它增加的协调成本。

下一章:同步与内存一致性。