Linux 生产者消费者模型
1.什么是生产者消费者模型
生产者消费者模型(consumer producter),也叫做cp问题。

在日常生活中,都会有超市的存在,超市上面摆放着非常非常多的货物,这些货物都是由工厂生产出来,之后拉到超市,摆放到货架上面,由顾客来超市时,就可以购买这些货物。
那么生产货物的工厂就相当于是生产者,那些购买货物的顾客就相当于是消费者,超市就是一个共享内存缓冲区,用来临时存放数据。当货架上面的货物满了就意味着工厂就应该先停止生产,等待货架上的货物被购课购买走了一些,在继续生产,同时,如果货架上面的货物空了,顾客也就应该停止购买了,等待货物上架了,再继续购买。
假设临近过年了,按照往年经验供应商预计顾客会对货物的需求量剧增,但是过年的时候,工厂的工人要放假无法生产货物了,而过年的时候,顾客对货物的需求量又很大,所以怎么办呢?所以多个供应商在过年前的前一个月加大工厂的产量,确保将超市的货架摆满,于是过年的时候,工厂的工人放假了,消费者就去超市中购买货物,此时生产者已经不生产货物了,但是消费者仍能购买到很多货物
所以超市作为一个大号的缓存,生产和消费就不必同步进行,可以将生产和消费的行为进行一定程度的解耦。
同时,因为超市是共享资源,就需要被保护起来,在多线程的环境下,多个生产者,多个消费者同时的访问数据会出现数据的不安全情况,所以需要使用锁来保证同一时间只有一个执行流操作缓冲区。
由此形成3种关系:
1. 生产者与生产者之间:互斥
多个生产者竞争向缓冲区放数据,同一时间只能有一个生产者操作。
2. 消费者与消费者之间:互斥
多个消费者竞争从缓冲区取数据,同一时间只能有一个消费者操作。
3. 生产者与消费者之间:互斥 + 同步
互斥:同一时间对于缓冲区的访问只有一个,当生产者访问时,消费者就不能访问,反之一样。
同步:缓冲区满了,生产者必须等待;缓冲区空了,消费者必须等待,保证按顺序访问。
为了方便理解,我们给生产者消费者模型定义一个321原则:
3:3个关系( 生产者与生产者之间,消费者与消费者之间,生产者与消费者之间)
2:2个角色(生产者和消费者)
1:1个场所(也就是生产者和消费者访问的缓冲区)
生产者消费者模型的优点:
1.解耦,生产者和消费者,不直接交互,生产和消费都与缓冲区交互。
2.支持忙闲不均,可以提前生产缓存,应对高峰期消费;消费慢时也能缓存多余数据,起到削峰填谷的作用。
3.高效?
2.为什么生产者消费者模型是高效的

首先我们应该完整的了解一下整个过程,工厂生产了货物,之后到超市摆放货物(获取数据,将数据写入缓冲区),消费者去超市消费货物,之后自己处理掉货物(从共享缓冲区取数据,加工处理数据)。
但是在多个线程的情况下,虽然说有多个生产者,多个消费者,可是为了保证超市(临界资源的安全),同一时间的情况下,只能由一个生产者或者一个消费者来超市放货物,或者消费货物的啊,那哪里的高效呢?
虽然同一时间只能由一个生产者或者一个消费者来超市放货物,或者消费货物的,但是其他线程也并没有闲着。
- 当其中一个生产者正在往缓冲区放数据时,其他生产者可以并发地去获取数据,互不等待,节省时间;
- 当其中一个消费者正在从缓冲区取数据时,其他消费者可以并发地处理已经拿到的数据,节省时间;
- 当一个生产者在往缓冲区放数据时,多个消费者可以同时在后台处理数据,生产与处理并行;
- 当一个消费者在从缓冲区取数据时,多个生产者可以同时在后台获取数据,获取与消费并行;
所以无论当前是哪个线程在操作缓冲区,其他所有生产者都可以并发获取数据、其他所有消费者都可以并发处理数据,全程大部分耗时操作都在异步并行执行,这也就是生产者消费者模型高效的地方。
3.使用queue来模拟阻塞队列的生产者消费者模型(单消费者单生产者)
1.BlockQueue.hpp
首先我们需要一个类,这个类里面的成员变量有一个queue队列(模拟缓冲区),一个变量(表示缓冲区的最大存放数据),一把锁,用来实现缓冲区的线程安全(只有一个线程可以访问到临界区的资源),两个条件变量,来保证,当缓冲区满了,生产者要停止向缓冲区中存放数据了,当缓冲区为空时,消费者要停止向缓冲区消费数据了。
template<class T>
class BlockQueue
{
private:
queue<T> q;
//缓冲区最大存放数据
int _maxcap;
pthread_mutex_t mutex;
//消费用到的条件变量
pthread_cond_t c_cond;
//生产用到的条件变量
pthread_cond_t p_cond;
};
1.构造函数
对这些变量进行初始化,给_maxcap设置默认值20。
成员变量都是内置类型,所以使用默认生成的析构函数即可。
static const int defaultnum =20;
public:
BlockQueue(int maxcap=defaultnum)
:_maxcap(maxcap)
{
pthread_mutex_init(&mutex,nullptr);
pthread_cond_init(&c_cond,&mutex);
pthread_cond_init(&p_cond,&mutex);
}
接下来需要实现两个接口,来完成生成和消费的动作,同时使用条件变量来达到到缓冲区数据满了停止向缓冲区输送数据,缓冲区为空,就停止消费。
2.push接口(完成生产动作)
那么BlockQueue的成员函数push就是创建节点插入阻塞队列,即入队列,那么由于要保证互斥,所以在访问临界资源queue之前要先申请锁,访问完成之后要释放锁
生产者线程什么时候要阻塞等待呢?当阻塞队列里面的数据满了之后,应该让生产者线程阻塞。
之后就是正常的生产数据,生产数据,代表着阻塞队列中一定有数据了,就可以唤醒消费者线程来消费了,完成了对临界资源的访问,之后解锁。
void push(const T&data)
{
pthread_mutex_lock(&mutex);
if(q.size()==_maxcap)
{
pthread_cond_wait(&p_cond,&mutex);
}
q.push(data);
pthread_cond_signal(&c_cond);
pthread_mutex_unlock(&mutex);
}
这里可能会有一个疑问:判断这一条件能不能存放在加锁之前呢?
肯定是不可以的,q.size()也是访问临界区的资源的,临界的资源就是需要加锁保护的。
3.pop接口(完成消费动作)
那么BlockQueue的成员函数pop就是将节点从阻塞队列弹出,即出队列,那么由于要保证互斥,所以在访问临界资源queue之前要先申请锁,访问完成之后要释放锁
消费者线程什么时候要阻塞等待呢?即当阻塞队列中没有节点的时候不能将节点从队列中出队,当阻塞队列的节点个数为0的时候,应该让消费者线程去消费者条件队列下去进行等待
那么走到接下来一定是if循环判断不成立,即节点的个数不为0,那么接下来可以大胆的出队,所以当出队后,消费者线程可以保证,阻塞队列中至少一定有一个节点空间可以被使用,即可以确保一个生产者线程生产一个节点到阻塞队列,所以消费者线程就去生产者条件变量下唤醒生产者条件变量,最后消费者线程释放锁,同时return返回出队节点即可
T pop()
{
pthread_mutex_lock(&mutex);
if(q.size()==0)
{
pthread_cond_wait(&c_cond,&mutex);
}
T data=q.front();
q.pop();
pthread_cond_signal(&p_cond);
pthread_mutex_unlock(&mutex);
return data;
}
这样子的pop和push接口是存在问题的,因为会有伪唤醒的情况存在。
假如生产的条件变量里面的队列满了,阻塞在条件变量中,消费者消费了一个后,原本操作应该是用 phread_cond_signal 唤醒一个生产来填补上一个空缺,但是 没有用 phread_cond_signal 而是用了 phread_cond_broadcast 唤醒了多个,那么就会导致伪唤醒,当一个线程去填补了这个空缺,这个时候生产已经满了,锁是只有一把的,但是线程有多个,伪唤醒的生产线程就会与消费线程去竞争锁,如果当伪唤醒的生产竞争到了,就会又生产,这个时候生产已经超过了最大值,就会出问题。
所以push和pop的判断条件应该是while,不是if,使用while判断,被伪唤醒的线程,被唤醒了会再次进行一次判断,条件不满足,再次被阻塞住。
2.main.cpp(单消费者单生产者)
主函数:
创建队列和两个线程,然后主线程守在门口等两个线程都运行完(实际上因为是死循环,线程永远不会结束,主线程一直在等待)。
int main()
{
BlockQueue<int>*bq=new BlockQueue<int>();
pthread_t c,p;
pthread_create(&c,nullptr,Consumer,bq);
pthread_create(&p,nullptr,Productor,bq);
pthread_join(c,nullptr);
pthread_join(p,nullptr);
}
消费者函数:不断从队列里拿数据并打印,队列空了就等着。
void*Consumer(void*args)
{
BlockQueue<int>*bq=static_cast<BlockQueue<int>*>(args);
while(true)
{
int data=bq->pop();
cout<<"消费了一个数据: "<<data<<endl;
}
}
生产者函数:
不断从0开始生成数字,每1秒放一个进队列,队列满了就等着。
void*Productor(void*args)
{
int data=0;
BlockQueue<int>*bq=static_cast<BlockQueue<int>*>(args);
while(true)
{
bq->push(data);
cout<<"生成了一个数据: "<<data<<endl;
data++;
sleep(1);
}
}
4.改进生产者消费者模型
将消费者生成者改为多线程版,同时对数据进行改进,不再是单单打印。
创建一个任务类
1.构造函数
我们要实现的任务是计算加,减,乘,除,取模运算,那么就要有对应的运算符,所以我们定义一个全局的opers用于存储运算符
既然是运算,那么就有可能出现除零错误,取模零错误,以及运算符错误,所以我们使用枚举定义对应的错误码\n接下来看类Task的私有成员变量,那么首先应该有左操作数data1_,右操作数data2_,操作符oper_,计算结果result_,错误码exitcode_
所以我们就在Task的构造函数中对这些成员变量进行初始化即可,这里的操作符需要传入,并且计算结果以及错误码默认为0
Task(int x,int y,char op)
:data1_(x)
,oper_(op)
,data2_(y)
,result_(0)
,exitcode_()
{}
2.run()
既然是封装的任务,那么就应该有任务的计算方法,这里为了便于调用我们使用仿函数,所以首先利用switch case语句对操作符oper_进行判断,然后计算即可
但是这里可能会出现除零错误,所以我们判断一下当运算符是除号的时候,如果第二个操作数data_为0,那么就会发生除零错误,这里我们就不进行运算了,直接设置错误码即可,用户看到操作码之后,就会知道出现错误了,结果不可靠
同样的也有可能会出现取模零错误,同样的方法和除零错误一样,如果第二个操作数data_为0,那么就会发生取模零错误,这里我们就不进行运算了,直接设置错误码即可
但是有可能用户传入的运算符是其它值,所以这里如果运算符不是加,减,乘,除,取模,那么就会进入default,那么同样的我们设置错误码即可
void run()
{
switch(oper_)
{
case '+':
result_=data1_+data2_;
break;
case '-':
result_=data1_-data2_;
break;
case '*':
result_=data1_*data2_;
break;
case '/':
{
if(data2_==0) exitcode_=ModZero;
result_=data1_/data2_;
}
break;
case '%':
{
if(data2_==0) exitcode_=DivZero;
else result_=data1_%data2_;
}
break;
default:
exitcode_=Unknown;
break;
}
}
void operator()()
{
run();
}
3.GetTask
那么生产者生产了什么样的任务到阻塞队列,还应该打印给用户看,所以Task类还应该提供一个获取任务,即将任务组成字符串进行返回即可
string GetTask()
{
string r=to_string(data1_);
r+=oper_;
r+=to_string(data2_);
r+="=?";
return r;
}
4.GetResult
那么生产者生产了什么样的任务到阻塞队列,还应该打印给用户看,所以Task类还应该提供一个获取任务,即将任务组成字符串进行返回即可
string GetResult()
{
string r=to_string(data1_);
r+=oper_;
r+=to_string(data2_);
r+="=";
r+=to_string(result_);
r+="[code";
r+=to_string(exitcode_);
r+="]";
return r;
}
对main函数进行改变
创建多线程,同时传入的数据变为Task类型
#include "BlockQueue.hpp"
#include "task.hpp"
void*Consumer(void*args)
{
BlockQueue<Task>*bq=static_cast<BlockQueue<Task>*>(args);
while(true)
{
Task t=bq->pop();
t();
cout<<"处理了一个任务: "<<t.GetResult()<<endl;
sleep(1);
}
}
void*Productor(void*args)
{
BlockQueue<Task>*bq=static_cast<BlockQueue<Task>*>(args);
while(true)
{
int data1=rand()%10;//[0,9]
usleep(10);
int data2=rand()%10;
char op=opers[rand()%opers.size()];
Task t(data1,data2,op);
cout<<"生产了一个任务: "<<t.GetTask()<<endl;
bq->push(t);
// sleep(1);
}
}
int main()
{
srand(time(nullptr));
BlockQueue<Task>*bq=new BlockQueue<Task>();
pthread_t c[5],p[5];
for(int i=0;i<5;i++)
{
pthread_create(c+i,nullptr,Consumer,bq);
}
for(int i=0;i<5;i++)
{
pthread_create(p+i,nullptr,Productor,bq);
}
for(int i=0;i<5;i++)
{
pthread_join(c[i],nullptr);
}
for(int i=0;i<5;i++)
{
pthread_join(p[i],nullptr);
}
delete bq;
}
5.源码展示
BlockQueue.hpp
#include <mutex>
#include <unistd.h>
#include <iostream>
#include <pthread.h>
#include <queue>
using namespace std;
template<class T>
class BlockQueue
{
static const int defaultnum =20;
public:
BlockQueue(int maxcap=defaultnum)
:_maxcap(maxcap)
{
pthread_mutex_init(&mutex,nullptr);
pthread_cond_init(&c_cond,nullptr);
pthread_cond_init(&p_cond,nullptr);
}
void push(const T&data)
{
pthread_mutex_lock(&mutex);
while(q.size()==_maxcap)
{
pthread_cond_wait(&p_cond,&mutex);
}
q.push(data);
pthread_cond_signal(&c_cond);
pthread_mutex_unlock(&mutex);
}
T pop()
{
pthread_mutex_lock(&mutex);
while(q.size()==0)
{
pthread_cond_wait(&c_cond,&mutex);
}
T data=q.front();
q.pop();
pthread_cond_signal(&p_cond);
pthread_mutex_unlock(&mutex);
return data;
}
private:
queue<T> q;
//缓冲区最大存放数据
int _maxcap;
pthread_mutex_t mutex;
//消费用到的条件变量
pthread_cond_t c_cond;
//生产用到的条件变量
pthread_cond_t p_cond;
};
main.cpp
#include "BlockQueue.hpp"
#include "task.hpp"
void*Consumer(void*args)
{
BlockQueue<Task>*bq=static_cast<BlockQueue<Task>*>(args);
while(true)
{
Task t=bq->pop();
t();
cout<<"处理了一个任务: "<<t.GetResult()<<endl;
sleep(1);
}
}
void*Productor(void*args)
{
BlockQueue<Task>*bq=static_cast<BlockQueue<Task>*>(args);
while(true)
{
int data1=rand()%10;//[0,9]
usleep(10);
int data2=rand()%10;
char op=opers[rand()%opers.size()];
Task t(data1,data2,op);
cout<<"生产了一个任务: "<<t.GetTask()<<endl;
bq->push(t);
// sleep(1);
}
}
int main()
{
srand(time(nullptr));
BlockQueue<Task>*bq=new BlockQueue<Task>();
pthread_t c[5],p[5];
for(int i=0;i<5;i++)
{
pthread_create(c+i,nullptr,Consumer,bq);
}
for(int i=0;i<5;i++)
{
pthread_create(p+i,nullptr,Productor,bq);
}
for(int i=0;i<5;i++)
{
pthread_join(c[i],nullptr);
}
for(int i=0;i<5;i++)
{
pthread_join(p[i],nullptr);
}
delete bq;
}
task.hpp
#pragma once
#include<iostream>
#include<string>
using namespace std;
string opers="+-*/%";
enum{
DivZero=1,
ModZero,
Unknown
};
class Task
{
public:
Task(int x,int y,char op)
:data1_(x)
,oper_(op)
,data2_(y)
,result_(0)
,exitcode_()
{}
void run()
{
switch(oper_)
{
case '+':
result_=data1_+data2_;
break;
case '-':
result_=data1_-data2_;
break;
case '*':
result_=data1_*data2_;
break;
case '/':
{
if(data2_==0) exitcode_=ModZero;
result_=data1_/data2_;
}
break;
case '%':
{
if(data2_==0) exitcode_=DivZero;
else result_=data1_%data2_;
}
break;
default:
exitcode_=Unknown;
break;
}
}
void operator()()
{
run();
}
string GetResult()
{
string r=to_string(data1_);
r+=oper_;
r+=to_string(data2_);
r+="=";
r+=to_string(result_);
r+="[code";
r+=to_string(exitcode_);
r+="]";
return r;
}
string GetTask()
{
string r=to_string(data1_);
r+=oper_;
r+=to_string(data2_);
r+="=?";
return r;
}
~Task()
{}
private:
int data1_;
int data2_;
char oper_;
int result_;
int exitcode_;
};
更多推荐




所有评论(0)