十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

生产者—消费者模型详解:从 PV 操作到 Linux pthread 实现

生产者—消费者模型详解:从 PV 操作到 Linux pthread 实现 生产者—消费者模型是操作系统中最经典的线程同步与互斥问题之一也是理解互斥锁、信号量、条件变量、阻塞队列以及线程池的基础。它看起来只是“一个线程放数据一个线程拿数据”但真正需要解决的是共享缓冲区如何保证线程安全、生产者什么时候应该阻塞、消费者什么时候应该阻塞以及线程应该如何被正确唤醒。一、什么是生产者—消费者问题假设系统中存在一块共享缓冲区共享缓冲区 ┌───────────────────┐ 生产者 ──→ [ ][ ][ ][ ][ ] ──→ 消费者 └───────────────────┘ 最大容量 N生产者负责产生数据并放入缓冲区Producer ↓ 生产数据 ↓ 放入 Buffer消费者负责从缓冲区取出数据并处理Buffer ↓ 取出数据 ↓ Consumer例如生产者 网络线程接收请求 ↓ 任务队列 ↓ 消费者 工作线程处理请求这实际上就是大量服务器线程池的基本结构。如果缓冲区是有界缓冲区例如最多只能保存 5 个元素那么还必须满足缓冲区未满 ↓ 生产者才能继续生产 缓冲区非空 ↓ 消费者才能继续消费因此Buffer 满了 → Producer 必须等待 Buffer 空了 → Consumer 必须等待这里同时存在两个问题互斥问题多个线程可能同时访问共享缓冲区因此同一时刻必须限制对临界区的并发修改。同步问题生产者和消费者之间存在执行顺序约束生产者 Buffer 没空间 → 等待消费者 消费者 Buffer 没数据 → 等待生产者所以生产者—消费者问题不是单纯的互斥问题而是“互斥 同步”问题。二、经典 PV 信号量模型操作系统教材和 408 中最经典的写法使用三个信号量mutex 1 empty N full 0其中mutex → 控制对缓冲区的互斥访问 empty → 当前还有多少个空位置 full → 当前有多少个产品假设缓冲区容量N 5初始状态Buffer [ ][ ][ ][ ][ ] empty 5 full 0 mutex 1生产者执行Producer() { 生产一个产品; P(empty); P(mutex); 将产品放入缓冲区; V(mutex); V(full); }消费者执行Consumer() { P(full); P(mutex); 从缓冲区取出产品; V(mutex); V(empty); 消费产品; }整个关系可以理解成empty ↓ Producer → [ Buffer ] → Consumer ↑ full mutex保护Buffer为什么生产者必须先P(empty)再P(mutex)这个顺序非常重要。正确P(empty) P(mutex) 生产 V(mutex) V(full)如果错误地写成P(mutex) P(empty)假设缓冲区已经满了empty 0生产者先P(mutex)把缓冲区锁住。然后P(empty)发现没有空间于是生产者阻塞。但是消费者想消费P(mutex)发现 mutex 被生产者占着于是消费者也阻塞。最终生产者 等待消费者腾出空间 消费者 等待生产者释放 mutex形成等待 Producer ─────→ Consumer ↑ │ └──────────────┘ 等待 DEADLOCK所以必须先判断资源条件再进入临界区P(empty/full) ↓ P(mutex) ↓ 访问 Buffer这是生产者—消费者 PV 操作中非常经典的考点。三、Linux 中如何使用 mutex 条件变量实现在真正的 Linux 多线程程序中生产者—消费者模型经常使用pthread_mutex_t pthread_cond_t实现。假设我们有mutex not_empty not_full三者作用分别是mutex → 保护共享队列 not_empty → Buffer 非空这个条件 not_full → Buffer 未满这个条件生产者逻辑加 mutex while(Buffer 满) { 等待 not_full } push 数据 通知 not_empty 释放 mutex消费者逻辑加 mutex while(Buffer 空) { 等待 not_empty } pop 数据 通知 not_full 释放 mutex注意条件变量本身并不是锁。mutex 解决 “谁可以访问共享队列” condition variable 解决 “什么时候才应该继续执行”这两个概念一定要区分。例如消费者发现queue.empty()它并不是while (queue.empty()) { }这样死循环等待。否则就是Busy Waiting 忙等待CPU 会不断检查empty? empty? empty? empty? ...浪费 CPU。正确方式是pthread_cond_wait(not_empty, mutex);让线程真正进入睡眠等待。pthread_cond_wait()有一个特别重要的行为它会原子地释放 mutex 并进入等待线程被唤醒以后又会在返回之前重新取得 mutex。所以pthread_mutex_lock() ↓ pthread_cond_wait() │ ├── 自动释放 mutex │ ├── 当前线程休眠 │ ├── 被 signal 唤醒 │ └── 重新获得 mutex ↓ 继续执行这也是条件变量必须和 mutex 配合使用的重要原因。四、为什么pthread_cond_wait()外面必须使用while不能只用if这是生产者—消费者代码中最重要的细节之一。错误写法if (queue.empty()) { pthread_cond_wait(not_empty, mutex); } data queue.front();更正确的写法是while (queue.empty()) { pthread_cond_wait(not_empty, mutex); }原因是线程被唤醒不代表条件现在一定成立。假设有两个消费者Consumer A Consumer B现在queue.empty()所以A 等待 B 等待生产者生产一个数据queue [100]然后唤醒消费者。假设由于调度等原因多个消费者竞争 mutexA 先拿到A pop 100于是queue 又空了B 随后获得 mutex。如果 B 使用if它可能直接继续向下执行queue.front();但此时queue.empty() true于是程序出错。而使用whileB 被唤醒后会重新判断queue.empty() ?如果仍然为空继续 wait所以应该记住条件变量 通知只是提醒你“条件可能发生变化” 真正判断条件是否成立 必须检查共享状态本身因此经典写法永远是while (!条件成立) { pthread_cond_wait(...); }而不是if (!条件成立) { pthread_cond_wait(...); }五、完整 pthread 生产者—消费者代码下面实现一个有界阻塞队列支持多个生产者和多个消费者。#include iostream #include queue #include pthread.h #include unistd.h using namespace std; class BlockQueue { public: explicit BlockQueue(size_t capacity 5) : _capacity(capacity), _stop(false) { pthread_mutex_init(_mutex, nullptr); pthread_cond_init(_notEmpty, nullptr); pthread_cond_init(_notFull, nullptr); } ~BlockQueue() { pthread_mutex_destroy(_mutex); pthread_cond_destroy(_notEmpty); pthread_cond_destroy(_notFull); } // 生产数据 bool push(int data) { pthread_mutex_lock(_mutex); // 队列满了生产者等待 while (_queue.size() _capacity !_stop) { pthread_cond_wait(_notFull, _mutex); } // 队列已经停止 if (_stop) { pthread_mutex_unlock(_mutex); return false; } // 临界区向共享队列中插入数据 _queue.push(data); // Buffer 已经非空可以唤醒消费者 pthread_cond_signal(_notEmpty); pthread_mutex_unlock(_mutex); return true; } // 消费数据 bool pop(int data) { pthread_mutex_lock(_mutex); // 队列为空消费者等待 while (_queue.empty() !_stop) { pthread_cond_wait(_notEmpty, _mutex); } /* * stop 后仍然允许消费者把队列中的 * 剩余数据全部消费完。 */ if (_queue.empty() _stop) { pthread_mutex_unlock(_mutex); return false; } // 临界区取出数据 data _queue.front(); _queue.pop(); // Buffer 已经存在空位置可以唤醒生产者 pthread_cond_signal(_notFull); pthread_mutex_unlock(_mutex); return true; } // 停止阻塞队列 void stop() { pthread_mutex_lock(_mutex); _stop true; /* * 所有正在等待的生产者、消费者 * 都必须被唤醒否则可能永久阻塞。 */ pthread_cond_broadcast(_notEmpty); pthread_cond_broadcast(_notFull); pthread_mutex_unlock(_mutex); } private: queueint _queue; size_t _capacity; bool _stop; pthread_mutex_t _mutex; // 消费者等待队列非空 pthread_cond_t _notEmpty; // 生产者等待队列未满 pthread_cond_t _notFull; }; /*-------------------------------- 生产者 --------------------------------*/ struct ThreadData { int id; BlockQueue* queue; }; void* producer(void* args) { ThreadData* td static_castThreadData*(args); for (int i 0; i 10; i) { int data td-id * 100 i; if (!td-queue-push(data)) { break; } cout Producer td-id produce: data endl; usleep(100000); } return nullptr; } /*-------------------------------- 消费者 --------------------------------*/ void* consumer(void* args) { ThreadData* td static_castThreadData*(args); int data; while (td-queue-pop(data)) { cout Consumer td-id consume: data endl; usleep(200000); } return nullptr; } /*-------------------------------- main --------------------------------*/ int main() { BlockQueue queue(5); pthread_t producers[2]; pthread_t consumers[2]; ThreadData producerArgs[2]; ThreadData consumerArgs[2]; // 创建两个消费者 for (int i 0; i 2; i) { consumerArgs[i].id i; consumerArgs[i].queue queue; pthread_create( consumers[i], nullptr, consumer, consumerArgs[i]); } // 创建两个生产者 for (int i 0; i 2; i) { producerArgs[i].id i; producerArgs[i].queue queue; pthread_create( producers[i], nullptr, producer, producerArgs[i]); } // 等待所有生产者结束 for (int i 0; i 2; i) { pthread_join(producers[i], nullptr); } /* * 生产者全部退出后 * 不再产生新数据。 */ queue.stop(); // 等待消费者消费剩余数据并退出 for (int i 0; i 2; i) { pthread_join(consumers[i], nullptr); } return 0; }编译g producer_consumer.cpp -o producer_consumer -pthread这里最核心的数据结构是BlockQueue ┌────────────────┐ Producer ──→│ queueint │──→ Consumer │ │ │ capacity 5 │ └────────────────┘ ↑ ↑ │ │ notFull notEmpty mutex ↓ 保护整个 queue注意notEmpty 和 notFull不是两把锁。真正的锁只有_mutex两个条件变量只是分别对应_notEmpty → queue.size() 0 _notFull → queue.size() capacity因为两个条件都依赖同一个_queue的状态所以使用同一把 mutex保护队列状态是最自然、最容易保证正确性的设计。六、这份代码真正解决了哪些并发问题首先是互斥访问。假设Producer A Producer B Consumer C三者同时访问_queue如果完全不加锁可能出现A 正在修改 queue B 同时修改 queue C 同时 pop而std::queue本身并不保证这种并发修改安全。所以必须pthread_mutex_lock(_mutex);保证同一时刻 只有一个线程 修改 queue其次是线程同步。如果队列已经满[1][2][3][4][5] capacity 5生产者不是不停地轮询while (_queue.size() 5) { }而是pthread_cond_wait(_notFull, _mutex);进入阻塞。消费者拿走数据[2][3][4][5][ ]再pthread_cond_signal(_notFull);通知一个生产者现在有空位置了消费者也是完全相同的道理。所以整个模型实际是queue 满 ↓ Producer 睡眠 │ │ Consumer 消费 ↓ signal(notFull) ↓ Producer 被唤醒 queue 空 ↓ Consumer 睡眠 │ │ Producer 生产 ↓ signal(notEmpty) ↓ Consumer 被唤醒另外一个非常容易遗漏的问题是线程退出。如果消费者正在pthread_cond_wait()而生产者已经全部退出再也没有人生产数据生产者全部死亡 queue.empty() 消费者正在 wait那么消费者可能永远无法退出。所以工程代码通常需要bool stop;停止时pthread_cond_broadcast(_notEmpty); pthread_cond_broadcast(_notFull);把所有阻塞线程唤醒。消费者再检查if (_queue.empty() _stop) { return false; }这样线程才能安全结束。七、生产者—消费者模型和线程池有什么关系生产者—消费者并不是只存在于操作系统教材里。很多线程池本质上就是生产者—消费者模型。例如服务器主线程 / 网络线程 │ Producer │ ▼ ┌──────────────────┐ │ Task Queue │ └──────────────────┘ │ │ │ ▼ ▼ ▼ Worker1 Worker2 Worker3 │ │ │ └── Consumers ─┘主线程收到任务taskQueue.push(task);工作线程while (true) { task taskQueue.pop(); task(); }所以线程池 任务队列 若干 Consumer Worker 生产者提交任务你现在学的pthread mutex condition variable producer-consumer其实很快就可以直接组合成一个自己的线程池。这也是为什么后端/C 面试经常从生产者消费者继续追问阻塞队列 ↓ 线程池 ↓ 高并发服务器公开面经中也能看到这种联系字节的后台/C 面试出现过现场实现生产者—消费者并继续追问生产者和消费者是否应该共用一把锁腾讯 C 后台面试则经常进一步问线程池如何设计以及如何做高性能线程池。八、大厂面试中生产者—消费者常见追问生产者—消费者在面试中一般不会只问“什么是生产者—消费者”更常见的是从简单概念一路往下追。字节后台岗位曾直接要求实现生产者—消费者模型也出现过“队列为空时消费者怎么办”“生产者和消费者是否应该使用同一把锁”等追问腾讯公开面经中出现过“两线程生成随机字符串一个线程负责消费打印”的多线程题也有项目中直接追问生产者—消费者模型的案例。需要注意这些是公开候选人面经并不是腾讯或字节官方题库。面试时建议重点准备下面这些问题。1. 什么是生产者—消费者模型可以回答多个执行流通过共享缓冲区交换数据。生产者负责向缓冲区写入数据消费者负责取出数据缓冲区属于临界资源因此需要互斥访问同时生产者与消费者之间还存在“缓冲区不能满写、不能空读”的同步约束所以生产者—消费者模型本质上是互斥与同步结合的问题。2. 为什么只有 mutex 不够因为mutex 只能解决 不能同时操作 queue 但解决不了 queue 空了消费者怎么办 queue 满了生产者怎么办如果只有 mutex就可能只能不断while (queue.empty()) { }形成忙等待。所以通常还需要Semaphore 或 Condition Variable实现线程阻塞和唤醒。3. 为什么需要两个条件变量因为我们实际上存在两个不同的条件消费者关心 queue ! empty 生产者关心 queue ! full所以notEmpty → 唤醒消费者 notFull → 唤醒生产者一把 mutex 保护共享队列两个 condition variable 分别等待不同条件。4. 为什么pthread_cond_wait()必须配合 mutex因为检查条件 进入睡眠必须正确衔接。pthread_cond_wait()会在等待时原子地释放 mutex并在返回前重新获得 mutex从而避免线程在“检查条件”和“真正睡眠”之间发生典型的通知丢失竞争。5. 为什么必须用while不能用if因为被唤醒 ≠ 条件一定成立线程醒来以后必须重新检查共享状态所以while (queue.empty()) { pthread_cond_wait(...); }比if (queue.empty()) { pthread_cond_wait(...); }更加正确。这是非常高频的并发追问。6.signal和broadcast有什么区别pthread_cond_signal()通常唤醒至少一个等待线程适合push 一个任务 → 唤醒一个消费者而pthread_cond_broadcast()用于唤醒所有等待者。例如程序准备退出此时所有消费者都必须知道不用再等了所以通常stop true; pthread_cond_broadcast(...);7. 生产者和消费者可以使用两把 mutex 吗对于普通的std::queue有界阻塞队列通常使用一把 mutex保护完整的队列状态更加合理。因为push pop size empty full本质上属于同一个共享数据结构的不变量。如果简单粗暴地把Producer → mutex1 Consumer → mutex2分开并不能保证std::queue的并发安全因为生产和消费仍然可能同时修改同一个容器内部状态。当然高性能队列可以采用头尾分离锁 CAS Atomic 无锁 RingBuffer等更加复杂的结构但那已经属于更高级的并发数据结构设计问题。8. 如果让你设计线程池生产者—消费者怎么用标准回答可以从submit() ↓ 任务入队 ↓ TaskQueue ↓ condition_variable ↓ 唤醒 Worker ↓ Worker pop task ↓ 执行任务开始。然后面试官可能继续追问线程池怎么停止 任务队列满了怎么办 Worker 数量怎么确定 如何防止惊群 如何实现优雅退出 任务抛异常怎么办 如何支持返回值 如何动态扩容线程 如何降低锁竞争 能不能实现无锁任务队列这时生产者—消费者模型就已经从一道操作系统基础题变成真正的并发编程与服务器设计题了。总结生产者—消费者模型真正需要掌握的不是死记代码而是下面这条逻辑共享 Buffer │ 是临界资源 ↓ mutex │ 保证同一时刻安全访问 │ ┌──────────┴──────────┐ │ │ Buffer 空 Buffer 满 │ │ Consumer 等待 Producer 等待 │ │ notEmpty notFull ↑ ↑ │ │ Producer 唤醒 Consumer 唤醒因此最终只需要牢牢记住mutex 解决“能不能同时访问”的互斥问题condition variable / semaphore 解决“什么时候可以继续执行”的同步问题。经典 PV 模型Producer: P(empty) P(mutex) 生产 V(mutex) V(full)Consumer: P(full) P(mutex) 消费 V(mutex) V(empty)Linux pthread 模型Producer: lock while(full) wait(notFull) push signal(notEmpty) unlockConsumer: lock while(empty) wait(notEmpty) pop signal(notFull) unlock
返回列表