【C++】线程安全队列(二):生产者消费者模型与 condition_variable 阻塞等待 一、从线程安全队列到生产者消费者模型上一篇中使用std::mutex配合std::deque实现了一个简单的线程安全队列。它解决的核心问题是多个线程同时操作队列 ↓ 使用 mutex 进行保护 ↓ 同一时刻只有一个线程能够修改队列 ↓ 避免数据竞争但是这种简单队列还有一个问题。假设消费者不断从队列中取数据int value; while (queue.next(value)) { // 处理数据 }如果消费者执行到这里时队列恰好为空消费者 ↓ 检查队列 ↓ 队列为空 ↓ next() 返回 false ↓ 消费者结束可是生产者可能只是暂时还没有来得及添加数据。例如消费者启动 ↓ 发现队列为空 ↓ 直接退出 过了一会 生产者 ↓ 添加任务这样新加入的数据就没有消费者处理了。一种简单粗暴的方法是让消费者一直检查while (true) { if (!queue.empty()) { // 获取任务 } }但是这种方式会导致线程不停循环检查队列 ↓ 没有数据 ↓ 继续检查 ↓ 还是没有 ↓ 继续检查……即使没有任何任务CPU 也一直在工作这种情况通常称为忙等待Busy Waiting。更合理的方式应该是队列有数据 ↓ 消费者取数据 队列为空 ↓ 消费者睡眠等待 生产者加入数据 ↓ 通知消费者 消费者被唤醒 ↓ 继续取数据这就是std::condition_variable的主要作用。因此一个比较完整的生产者消费者队列通常需要std::queueT _queue; // 保存数据 std::mutex _queueLock; // 保护队列 std::condition_variable _condition; // 控制线程等待和唤醒三者之间可以简单理解为┌──────────────┐ 生产者 ──Push()──→ │ queue │ ──Pop()──→ 消费者 └──────────────┘ ↑ mutex ↑ condition_variable 等待 / 唤醒线程mutex负责解决多个线程能不能同时修改队列而condition_variable解决的是队列没有数据时消费者应该怎么办二、ProducerConsumerQueue 的基本结构先来看一个生产者消费者队列需要哪些成员变量#include condition_variable #include mutex #include queue #include atomic templatetypename T class ProducerConsumerQueue { private: std::mutex _queueLock; // 保护队列 std::queueT _queue; // 保存生产者产生的数据 std::condition_variable _condition; // 控制消费者等待和唤醒 std::atomicbool _shutdown; // 队列是否停止工作 public: ProducerConsumerQueue() : _shutdown(false) {} };其中_queue是真正保存数据的地方std::queueT _queue;例如ProducerConsumerQueueint queue;内部就相当于std::queueint _queue;队列遵循 FIFOFirst In First Out 先进先出例如push(10) push(20) push(30) 队头 队尾 ↓ ↓ 10 → 20 → 30消费者第一次取出10第二次取出20然后是30第二个成员std::mutex _queueLock;用来保护_queue。生产者执行_queue.push(value);消费者执行_queue.pop();这些操作不能被多个线程同时执行因此需要_queueLock。第三个成员std::condition_variable _condition;它负责消费者线程的等待 唤醒最后std::atomicbool _shutdown;用于表示整个队列是否准备结束。初始状态_shutdown false;表示队列正常运行调用取消操作以后_shutdown true;表示队列准备停止这几个成员组合起来就形成了一个比较典型的生产者消费者队列ProducerConsumerQueue ┌──────────────────────────────┐ │ │ │ std::queueT _queue │ │ ↑ │ │ std::mutex _queueLock │ │ │ │ condition_variable │ │ 等待 / 唤醒 │ │ │ │ atomicbool _shutdown │ │ 控制退出 │ │ │ └──────────────────────────────┘三、Push生产者添加数据并唤醒消费者生产者向队列添加数据可以实现为void Push(const T value) { std::lock_guardstd::mutex lock(_queueLock); // 获取互斥锁 _queue.push(value); // 数据加入队列 _condition.notify_one(); // 唤醒一个等待中的消费者 }首先std::lock_guardstd::mutex lock(_queueLock);给队列加锁。然后_queue.push(value);将数据加入队列。例如queue.Push(10);队列变成┌────┐ │ 10 │ └────┘再次queue.Push(20);变成10 → 20真正值得注意的是最后一句_condition.notify_one();它的意思不是“通知队列”而是唤醒一个正在_condition上等待的线程。假设现在消费者因为队列为空进入等待消费者线程 ↓ 队列为空 ↓ condition_variable ↓ 进入等待状态此时生产者queue.Push(100);执行过程生产者获得锁 ↓ push(100) ↓ 队列中出现数据 ↓ notify_one() ↓ 唤醒一个消费者消费者被唤醒以后继续尝试获取数据queue ↑ 生产者 ──push()──→ [100] │ notify_one() ↓ 消费者 ↓ 唤醒这里有两个比较容易混淆的函数notify_one(); notify_all();notify_one()唤醒一个等待线程。notify_all()唤醒所有等待线程。生产者每次通常只添加一个任务因此_condition.notify_one();一般就足够了。例如三个消费者都在等待消费者1等待 消费者2等待 消费者3等待生产者加入一个任务queue.Push(100);调用notify_one();可能只有消费者2被唤醒消费者1继续等待 消费者2被唤醒 → 获取100 消费者3继续等待这样避免了明明只有一个任务却把所有消费者全部叫醒。另外有些代码中可能会看到_queue.push(std::move(value));如果Push()的参数本身是const T value那么这里的std::move(value)实际得到的是const T很多类型最终仍然会发生拷贝而不是真正的移动。如果希望真正支持移动可以额外提供void Push(T value) { std::lock_guardstd::mutex lock(_queueLock); _queue.push(std::move(value)); _condition.notify_one(); }不过理解生产者消费者模型时暂时把重点放在push notify_one这一组操作上即可。四、Pop 与 WaitAndPop普通取数据和阻塞等待的区别消费者取数据可以有两种方式。第一种是普通的Pop()bool Pop(T value) { std::lock_guardstd::mutex lock(_queueLock); if (_queue.empty() || _shutdown) return false; value _queue.front(); _queue.pop(); return true; }首先获取锁std::lock_guardstd::mutex lock(_queueLock);然后判断if (_queue.empty() || _shutdown) return false;只要队列为空或者队列已经停止就不继续取数据。如果存在数据value _queue.front(); _queue.pop();先获取队头value _queue.front();再删除队头_queue.pop();例如10 → 20 → 30执行queue.Pop(value);之后value 10 队列 20 → 30这种Pop()有一个特点没数据就直接返回不会等待。真正体现生产者消费者模型的是WaitAndPop()void WaitAndPop(T value) { std::unique_lockstd::mutex lock(_queueLock); while (_queue.empty() !_shutdown) _condition.wait(lock); if (_queue.empty() || _shutdown) return; value _queue.front(); _queue.pop(); }这里首先出现std::unique_lockstd::mutex lock(_queueLock);为什么这里不用前面的std::lock_guardstd::mutex而要使用std::unique_lockstd::mutex因为后面需要执行_condition.wait(lock);condition_variable::wait()在等待过程中需要释放 mutex ↓ 线程睡眠 ↓ 收到通知 ↓ 重新获取 mutex这种“中途释放锁、醒来以后重新加锁”的操作需要unique_lock配合完成。最关键的代码就是while (_queue.empty() !_shutdown) _condition.wait(lock);可以先把条件翻译成人话队列为空 并且 程序还没有关闭那么_condition.wait(lock);消费者进入等待。假设当前_queue.empty() true _shutdown false那么消费者执行进入 wait ↓ 释放 mutex ↓ 消费者线程睡眠为什么wait()必须把锁释放因为如果消费者睡着以后还一直拿着_queueLock消费者拿着锁睡觉 ↓ 生产者 Push() ↓ 想获得 _queueLock ↓ 拿不到 ↓ 无法添加数据这样生产者永远无法把数据放进去消费者也永远等不到数据。所以_condition.wait(lock);内部非常重要的一件事就是等待时自动释放互斥锁。完整流程可以理解为消费者执行 WaitAndPop() ↓ 获得 _queueLock ↓ 检查队列 ↓ 队列为空 ↓ wait(lock) ↓ 释放 _queueLock ↓ 消费者睡眠 ↓ -------------------------- ↓ 生产者获得 _queueLock ↓ push() ↓ notify_one() ↓ -------------------------- ↓ 消费者被唤醒 ↓ 重新竞争 _queueLock ↓ 获得锁 ↓ 重新检查条件 ↓ front() ↓ pop()这里还有一个非常重要的细节while (_queue.empty() !_shutdown)为什么是while而不是if (_queue.empty() !_shutdown)因为条件变量存在虚假唤醒Spurious Wakeup。线程被唤醒并不意味着队列一定有数据因此醒来之后还必须重新检查_queue.empty()所以通常应该写成while (条件不满足) condition.wait(lock);C 还提供了更加常见的谓词版本_condition.wait(lock, [this]() { return !_queue.empty() || _shutdown; });这句可以理解为一直等待直到队列中有数据或者整个队列准备退出。实际上它内部的思想和while (_queue.empty() !_shutdown) _condition.wait(lock);基本一致。使用谓词以后WaitAndPop()可以写成void WaitAndPop(T value) { std::unique_lockstd::mutex lock(_queueLock); _condition.wait(lock, [this]() { return !_queue.empty() || _shutdown; }); if (_queue.empty() || _shutdown) return; value _queue.front(); _queue.pop(); }这也是实际编写条件变量代码时很常见的一种写法。五、Cancel如何让正在等待的消费者安全退出最后还需要解决一个问题如果程序准备结束但是消费者还在wait()里面睡觉怎么办假设现在有三个消费者消费者1 → wait() 消费者2 → wait() 消费者3 → wait()此时程序准备退出。如果什么都不做这些线程可能仍然处于等待状态。因此需要提供void Cancel()通知整个队列结束。一个基本实现为void Cancel() { { std::lock_guardstd::mutex lock(_queueLock); _shutdown true; } _condition.notify_all(); }首先_shutdown true;告诉消费者队列已经停止工作然后_condition.notify_all();把所有正在等待的消费者全部唤醒。为什么这里不是notify_one();因为程序准备退出时希望消费者1 消费者2 消费者3全部醒来并退出而不是只叫醒其中一个。执行过程Cancel() ↓ _shutdown true ↓ notify_all() ↓ ┌─────────┬─────────┬─────────┐ ↓ ↓ ↓ 消费者1 消费者2 消费者3 唤醒 唤醒 唤醒 ↓ ↓ ↓ 发现 shutdown true ↓ ↓ ↓ 退出 退出 退出把前面的内容组合起来一个简单的生产者消费者队列可以写成#include atomic #include condition_variable #include mutex #include queue templatetypename T class ProducerConsumerQueue { private: std::mutex _queueLock; // 保护队列 std::queueT _queue; // 保存数据 std::condition_variable _condition; // 控制等待和唤醒 std::atomicbool _shutdown; // 是否停止 public: ProducerConsumerQueue() : _shutdown(false) {} // 生产者添加数据 void Push(const T value) { { std::lock_guardstd::mutex lock(_queueLock); _queue.push(value); } _condition.notify_one(); } // 普通取数据没有数据直接返回 false bool Pop(T value) { std::lock_guardstd::mutex lock(_queueLock); if (_queue.empty() || _shutdown) return false; value _queue.front(); _queue.pop(); return true; } // 阻塞取数据没有数据就等待 bool WaitAndPop(T value) { std::unique_lockstd::mutex lock(_queueLock); _condition.wait(lock, [this]() { return !_queue.empty() || _shutdown; }); if (_shutdown _queue.empty()) return false; value _queue.front(); _queue.pop(); return true; } // 停止队列并唤醒所有消费者 void Cancel() { { std::lock_guardstd::mutex lock(_queueLock); _shutdown true; } _condition.notify_all(); } };然后可以创建一个简单的生产者消费者程序#include iostream #include thread int main() { ProducerConsumerQueueint queue; // 消费者没有数据时会阻塞等待 std::thread consumer([]() { int value; while (queue.WaitAndPop(value)) { std::cout consumer : value std::endl; } }); // 生产者不断向队列中添加数据 std::thread producer([]() { queue.Push(10); queue.Push(20); queue.Push(30); queue.Push(40); queue.Cancel(); }); producer.join(); consumer.join(); return 0; }整个程序的工作流程就是消费者启动 ↓ WaitAndPop() ↓ 队列为空 ↓ wait() 睡眠 ↓ ════════════════════ ↓ 生产者 Push(10) ↓ notify_one() ↓ ════════════════════ ↓ 消费者被唤醒 ↓ 取出10 ↓ 继续 WaitAndPop()因此相比上一篇简单的deque mutex这一篇又增加了condition_variable最终形成mutex ↓ 生产者 → 线程安全队列 → 消费者 │ ↑ └──── notify_one ───────┘ │ wait()这一部分最需要掌握的其实就是下面几组对应关系std::mutex // 保护共享队列 std::lock_guardstd::mutex // 简单加锁 std::unique_lockstd::mutex // 配合 condition_variable 使用以及_condition.wait(lock); // 消费者等待 _condition.notify_one(); // 唤醒一个消费者 _condition.notify_all(); // 唤醒所有消费者最终把整个生产者消费者模型概括成一句话就是生产者负责向队列中放数据消费者负责从队列中取数据mutex 保证队列访问安全condition_variable 负责在没有数据时让消费者休眠并在新数据到来后将其唤醒。理解这一套流程之后再继续学习线程池会非常自然因为线程池中的提交任务 ↓ 任务队列 ↓ 工作线程等待 ↓ 有任务后唤醒 ↓ 执行任务本质上就是生产者消费者模型的一种典型应用。0voice · GitHub