多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

【C++】手写线程安全队列:mutex + condition_variable 实现生产者消费者模型

【C++】手写线程安全队列:mutex + condition_variable 实现生产者消费者模型 前面已经分别介绍过互斥锁和条件变量。单独看这些知识点可能比较零散而线程安全队列正好可以把它们串起来。普通的std::queue并不是线程安全的std::queueint tasks;如果一个线程正在tasks.push(10);另一个线程同时tasks.pop();就可能发生数据竞争。因此我们希望封装一个这样的队列ThreadSafeQueueint queue; queue.push(100); int value; queue.wait_and_pop(value);其中push()可以被生产者线程调用wait_and_pop()可以让消费者在没有数据时自动等待。一、普通 queue 为什么不能直接多线程使用先来看一个简单例子#include queue #include thread std::queueint tasks; void Producer() { for (int i 0; i 1000; i) { tasks.push(i); } } void Consumer() { while (!tasks.empty()) { int value tasks.front(); tasks.pop(); } }如果生产者和消费者同时运行std::thread t1(Producer); std::thread t2(Consumer);它们可能同时访问tasks而std::queue本身不会自动加锁。例如消费者执行if (!tasks.empty())刚判断队列不为空生产者或其他消费者就可能修改队列。更典型的问题是if (!tasks.empty()) { int value tasks.front(); tasks.pop(); }这三步并不是一个不可分割的整体检查队列 ↓ 读取队首 ↓ 删除队首因此必须使用互斥锁保护整个操作过程。最简单的写法std::mutex mutex; void Push(int value) { std::lock_guardstd::mutex lock(mutex); tasks.push(value); }取数据时同样加锁bool Pop(int value) { std::lock_guardstd::mutex lock(mutex); if (tasks.empty()) { return false; } value tasks.front(); tasks.pop(); return true; }这样可以保证同一时刻只有一个线程修改队列。不过还有一个问题队列为空时消费者应该怎么办如果不停调用while (!Pop(value)) { }线程会一直空转占用 CPU。所以还需要condition_variable。二、先实现 push 和 try_pop我们先搭建一个基础线程安全队列#include condition_variable #include mutex #include queue templatetypename T class ThreadSafeQueue { public: void push(const T value) { { std::lock_guardstd::mutex lock(mutex_); queue_.push(value); } condition_.notify_one(); } bool try_pop(T value) { std::lock_guardstd::mutex lock(mutex_); if (queue_.empty()) { return false; } value queue_.front(); queue_.pop(); return true; } private: std::queueT queue_; std::mutex mutex_; std::condition_variable condition_; };先看void push(const T value)内部首先std::lock_guardstd::mutex lock(mutex_); queue_.push(value);保证多个生产者不能同时破坏队列内部结构。加入数据以后condition_.notify_one();通知一个正在等待数据的消费者队列中已经有新数据了可以起来检查了。这里故意把notify_one()放在锁作用域外{ std::lock_guardstd::mutex lock(mutex_); queue_.push(value); } condition_.notify_one();而不是std::lock_guardstd::mutex lock(mutex_); queue_.push(value); condition_.notify_one();后一种通常也能保证正确性但唤醒消费者以后消费者还需要重新获取mutex_。如果生产者此时仍然持有锁消费者被唤醒 ↓ 想获取mutex ↓ 生产者还没有释放 ↓ 消费者继续等待所以通常先完成共享数据修改并释放锁再进行通知会更自然。try_pop()则表示尝试取出一个数据如果当前没有数据就立即返回。bool try_pop(T value) { std::lock_guardstd::mutex lock(mutex_); if (queue_.empty()) { return false; } value queue_.front(); queue_.pop(); return true; }使用int value; if (queue.try_pop(value)) { std::cout value \n; } else { std::cout 队列为空\n; }它不会等待。因此try_pop()比较适合有任务就处理 没任务就去做其他事情而线程池工作线程通常需要的是没有任务就睡眠 有任务再醒来这就需要wait_and_pop()。三、wait_and_pop 为什么需要 unique_lock实现void wait_and_pop(T value) { std::unique_lockstd::mutex lock(mutex_); condition_.wait(lock, [this]() { return !queue_.empty(); }); value queue_.front(); queue_.pop(); }这里最重要的一句是condition_.wait(lock, [this]() { return !queue_.empty(); });它表示只要队列为空就继续等待队列不为空以后才继续向下执行。大致等价于while (queue_.empty()) { condition_.wait(lock); }假设队列为空消费者获得mutex_ ↓ 发现queue_为空 ↓ wait释放mutex_ ↓ 消费者睡眠之后生产者queue.push(100);内部执行获得mutex_ ↓ queue_.push(100) ↓ 释放mutex_ ↓ notify_one()消费者被唤醒重新获得mutex_ ↓ 再次检查queue_.empty() ↓ 发现不为空 ↓ 取出数据这里必须使用std::unique_lockstd::mutex而不是std::lock_guardstd::mutex因为wait()睡眠时需要临时unlock醒来以后还需要lockunique_lock支持这种灵活控制。因此std::unique_lockstd::mutex lock(mutex_); condition_.wait(lock, [this]() { return !queue_.empty(); });是线程安全队列中非常经典的一种写法。四、完整线程安全队列实现把几个接口组合起来#include condition_variable #include mutex #include queue templatetypename T class ThreadSafeQueue { public: ThreadSafeQueue() default; void push(const T value) { { std::lock_guardstd::mutex lock(mutex_); queue_.push(value); } condition_.notify_one(); } void push(T value) { { std::lock_guardstd::mutex lock(mutex_); queue_.push(std::move(value)); } condition_.notify_one(); } bool try_pop(T value) { std::lock_guardstd::mutex lock(mutex_); if (queue_.empty()) { return false; } value std::move(queue_.front()); queue_.pop(); return true; } void wait_and_pop(T value) { std::unique_lockstd::mutex lock(mutex_); condition_.wait(lock, [this]() { return !queue_.empty(); }); value std::move(queue_.front()); queue_.pop(); } bool empty() const { std::lock_guardstd::mutex lock(mutex_); return queue_.empty(); } size_t size() const { std::lock_guardstd::mutex lock(mutex_); return queue_.size(); } private: std::queueT queue_; mutable std::mutex mutex_; std::condition_variable condition_; };这里同时提供了两个push()void push(const T value); void push(T value);第一个处理左值int value 100; queue.push(value);第二个支持移动queue.push(100);或者std::string str hello; queue.push(std::move(str));取数据时也使用value std::move(queue_.front());避免某些较大对象发生不必要的复制。另外bool empty() const是const成员函数但里面需要锁住mutex_所以互斥锁声明为mutable std::mutex mutex_;mutable表示即使当前对象是const这个成员仍然允许修改。而加锁和解锁本身会修改 mutex 的内部状态所以需要mutable。五、生产者消费者完整示例下面创建两个生产者和两个消费者。#include iostream #include thread ThreadSafeQueueint queue; void Producer(int start) { for (int i 0; i 5; i) { int value start i; queue.push(value); std::cout 生产 value \n; } } void Consumer() { for (int i 0; i 5; i) { int value; queue.wait_and_pop(value); std::cout 消费 value \n; } } int main() { std::thread producer1(Producer, 100); std::thread producer2(Producer, 200); std::thread consumer1(Consumer); std::thread consumer2(Consumer); producer1.join(); producer2.join(); consumer1.join(); consumer2.join(); return 0; }生产者不断queue.push(value);消费者不断queue.wait_and_pop(value);如果队列中存在数据消费者直接取出如果队列为空消费者进入wait ↓ 释放mutex ↓ 进入睡眠生产者加入任务push数据 ↓ notify_one ↓ 消费者被唤醒这就是最基本的生产者—消费者模型。整个结构其实已经非常接近线程池中的任务队列外部线程 ↓ 提交任务 ↓ ThreadSafeQueue ↓ condition_variable通知 ↓ Worker线程醒来 ↓ wait_and_pop取任务 ↓ 执行任务如果把ThreadSafeQueueint改成ThreadSafeQueuestd::functionvoid()队列里面保存的就不再是数字而是一个个真正可以执行的任务ThreadSafeQueuestd::functionvoid() tasks;提交任务tasks.push([]() { std::cout 执行任务\n; });工作线程std::functionvoid() task; tasks.wait_and_pop(task); task();这样就已经搭出了一个简化线程池最核心的结构任务 ↓ 线程安全队列 ↓ condition_variable ↓ Worker ↓ 执行task()这一篇最需要掌握的是std::queue本身不保证线程安全 mutex负责保护队列内部数据 try_pop没有数据时立即返回 wait_and_pop没有数据时让线程睡眠 condition_variable负责通知等待线程 wait需要unique_lock因为等待期间必须释放mutex push完成数据修改后再notify_one 线程安全队列是生产者消费者模型和线程池的重要基础。0voice · GitHub
返回列表