Skip to content

并发与多线程

下文的队列与线程池示例使用 C++17(std::optional、类模板实参推导)。

并发正确性首先依赖于明确的所有权、生命周期和同步关系,再考虑性能。若多个线程并发访问同一内存位置,且至少一个操作是写入,必须使用锁、原子操作或其他同步机制;否则就是数据竞争,在 C++ 中行为未定义。

mutex、自旋锁与 atomic

工具特征适用场景
std::mutex未获取到锁时线程可被阻塞/睡眠,不浪费 CPU 忙等临界区可能较长、竞争不可预测、涉及 I/O 或复杂不变量
spinlock取锁失败时循环尝试,持续占用 CPU临界区极短、等待时间远小于线程调度开销且不能睡眠的场景;用户态应慎用
std::atomic<T>对单个原子对象的读改写不可分割,可指定内存序标志位、计数器、简单状态机;复杂多变量不变量通常仍需锁或串行化

std::mutex

cpp
std::mutex mutex;
int counter = 0;

void increment() {
    std::lock_guard<std::mutex> guard(mutex); // RAII:离开作用域自动解锁
    ++counter;
}

优先用 std::lock_guard;需要手动解锁、延迟加锁或配合条件变量时使用 std::unique_lock。不要在持锁期间执行长时间计算、网络 I/O 或调用未知外部代码,否则竞争会放大,甚至引入死锁。

自旋锁

下面是教学用的简单实现:

cpp
#include <atomic>

class SpinLock {
public:
    void lock() noexcept {
        while (flag_.test_and_set(std::memory_order_acquire)) {
            // 忙等;生产环境可加入 pause/yield/backoff 降低竞争开销
        }
    }

    void unlock() noexcept {
        flag_.clear(std::memory_order_release);
    }

private:
    std::atomic_flag flag_ = ATOMIC_FLAG_INIT;
};

自旋锁等待时不会让出 CPU。若持锁线程被抢占,其他核心可能空转很久;在单核、临界区长或高竞争场景下尤其糟糕。普通业务代码通常优先使用 std::mutex

原子操作不等于“整个业务原子”

cpp
std::atomic<bool> stopped{false};
std::atomic<int> request_count{0};

stopped.store(true, std::memory_order_release);
request_count.fetch_add(1, std::memory_order_relaxed);

atomic 能安全处理一个标志或计数器,但不能自动维护多个字段之间的整体关系。例如“库存大于 0 才扣减库存并创建订单”涉及多个状态与副作用,通常需要 mutex、数据库事务、CAS 循环配合清晰状态机,或将操作串行化。内存序也不应随意优化:不了解同步关系时,先用默认的顺序一致性(seq_cst)或在有明确发布-获取关系时使用 release/acquire

std::memory_order 如何使用

原子操作保证对这个原子对象的读改写不会产生数据竞争;内存序进一步规定它与其他内存读写之间的可见性和重排序边界。它解决的问题是:一个线程写好普通数据后,如何让另一个线程在看到“已就绪”标志时,也能可靠看到这份数据。

内存序含义典型用途
relaxed仅保证该原子对象的原子性和单对象修改顺序;不建立其他数据的同步关系统计计数器、无需关联其他状态的指标
acquire阻止后续读写被重排到该操作之前;若读到了匹配 release 发布的值,可看到发布前的写入消费“已就绪”标志、取走已发布节点
release阻止此前读写被重排到该操作之后;将此前写入发布给 acquire 一方发布对象、设置就绪标志、入队完成
acq_rel同时具有 acquire 与 release 语义,只能用于读改写(RMW)操作CAS 既取得旧状态又发布新状态
seq_cstacquire/release 之外,再为所有顺序一致原子操作提供单一全局顺序默认选择;正确性优先、尚未证明可放宽时

relaxed:只用于独立原子状态

cpp
#include <atomic>
#include <cstdint>

std::atomic<std::uint64_t> request_count{0};

void on_request() {
    request_count.fetch_add(1, std::memory_order_relaxed);
}

这里我们只关心计数值本身的原子递增,并不依靠“计数变成某个值”去读取其他普通变量,因此 relaxed 合适。它不会让其他非原子数据自动可见,也不能用来发布对象:

cpp
struct Message { int payload = 0; };
Message message;
std::atomic<bool> ready{false};

// 错误示意:消费者看到 ready 后,仍没有语言层面的保证能看到 message 的完整写入。
message.payload = 42;
ready.store(true, std::memory_order_relaxed);

release + acquire:发布—获取

cpp
#include <atomic>

struct Message {
    int payload = 0; // 非原子数据
};

Message message;
std::atomic<bool> ready{false};

// 生产者线程
message.payload = 42;
ready.store(true, std::memory_order_release); // 发布此前对 message 的写入

// 消费者线程
while (!ready.load(std::memory_order_acquire)) {
    // 等待;实际程序通常应使用条件变量、信号量或退避,避免空转
}
int value = message.payload; // 保证读到 42

当 acquire load 读到了该 release store(或其 release sequence)发布的值时,两者建立 synchronizes-with 关系,进而使生产者对 message 的写入 happens-before 消费者之后的读取。前提是发布后不再与消费者并发修改同一份 message;若有重复发布、复用槽位或多方读写,仍需完整的状态协议。

这正是前文 SPSC 环形队列在生产者写入槽位后以 release 更新 tail、消费者以 acquire 读取 tail 的原因;消费者释放槽位时同理用 release 更新 head

CAS 与 acq_rel

CAS(compare_exchange_weak/strong)和 fetch_add 等读改写操作可使用 acq_rel:成功时既获取之前发布的状态,又发布当前更新。

cpp
std::atomic<int> state{0};
int expected = 0;

if (state.compare_exchange_strong(expected, 1,
                                  std::memory_order_acq_rel,
                                  std::memory_order_acquire)) {
    // 成功:获得旧状态并发布新状态
} else {
    // 失败:expected 被更新为实际值;失败内存序不能是 release 或 acq_rel
}

不要机械地给每个 CAS 都写 acq_rel:若操作仅用于统计,relaxed 即可;若只发布新数据,可用 release;若只消费已发布数据,可用 acquire。选择应由状态机的发布/读取关系决定。

实践准则与常见误区

  1. 默认使用 seq_cst,只有在测量到性能瓶颈且能证明同步关系时再放宽;
  2. store 不能使用 acquireload 不能使用 release;RMW 操作可使用上述全部常用内存序;
  3. volatile 不能代替 atomic,它不提供线程同步和原子性;
  4. 不要用内存屏障或复杂内存序“碰运气”;std::atomic_thread_fence 属于更高级工具,应先明确需要建立的 happens-before 关系;
  5. memory_order_consume 在主流编译器中的支持和实现不可靠,实践中通常使用 acquire 代替。

volatile 不等于 atomic

volatile 的作用是让对该对象的读写成为编译器不可随意省略或合并的可观察访问;它保证读改写原子性、不建立线程间 happens-before、也不替代 mutex 或 std::atomic

需求应使用的工具
普通多线程共享状态、计数器、发布—获取std::atomic 或 mutex
内存映射硬件寄存器(MMIO)通常由平台/驱动 API 指定的 volatile 访问与硬件屏障
信号处理遵从平台和标准的严格限制;不要把 volatile 当作通用线程同步
防止编译器消除普通业务代码应改正测试/程序设计,而不是滥用 volatile

“其他 CPU 或其他进程会修改共享内存”不是使用 volatile 的理由;仍须按并发协议使用原子操作、锁和必要的跨进程同步原语。volatile 适用于特殊硬件/底层接口语义,现代 C++ 应用层并发默认选择 atomic

条件变量

条件变量让线程在条件不满足时原子地释放 mutex 并休眠;被唤醒后重新持锁再检查条件。

cpp
std::mutex mutex;
std::condition_variable cv;
std::queue<Message> queue;
bool stopped = false;

std::unique_lock<std::mutex> lock(mutex);
cv.wait(lock, [&] {
    return stopped || !queue.empty();
});

必须使用带谓词的 wait,或手写 while 循环:条件变量可能发生虚假唤醒,而且即使被正常唤醒,其他消费者也可能先一步取走数据。notify_one() 通常唤醒一个等待者,notify_all() 唤醒所有等待者;修改受保护条件后再通知是常见写法。

有界生产者—消费者队列

网络服务常把不同职责拆开:

text
网络线程 → 有界请求队列 → 业务线程 → 有界响应队列 → 网络线程

有界队列限制缓冲数量,而不是无限积压请求。它的价值包括:限制内存、将过载向上游传递(背压)、限制排队造成的尾延迟,并为超时/拒绝策略提供边界。

一个可关闭的阻塞队列

cpp
#include <condition_variable>
#include <cstddef>
#include <mutex>
#include <optional>
#include <queue>
#include <stdexcept>
#include <utility>

// close() 后不再接受 push;已入队元素仍可被 pop 取完。
template <typename T>
class BlockingQueue {
public:
    explicit BlockingQueue(std::size_t capacity) : capacity_(capacity) {
        if (capacity == 0) {
            throw std::invalid_argument("capacity must be positive");
        }
    }

    // 队列满时阻塞;关闭后返回 false。
    bool push(T value) {
        std::unique_lock lock(mutex_);
        not_full_.wait(lock, [&] { return closed_ || queue_.size() < capacity_; });
        if (closed_) return false;

        queue_.push(std::move(value));
        lock.unlock();
        not_empty_.notify_one();
        return true;
    }

    // 队列为空时阻塞;关闭且已取空时返回 std::nullopt。
    std::optional<T> pop() {
        std::unique_lock lock(mutex_);
        not_empty_.wait(lock, [&] { return closed_ || !queue_.empty(); });
        if (queue_.empty()) return std::nullopt;

        T value = std::move(queue_.front());
        queue_.pop();
        lock.unlock();
        not_full_.notify_one();
        return value;
    }

    void close() {
        std::lock_guard lock(mutex_);
        closed_ = true;
        not_empty_.notify_all();
        not_full_.notify_all();
    }

private:
    const std::size_t capacity_;
    std::mutex mutex_;
    std::condition_variable not_empty_;
    std::condition_variable not_full_;
    std::queue<T> queue_;
    bool closed_ = false;
};

实际服务还要决定队列满时的策略:阻塞提交者、立即拒绝(try_push)、等待一小段时间后超时、丢弃低优先级任务,或执行降级。上例将 capacity == 0 视为非法参数并抛出异常。

SPSC、SPMC、MPSC 与 MPMC

队列模型的区别在于入队位置(tail)和出队位置(head)分别有多少线程竞争:

模型生产者消费者竞争点常见实现
SPSC11无索引写竞争无锁有界环形队列
SPMC1N多消费者竞争 headmutex 队列;高级实现用消费者 CAS
MPSCN1多生产者竞争 tailmutex 队列;高级实现用生产者 CAS
MPMCNNheadtail 都竞争有界 mutex 队列;高级实现用每槽位序列号 + CAS

前文的 BlockingQueue 是最稳妥的 MPMC 实现:多个生产者和消费者均可调用它。mutex 只保护短暂的入队/出队操作;消费者拿到任务后释放锁再执行业务,所以业务工作仍可并行。

SPSC:无锁有界环形队列

SPSC 中生产者独占写 tail,消费者独占写 head,无需 CAS 抢占位置。下面是非阻塞的简化 C++17 实现;它空出一个槽位以区分“队空”和“队满”,实际容量为 Capacity - 1

cpp
#include <array>
#include <atomic>
#include <cstddef>
#include <optional>
#include <utility>

template <typename T, std::size_t Capacity>
class SpscRingQueue {
    static_assert(Capacity > 1);

public:
    bool try_push(T value) {
        const auto tail = tail_.load(std::memory_order_relaxed);
        const auto next = (tail + 1) % Capacity;

        // acquire:确认消费者已完成该槽位的读取。
        if (next == head_.load(std::memory_order_acquire)) return false;

        slots_[tail].emplace(std::move(value));
        // release:发布槽位数据,消费者随后才能读取新的 tail。
        tail_.store(next, std::memory_order_release);
        return true;
    }

    std::optional<T> try_pop() {
        const auto head = head_.load(std::memory_order_relaxed);
        // acquire:确认生产者已完整写入该槽位。
        if (head == tail_.load(std::memory_order_acquire)) return std::nullopt;

        T value = std::move(*slots_[head]);
        slots_[head].reset();
        // release:通知生产者该槽位现在可以复用。
        head_.store((head + 1) % Capacity, std::memory_order_release);
        return value;
    }

private:
    std::array<std::optional<T>, Capacity> slots_;
    // 64 是常见缓存行大小;实际项目应按平台并用 profiler 验证。
    alignas(64) std::atomic<std::size_t> head_{0};
    alignas(64) std::atomic<std::size_t> tail_{0};
};

这里 release/acquire 建立了“写入元素 → 发布 tail → 读取元素”的可见性关系。该实现只能用于严格的单生产者和单消费者;多个线程同时调用 try_pushtry_pop 都会破坏正确性。

SPMC / MPSC / MPMC:优先选择有界 mutex 队列

  • SPMC:生产者只有一个,但多个消费者要争抢 head,必须保证每个任务只被一个消费者取得。mutex 队列中,消费者在锁内 pop,拿到任务后在锁外并行处理即可。
  • MPSC:多个生产者要争抢 tail,消费者只有一个;常见于多个业务线程向单日志线程或事件循环递交任务。
  • MPMC:两侧均竞争。前文的 mutex + not_empty + not_full + closed 队列是校招和大多数业务场景的首选回答。

无锁 SPMC/MPSC/MPMC 不能只对索引做 fetch_add:还需处理队空/队满、元素发布、环形槽位复用及竞争失败重试。正确的有界 MPMC 通常为每个槽位维护序列号,并用 CAS 争抢入队/出队位置;还要考虑内存序、ABA、对象回收和伪共享。因此没有明确性能瓶颈时,不应贸然手写无锁 MPMC。

游戏服务器:按实体分片保证顺序

普通 MPMC 队列只保证任务不会被重复取走,不能保证同一房间的任务完成顺序。若同一个玩家或房间必须串行处理,可按 player_idroom_id 路由到固定 worker:

cpp
const std::size_t worker_index =
    std::hash<RoomId>{}(room_id) % worker_count;
text
不同房间 → 不同 worker → 可并行
同一房间 → 固定 worker → 队列顺序消费 → 串行处理

若只有一个网络投递线程,每个分片队列可采用 SPSC;若多个网络线程会投递到同一分片队列,则至少需要 MPSC。哈希分片只保证同 key 路由一致;如果消息从多个来源并发到达,还应由协议序列号或单一入口定义输入顺序。

线程池

线程池复用固定/受控数量的工作线程,避免为每个任务反复创建线程。典型组件:

  • 工作线程集合;
  • 有界任务队列;
  • 条件变量或信号量;
  • 停止标志与优雅关闭逻辑;
  • 队列满时的拒绝/超时策略;
  • 排队时间、活跃线程、拒绝数、任务耗时等监控指标。

最小实现

cpp
#include <functional>
#include <stdexcept>
#include <thread>
#include <vector>

class ThreadPool {
public:
    ThreadPool(std::size_t thread_count, std::size_t queue_capacity)
        : tasks_(queue_capacity) {
        if (thread_count == 0) {
            throw std::invalid_argument("thread_count must be positive");
        }
        workers_.reserve(thread_count);
        for (std::size_t i = 0; i < thread_count; ++i) {
            workers_.emplace_back([this] {
                while (auto task = tasks_.pop()) {
                    try {
                        (*task)();
                    } catch (...) {
                        // 生产环境应记录异常;绝不能让异常逃出线程入口函数。
                    }
                }
            });
        }
    }

    ~ThreadPool() {
        tasks_.close();
        for (auto& worker : workers_) worker.join();
    }

    ThreadPool(const ThreadPool&) = delete;
    ThreadPool& operator=(const ThreadPool&) = delete;

    // 简化接口:队列满时会阻塞;关闭后返回 false。
    bool submit(std::function<void()> task) {
        return tasks_.push(std::move(task));
    }

private:
    BlockingQueue<std::function<void()>> tasks_;
    std::vector<std::thread> workers_;
};

该示例强调生命周期:析构时先关闭队列,唤醒等待线程;工作线程处理完已入队任务后退出,再 join()。生产线程池往往还需要返回 future、任务取消、超时、优先级、动态扩缩容、异常上报和防止“工作线程向同一满队列递交任务导致自我死锁”的策略。

线程数不是越多越好:CPU 密集型任务通常从接近可用 CPU 核数开始压测;I/O 阻塞任务可设置更多线程,但上限应由连接数、下游容量、上下文切换与延迟目标共同决定。

伪共享(False Sharing)

CPU 缓存以 Cache Line 为单位维护一致性(许多机器为 64 字节,但不是语言保证)。若两个线程频繁写入不同变量,但变量恰好落在同一缓存行,缓存行会在核心间反复失效和转移,即使不存在逻辑数据竞争,性能也可能很差。

cpp
struct Counters {
    std::atomic<long> a{0};
    std::atomic<long> b{0}; // 可能与 a 落在同一缓存行
};

缓解方式:分离热点写变量、每线程维护局部计数后汇总、按实际平台缓存行大小对齐或填充。例如 C++17 可使用 std::hardware_destructive_interference_size(实现提供时)表达近似的隔离对齐需求:

cpp
#include <new>

struct alignas(std::hardware_destructive_interference_size) Counter {
    std::atomic<long> value{0};
};

要先用 profiler 验证问题。盲目 padding 会增加内存占用、降低缓存利用率。

无锁一定更快吗?

不一定。无锁算法避免了传统 mutex 的阻塞,但可能引入:

  • ABA 问题;
  • 节点/对象安全回收困难(hazard pointer、epoch 等);
  • 高竞争下 CAS 重试与原子操作争用;
  • Cache Line 在核心间抖动;
  • 更高的实现、验证和排障成本。

低竞争场景下,std::mutex 的路径通常很短,代码更简单、正确性更易证明,可能比复杂无锁结构更快。选择前应先明确吞吐、延迟、竞争程度和可维护性目标,并用真实负载压测。

使用 Markdown 与 VitePress 构建