什么是生产者消费者模型
生产者消费者模式通过一个容器来解决生产者和消费者的强耦合问题。生产者和消费者彼此之间不直接通信,而是通过阻塞队列进行通信。生产者生产完数据后不等待消费者处理,直接放入阻塞队列;消费者不从生产者处获取数据,而是从阻塞队列中取。阻塞队列相当于一个缓冲区,平衡了消费者和生产者的处理能力,实现了二者的解耦。

超市的现实例子
生活中买东西通常去超市而不是直接找供货商。假如需要买一桶方便面,直接找供货商可能不会成功,因为工厂生产是批量进行的,单件购买会导致成本过高且库存积压。现实生活中,供货商通过超市进行生产和消费的解耦。消费者(人)不需要直接向生产者(供货商)要数据,只需从超市(缓冲区)拿取即可。这样既平衡了处理能力,又避免了直接耦合带来的效率问题。
线程视角下的生产者消费者模型
生产者消费者模式本质上是线程间如何安全高效地进行通信。生产者负责生产数据的线程,消费者负责处理数据的线程,而超市则是一段具有特定结构的内存空间。由于生产者和消费者的数据通过这段共享内存空间通信,因此会产生各种并发问题:
- 消费者 VS 消费者:互斥。当资源充足时竞争不明显,但当资源稀缺(如只剩一桶方便面)时,消费者之间存在激烈的竞争关系,属于互斥。
- 生产者 VS 生产者:互斥。多个生产者同时向队列写入数据,存在竞争关系,需要互斥保护。
- 生产者 VS 消费者:互斥与同步。互斥体现在记录人员记录货物时消费者不可拿走;同步体现在队列为空时消费者需等待生产者供货,队列已满时生产者需等待消费者消费。
main 函数
int main() {
// 设置随机种子
srand(time(nullptr));
// 创建阻塞队列
BlockQueue<Task> *bq = new BlockQueue<Task>;
pthread_t c[3], p[5];
// 创建多生产者线程
for (int i = 0; i < 5; i++) {
pthread_create(p + i, nullptr, Producer, bq);
}
// 创建多消费者线程
for (int i = 0; i < 3; i++) {
pthread_create(c + i, nullptr, Consumer, bq);
}
// 等待线程结束
for (int i = 0; i < 5; i++) {
pthread_join(p[i], nullptr);
}
for (int i = 0; i < 3; i++) {
pthread_join(c[i], nullptr);
}
return 0;
}

从结果来看,一个简单的生产者消费者模型创建出来了。但代码中存在一个关于条件变量使用的常见误区,可能导致程序出现伪唤醒的情况。
生产者线程函数
void *Producer(void *args) {
BlockQueue<Task> *bq = (BlockQueue<Task> *)args;
std::string oper("+-*/%");
while (1) {
int x = rand() % 10;
int y = rand() % 10;
Task task(x, y, oper[rand() % 5]);
// 向阻塞队列中放入任务
// 如果队列已满,会在 push() 内部阻塞
bq->push(task);
std::cout << "生产了一个任务 : ";
task.getTask();
sleep(1);
}
}
消费者线程函数
void *Consumer(void *args) {
BlockQueue<Task> *bq = (BlockQueue<Task> *)args;
while (1) {
// 从阻塞队列中取任务
// 如果队列为空,会在 pop() 内部阻塞
Task task = bq->pop();
std::cout << "消耗了一个任务 : ";
task.run();
sleep(2);
}
}
Task 类定义
class Task {
public:
Task(int x, int y, char oper, int result = 0, int exitcode = 0)
: x_(x), y_(y), oper_(oper), result_(result), exitcode_(exitcode) {}
void run() {
switch (oper_) {
case '+': result_ = x_ + y_; break;
case '-': result_ = x_ - y_; break;
case '*': result_ = x_ * y_; break;
case '/':
if (y_ == 0) { exitcode_ = 1; }
else { result_ = x_ / y_; }
break;
case '%':
if (y_ == 0) { exitcode_ = 2; }
else { result_ = x_ % y_; }
break;
}
printf("%d %c %d = %d[%d]\n", x_, oper_, y_, result_, exitcode_);
}
void getTask() {
printf("%d %c %d = ?\n", x_, oper_, y_);
}
private:
int x_;
int y_;
char oper_;
int result_;
int exitcode_;
};
BlockQueue 阻塞队列实现
template <class T>
class BlockQueue {
public:
BlockQueue(int bqmax = 5) : bqmax_(bqmax) {
pthread_mutex_init(&mutex_, nullptr);
pthread_cond_init(&c_cond_, nullptr);
pthread_cond_init(&p_cond_, nullptr);
}
T pop() {
pthread_mutex_lock(&mutex_);
if (bq_.size() == 0) {
pthread_cond_wait(&c_cond_, &mutex_);
}
T top = bq_.front();
bq_.pop();
pthread_cond_signal(&p_cond_);
pthread_mutex_unlock(&mutex_);
return top;
}
void push(const T& in) {
pthread_mutex_lock(&mutex_);
if (bq_.size() == bqmax_) {
pthread_cond_wait(&p_cond_, &mutex_);
}
bq_.push(in);
pthread_cond_signal(&c_cond_);
pthread_mutex_unlock(&mutex_);
}
~BlockQueue() {
pthread_mutex_destroy(&mutex_);
pthread_cond_destroy(&c_cond_);
pthread_cond_destroy(&p_cond_);
}
private:
std::queue<T> bq_;
int bqmax_;
pthread_mutex_t mutex_;
pthread_cond_t c_cond_;
pthread_cond_t p_cond_;
};
为什么判断条件要先加锁?

void *getTicket(void *args) {
threadDate *td = (threadDate *)args;
while (1) {
pthread_mutex_lock(td->mutex_);
if (tickets > 0) {
usleep(1000);
printf("%s get a tickets , tickets : %d\n", td->threadname.c_str(), tickets);
tickets--;
pthread_mutex_unlock(td->mutex_);
} else {
pthread_mutex_unlock(td->mutex_);
break;
}
}
return nullptr;
}
阻塞队列有两个典型约束:
- 队列满时:生产者不能继续生产
- 队列空时:消费者不能继续消费
这是资源暂时不满足条件的情况。为了防止多消费者拿到同一个数据或多生产者造成数据混乱,在多线程操作时必须加锁。同时,当队列满或为空时需进行判断,若满足条件则阻塞。判断临界资源是否满足条件本身也是在访问临界资源,若不先加锁,多个线程可能同时进入判断,导致数据异常(如票数变负)。因此必须先加锁再判断。
伪唤醒问题及解决方案
当资源不满足时,当前线程挂起阻塞,直到资源就绪。持有锁的线程挂起时会释放锁,以便其他线程申请。条件变量的第二个参数即为互斥锁,用于在线程挂起时释放锁,唤醒后重新申请。
在多生产者多消费者场景下,假设队列已满,一个消费线程消费后调用 pthread_cond_broadcast 唤醒了多个生产者。其中一个获得锁并填充空位后,队列再次满。此时其他被唤醒的生产者若获得锁,可能会误以为资源可用而继续生产,导致数据溢出,这就是伪唤醒。
为避免伪唤醒,应使用循环判断条件是否满足:
T pop() {
pthread_mutex_lock(&mutex_);
while (bq_.size() == 0) {
pthread_cond_wait(&c_cond_, &mutex_);
}
T top = bq_.front();
bq_.pop();
pthread_cond_signal(&p_cond_);
pthread_mutex_unlock(&mutex_);
return top;
}
void push(const T& in) {
pthread_mutex_lock(&mutex_);
while (bq_.size() == bqmax_) {
pthread_cond_wait(&p_cond_, &mutex_);
}
bq_.push(in);
pthread_cond_signal(&c_cond_);
pthread_mutex_unlock(&mutex_);
}
即使生产者申请到锁,再次判断时若资源依旧不满足,条件变量会将其挂起并释放锁,从而避免伪唤醒。
生产者消费者模型是多线程编程中最基础也是最重要的模式。通过阻塞队列,生产者和消费者可以安全、高效地协作,同时避免资源竞争和伪唤醒问题。理解了互斥、同步和条件变量的配合,就能轻松应对线程安全设计和高并发场景。


