Linux之生產(chǎn)者消費(fèi)者模型用法解讀
一、生產(chǎn)者消費(fèi)者模型
舉個(gè)例子:

我們現(xiàn)實(shí)生活中:工廠,超市,人[消費(fèi)者]之間的關(guān)系就是一個(gè)典型的生產(chǎn)者消費(fèi)者模型
這個(gè)模型的基本工作流程是:
- 工廠制造商品供貨給超市,消費(fèi)者到超市購(gòu)買(mǎi)商品
現(xiàn)實(shí)生活中,為什么會(huì)優(yōu)化出這樣的一種模型?
為什么消費(fèi)者不直接去工廠買(mǎi)東西?
主要有以下3個(gè)原因:
為了效率
1.對(duì)于工廠來(lái)說(shuō),消費(fèi)者如果直接來(lái)買(mǎi)商品,一個(gè)消費(fèi)者一次購(gòu)買(mǎi)的商品數(shù)量非常有限,而且一個(gè)工廠的產(chǎn)品種類一般并不多,所以消費(fèi)者還得去不同工廠買(mǎi)東西,工廠為了防止賣(mài)不完,只能減少商品的生產(chǎn),即降低了工廠的生產(chǎn)效率,但是如果賣(mài)給超市,就可以一次性賣(mài)一大車(chē)貨物,多找?guī)讉€(gè)超市,就可以做到:生產(chǎn)多少賣(mài)多少了
2.對(duì)于消費(fèi)者,工廠的占地面積比較大,所以一般都距離居民地比較遠(yuǎn),而超市一般都比工廠占地面積小的多,離消費(fèi)者很近,所以消費(fèi)者去超市買(mǎi)東西比去工廠更快,而且一個(gè)工廠的產(chǎn)品種類一般并不多,但是超市售賣(mài)的產(chǎn)品種類繁多
有了超市就可以做到生產(chǎn)者和消費(fèi)者之間的解藕
1.一個(gè)工廠關(guān)門(mén)了,超市換一個(gè)工廠進(jìn)貨就行了,只要進(jìn)的貨還是一樣的,就絲毫不影響消費(fèi)者消費(fèi)。事實(shí)就是如此,我們消費(fèi)者從超市買(mǎi)東西時(shí),根本就不知道這個(gè)商品是哪個(gè)工廠生產(chǎn)的,我們也不關(guān)心消費(fèi)者只管購(gòu)買(mǎi)商品就行了
2.到超市消費(fèi)者是誰(shuí),是什么群體,工廠也不知道,也不關(guān)心,這個(gè)是超市該關(guān)心的,因?yàn)楣S只把商品賣(mài)給超市,超市賣(mài)給誰(shuí)和工廠無(wú)關(guān),工廠只管生產(chǎn)就行了
3.有了超市的存在
- 左側(cè)的工廠發(fā)生變化不影響右側(cè)的消費(fèi)者
- 右側(cè)的消費(fèi)者發(fā)生變化不影響左側(cè)的工廠
- 這不就
解藕了嗎?
支持忙閑不均
有了超市這個(gè)緩沖區(qū)的存在
1.在購(gòu)物潮之前,超市就可以通知工廠多生產(chǎn)一些商品,超市多進(jìn)貨,方便應(yīng)對(duì)更多消費(fèi)者
2.在囤積的商品多的時(shí)候,超市也可以搞活動(dòng)吸引消費(fèi)者,并且通知工廠生產(chǎn)地慢一點(diǎn)
超市本質(zhì)上就是工廠和消費(fèi)者之間的緩存
即:超市支持工廠“預(yù)加載”產(chǎn)品,消費(fèi)者可以一定程度“預(yù)定”產(chǎn)品
上述例子對(duì)應(yīng)到計(jì)算機(jī)中的生產(chǎn)者消費(fèi)者模型就是:
- ①工廠:生產(chǎn)者線程
- ②消費(fèi)者:消費(fèi)者線程
- ③超市:以某種數(shù)據(jù)結(jié)構(gòu)組織的內(nèi)存區(qū)域
- ④商品:數(shù)據(jù)
和線程安全再對(duì)應(yīng)一下:
- ①超市:
共享/臨界資源 - ②我們要研究生產(chǎn)者消費(fèi)者模型,就需要研究清楚多個(gè)生產(chǎn)者和多個(gè)消費(fèi)者之間的同步互斥關(guān)系!
一共有3種關(guān)系:
1.生產(chǎn)者和生產(chǎn)者之間:互斥
- 因?yàn)槿绻鄠€(gè)生產(chǎn)者同時(shí)向超市(共享資源)寫(xiě)數(shù)據(jù),可能會(huì)產(chǎn)生線程安全問(wèn)題
- 因?yàn)橥瑫r(shí)有多個(gè)線程對(duì)同一份共享資源進(jìn)行修改,很容易互相覆蓋,出現(xiàn)并發(fā)線程安全問(wèn)題 即:線程之間切換時(shí)可能導(dǎo)致數(shù)據(jù)不一致問(wèn)題
- 其次共享資源的空間就那么大,多個(gè)生產(chǎn)者線程同時(shí)去寫(xiě)的話,一個(gè)線程寫(xiě)的多了,另一個(gè)線程就只能少寫(xiě)一點(diǎn)了
- 或者線程a剛向共享資源中的一個(gè)位置寫(xiě)了一個(gè)1,線程b就跑過(guò)來(lái)把那個(gè)位置的1改成10了
- 所以生產(chǎn)者線程之間是競(jìng)爭(zhēng)關(guān)系
2.消費(fèi)者和消費(fèi)者之間:互斥
- 消費(fèi)者線程,雖然是進(jìn)入臨界資源讀取數(shù)據(jù),但是其實(shí)也會(huì)對(duì)共享資源進(jìn)行修改
- 即:線程拿走了一個(gè)數(shù)據(jù)a,其他線程看待這個(gè)這個(gè)數(shù)據(jù)a就是過(guò)期數(shù)據(jù)了(和到超市買(mǎi)東西一樣,買(mǎi)了一包方便面,超市就少了一包方便面)
- 既然會(huì)修改,那么多個(gè)消費(fèi)者線程同時(shí)進(jìn)入共享資源進(jìn)行修改的話,也可能會(huì)出現(xiàn)并發(fā)切換導(dǎo)致的線程安全問(wèn)題
- 比如:線程a和線程b同時(shí)訪問(wèn)共享資源,線程a先把數(shù)據(jù)X消費(fèi)走了,但是因?yàn)槭峭瑫r(shí)進(jìn)入線程b任認(rèn)為數(shù)據(jù)X還是有效的,也拿了數(shù)據(jù)X(如果共享資源是數(shù)組,那可以把數(shù)據(jù)X理解為數(shù)組中的一個(gè)元素)
- 就會(huì)導(dǎo)致一份數(shù)據(jù),被使用了兩次
3.生產(chǎn)者和消費(fèi)者之間:互斥+同步
- 因?yàn)椴还苁巧a(chǎn)者線程還是消費(fèi)者線程,訪問(wèn)共享資源時(shí),都會(huì)進(jìn)行修改,多線程并發(fā)修改共享資源,就可能會(huì)出現(xiàn)并發(fā)線程安全問(wèn)題
- 所以它們首先得是
互斥的 - 但是如果只有互斥,那么消費(fèi)者如果不知道有沒(méi)有數(shù)據(jù),就只能不斷地去輪詢檢測(cè),每次輪詢都要申請(qǐng)鎖,那生產(chǎn)者線程就可能很難搶到鎖,就一直生產(chǎn)不了數(shù)據(jù),消費(fèi)者線程就一直拿不到數(shù)據(jù),造成惡性循環(huán)
- 就可能會(huì)導(dǎo)致鎖的饑餓問(wèn)題
- 所以需要
同步關(guān)系來(lái)提高生產(chǎn)者和消費(fèi)者模型的效率 即: - 設(shè)置條件變量,讓消費(fèi)者線程在條件變量的等待隊(duì)列中等,當(dāng)生產(chǎn)者線程生產(chǎn)數(shù)據(jù)之后,才喚醒消費(fèi)者線程,去讀取數(shù)據(jù)
二、生產(chǎn)者消費(fèi)者模型的阻塞隊(duì)列版本
生產(chǎn)者消費(fèi)者模型一般會(huì)使用一個(gè)阻塞隊(duì)列來(lái)作為共享資源,進(jìn)而實(shí)現(xiàn)多線程協(xié)作
生產(chǎn)者消費(fèi)者模型的阻塞隊(duì)列的特點(diǎn):
- 如果隊(duì)列為空,那么一個(gè)消費(fèi)者線程如果來(lái)拿數(shù)據(jù),它就會(huì)被阻塞
- 如果隊(duì)列為滿,那么一個(gè)生產(chǎn)者線程如果還要向隊(duì)列里寫(xiě)數(shù)據(jù),它就會(huì)被阻塞
- 如果隊(duì)列不空也不滿,那么生產(chǎn)者就可以向隊(duì)列尾部寫(xiě)數(shù)據(jù),消費(fèi)者就可以向隊(duì)列頭部拿數(shù)據(jù)
阻塞隊(duì)列類的簡(jiǎn)單實(shí)現(xiàn)
成員變量
1.存儲(chǔ)數(shù)據(jù)的容器
直接使用STL的queue,因?yàn)閿?shù)據(jù)的類型不確定,所以阻塞隊(duì)列類是模板類
2.一把鎖
阻塞隊(duì)列自己會(huì)被所有線程看見(jiàn),所以它是共享資源,所以需要鎖來(lái)保護(hù)自己
要幾把鎖呢?
因?yàn)樗猩a(chǎn)者線程之間,所有消費(fèi)者線程之間,以及生產(chǎn)者和消費(fèi)者之間
都是互斥的,所以它們得用同一把鎖來(lái)實(shí)現(xiàn)互斥
3.生產(chǎn)者線程的條件變量
因?yàn)樵跐M足一定條件(比如:阻塞隊(duì)列滿了,或者阻塞隊(duì)列滿了4/5了等)時(shí),可以讓所有生產(chǎn)者線程暫時(shí)暫停生產(chǎn)(即去條件變量的等待隊(duì)列中阻塞)
4.消費(fèi)者線程的條件變量
因?yàn)樵跐M足一定條件(比如:阻塞隊(duì)列為空,或者阻塞隊(duì)列空了4/5了等)時(shí),可以讓所有消費(fèi)者線程暫時(shí)暫停消費(fèi)(即去條件變量的等待隊(duì)列中阻塞)
為什么要搞兩個(gè)條件變量?
一個(gè)條件變量雖然也可以實(shí)現(xiàn)生產(chǎn)者線程和消費(fèi)者線程之間的同步
但是實(shí)現(xiàn)起來(lái)非常麻煩,而且不能區(qū)分條件變量的等待隊(duì)列下的是生產(chǎn)者線程還是消費(fèi)者線程
而且
兩個(gè)條件變量可以很好地支持:
生產(chǎn)者消費(fèi)者模型的第3個(gè)優(yōu)點(diǎn):忙閑不均
5.int _cap:阻塞隊(duì)列的最大容量
6.int _csleep_num:在消費(fèi)者條件變量的等待隊(duì)列中等待的線程個(gè)數(shù)
7.int _psleep_num:在生產(chǎn)者條件變量的等待隊(duì)列中等待的線程個(gè)數(shù)
6和7成員變量的存在主要是為了方便實(shí)現(xiàn)線程之間的互相喚醒機(jī)制
(即生產(chǎn)者線程生產(chǎn)了之后,可以喚醒消費(fèi)者線程來(lái)消費(fèi),反之同理)
成員函數(shù)
- Equeue:生產(chǎn)數(shù)據(jù)
void Equeue(const T& in)
{
pthread_mutex_lock(&_mutex);
//生產(chǎn)者調(diào)用
while(IsFull())
{
_psleep_num++;
cout << "生產(chǎn)者, 進(jìn)入休眠了:" << _psleep_num << endl;
pthread_cond_wait(&_full_cond, &_mutex);
_psleep_num--;
}
//100% 隊(duì)列有空間
_q.push(in);
if(_csleep_num > 0)
{
pthread_cond_signal(&_empty_cond);
cout << "喚醒消費(fèi)者..." << endl;
}
pthread_mutex_unlock(&_mutex);
}
代碼細(xì)節(jié):
偽喚醒問(wèn)題的解決
- 即:判斷線程是否要進(jìn)入條件變量的等待隊(duì)列時(shí),判斷不能用if而要用while
- 不然就有可能出現(xiàn)偽喚醒問(wèn)題:即在條件變量下等待的線程,喚醒條件其實(shí)并不滿足
- 但是因?yàn)槌绦騿T編碼的問(wèn)題,可能意外被喚醒了
例如:
- 生產(chǎn)者消費(fèi)者模型中,因?yàn)樽枞?duì)列中沒(méi)有數(shù)據(jù),所以全部都5個(gè)消費(fèi)者線程在條件變量的等待隊(duì)列中等待
- 生產(chǎn)者線程生產(chǎn)了一個(gè)數(shù)據(jù),意外地把喚醒了多個(gè)消費(fèi)者線程
- 然后一個(gè)消費(fèi)者線程搶到鎖之后,把阻塞隊(duì)列中那唯一的一個(gè)數(shù)據(jù)搶走了,它解鎖之后
- 因?yàn)閱拘蚜硕鄠€(gè)消費(fèi)者線程
- 所以鎖可能又被一個(gè)消費(fèi)者線程搶到了,但是此時(shí)阻塞隊(duì)列中根本沒(méi)有數(shù)據(jù)!
此時(shí):
1.如果此時(shí)是使用if進(jìn)行“線程是否需要進(jìn)入條件變量的等待隊(duì)列"的判斷的這個(gè)被偽喚醒的線程,重新申請(qǐng)并拿到鎖之后,就直接"餓虎出籠"去肆意妄為了
2.如果是使用while進(jìn)行“線程是否需要進(jìn)入條件變量的等待隊(duì)列”的判斷的,這個(gè)被喚醒的線程,重新申請(qǐng)并拿到鎖之后,也還是不能直接出循環(huán),因?yàn)橐倥袛嘁幌卵h(huán)條件是否不滿足了
雖然循環(huán)條件是"線程需要進(jìn)入等待隊(duì)列"的條件,但是如果這個(gè)條件滿足,不就意味著線程不應(yīng)該被喚醒嗎?
- Pop:獲取并刪除數(shù)據(jù)
T Pop()
{
//消費(fèi)者調(diào)用
pthread_mutex_lock(&_mutex);
while(IsEmpty())
{
_csleep_num++;
pthread_cond_wait(&_empty_cond, &_mutex);
_csleep_num--;
}
T data = _q.front();
_q.pop();
if(_psleep_num > 0)
{
pthread_cond_signal(&_full_cond);
cout << "喚醒生產(chǎn)者.." << endl;
}
pthread_mutex_unlock(&_mutex);
return data;
}
- IsEmpty:阻塞隊(duì)列是否為空
- IsFull:阻塞隊(duì)列是否為滿
bool IsFull()
{
return _q.size() >= _cap;
}
bool IsEmpty()
{
return _q.empty();
}
源碼
#pragma once
#include <iostream>
#include <pthread.h>
#include <queue>
#include <string>
using namespace std;
int defalutcap = 5;
template<typename T>
class BlockQueue
{
private:
bool IsFull()
{
return _q.size() >= _cap;
}
bool IsEmpty()
{
return _q.empty();
}
public:
BlockQueue(int cap = defalutcap)
: _cap(cap),
_csleep_num(0),
_psleep_num(0)
{
pthread_mutex_init(&_mutex, nullptr);
pthread_cond_init(&_full_cond, nullptr);
pthread_cond_init(&_empty_cond, nullptr);
}
void Equeue(const T& in)
{
pthread_mutex_lock(&_mutex);
//生產(chǎn)者調(diào)用
while(IsFull())
{
_psleep_num++;
cout << "生產(chǎn)者, 進(jìn)入休眠了:" << _psleep_num << endl;
pthread_cond_wait(&_full_cond, &_mutex);
_psleep_num--;
}
//100% 隊(duì)列有空間
_q.push(in);
if(_csleep_num > 0)
{
pthread_cond_signal(&_empty_cond);
cout << "喚醒消費(fèi)者..." << endl;
}
pthread_mutex_unlock(&_mutex);
}
T Pop()
{
//消費(fèi)者調(diào)用
pthread_mutex_lock(&_mutex);
while(IsEmpty())
{
_csleep_num++;
pthread_cond_wait(&_empty_cond, &_mutex);
_csleep_num--;
}
T data = _q.front();
_q.pop();
if(_psleep_num > 0)
{
pthread_cond_signal(&_full_cond);
cout << "喚醒生產(chǎn)者.." << endl;
}
pthread_mutex_unlock(&_mutex);
return data;
}
~BlockQueue()
{
pthread_mutex_destroy(&_mutex);
pthread_cond_destroy(&_full_cond);
pthread_cond_destroy(&_empty_cond);
}
private:
//臨界資源
queue<T> _q;
//大小
int _cap;
pthread_mutex_t _mutex;
pthread_cond_t _full_cond;
pthread_cond_t _empty_cond;
int _csleep_num;//消費(fèi)者休眠的個(gè)數(shù)
int _psleep_num;//生產(chǎn)者休眠的個(gè)數(shù)
};
三、POSIX信號(hào)量
POSIX信號(hào)量和SystemV信號(hào)量作?相同,都是?于同步操作,達(dá)到?沖突的訪問(wèn)共享資源?的。但POSIX可以?于線程間同步。
初始化信號(hào)量
sem_init
作用: 用于初始化一個(gè)未命名的POSIX信號(hào)量(也稱為匿名信號(hào)量),通常用于線程間同步或共享內(nèi)存的進(jìn)程間同步。
#include <semaphore.h> int sem_init(sem_t *sem, int pshared, unsigned int value);
- sem_t *sem: 指向要初始化的信號(hào)量對(duì)象的指針。
- pshared: 0表?線程間共享,?零表?進(jìn)程間共享
- value:信號(hào)量初始值
返回值
- 成功:返回 0
- 失?。悍橇?/li>
銷毀信號(hào)量
sem_destroy
int sem_destroy(sem_t *sem);
等待信號(hào)量
sem_wait
是 POSIX 信號(hào)量的 P操作(等待/獲取信號(hào)量),用于對(duì)信號(hào)量進(jìn)行原子減1操作。
它的主要作用是:
- 如果信號(hào)量值 > 0:立即將其減 1,線程繼續(xù)執(zhí)行。
- 如果信號(hào)量值 = 0:線程阻塞,直到信號(hào)量值變?yōu)檎龜?shù)(其他線程或進(jìn)程調(diào)用 sem_post 釋放資源)。
int sem_wait(sem_t *sem); //P()
發(fā)布信號(hào)量
sem_post
是 POSIX 信號(hào)量的 V操作(釋放/增加信號(hào)量),用于對(duì)信號(hào)量進(jìn)行原子加1操作。它的主要作用是:
- 將信號(hào)量的值 +1,表示釋放一個(gè)資源。
- 如果有線程阻塞在 sem_wait,則喚醒其中一個(gè)線程(取決于系統(tǒng)調(diào)度策略)。
int sem_post(sem_t *sem);//V()
四、基于環(huán)形隊(duì)列的?產(chǎn)消費(fèi)模型
上?節(jié)?產(chǎn)者-消費(fèi)者的例?是基于queue的,其空間可以動(dòng)態(tài)分配,現(xiàn)在基于固定??的環(huán)形隊(duì)列重寫(xiě)這個(gè)程序(POSIX信號(hào)量):
Sem的封裝
#include <iostream>
#include <semaphore.h>
#include <pthread.h>
namespace SemMoudle
{
const int defaultvalue = 1;
class Sem
{
public:
Sem(unsigned int sem_vlaue = defaultvalue)
{
sem_init(&_sem, 0, sem_vlaue);
}
void P()
{
//等待信號(hào)量,會(huì)將信號(hào)量的值減1
int n = sem_wait(&_sem);//原子的
(void)n;
}
void V()
{
//發(fā)布信號(hào)量,釋放資源,會(huì)將信號(hào)量的值加1
int n = sem_post(&_sem);//原子的
(void)n;
}
~Sem()
{
sem_destroy(&_sem);
}
private:
sem_t _sem;
};
}
Mutex的封裝
#pragma once
#include <iostream>
#include <pthread.h>
namespace MutexModue
{
class Mutex
{
public:
Mutex()
{
pthread_mutex_init(&_mutex, nullptr);
}
void Lock()
{
int n = pthread_mutex_lock(&_mutex);
(void)n;
}
void Unlock()
{
int n = pthread_mutex_unlock(&_mutex);
(void)n;
}
~Mutex()
{
pthread_mutex_destroy(&_mutex);
}
pthread_mutex_t* Get()
{
return &_mutex;
}
private:
pthread_mutex_t _mutex;
};
class LockGuard
{
public:
LockGuard(Mutex& mutex):_mutex(mutex)
{
_mutex.Lock();
}
~LockGuard()
{
_mutex.Unlock();
}
private:
Mutex& _mutex;
};
}
環(huán)形隊(duì)列
#pragma once
#include <iostream>
#include <vector>
#include "Sem.hpp"
#include "Mutex.hpp"
using namespace std;
static const int gcap = 5;
using namespace SemMoudle;
using namespace MutexModue;
//環(huán)形隊(duì)列
template<typename T>
class RingQueue
{
public:
RingQueue(int cap = gcap)
:_cap(cap),
_rq(cap),
_blank_sem(cap),
_p_step(0),
_data_sem(0),
_c_step(0)
{}
void Equeue(const T& in)
{
//生產(chǎn)者
//1.申請(qǐng)信號(hào)量,空位置信號(hào)量
_blank_sem.P();
{
LockGuard lockguard(_pmutex);
//2.生產(chǎn)
_rq[_p_step] = in;
//3.更新下標(biāo)
++_p_step;
//4.維持環(huán)形特性
_p_step %= _cap;
}
_data_sem.V();
}
void Pop(T* out)
{
//消費(fèi)者
//1.申請(qǐng)信號(hào)量,數(shù)據(jù)信號(hào)量
_data_sem.P();
{
LockGuard lockguard(_cmutex);
//2.消費(fèi)
*out = _rq[_c_step];
//3.更新下標(biāo)
++_c_step;
//4.維持環(huán)形特性
_c_step %= _cap;
}
_blank_sem.V();
}
private:
vector<T> _rq;
int _cap;
//生產(chǎn)者
Sem _blank_sem;//空位置
int _p_step;
//消費(fèi)者
Sem _data_sem;//數(shù)據(jù)
int _c_step;
//維護(hù)多生產(chǎn),多消費(fèi),2把鎖
Mutex _cmutex;
Mutex _pmutex;
};
總結(jié)
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
相關(guān)文章
Linux中創(chuàng)建,復(fù)制和刪除文件及目錄的命令詳解
這篇文章主要為大家詳細(xì)介紹了Linux中創(chuàng)建,復(fù)制和刪除文件及目錄的相關(guān)命令,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以跟隨小編一起學(xué)習(xí)一下2026-02-02
Linux?Docker安裝Jenkins并實(shí)現(xiàn)Maven工程自動(dòng)化部署過(guò)程
這段文章詳細(xì)介紹了在Linux系統(tǒng)上安裝和配置Jenkins的過(guò)程,包括關(guān)閉防火墻、訪問(wèn)Jenkins界面、解鎖Jenkins容器、安裝Maven插件、配置Jenkins連接碼云倉(cāng)庫(kù)以及設(shè)置構(gòu)建后操作等步驟2026-06-06
CentOS 7中 Minimal 安裝JDK 1.8的教程
這篇文章主要介紹了CentOS 7 Minimal 安裝JDK 1.8的教程,非常不錯(cuò),具有參考借鑒價(jià)值 ,需要的朋友可以參考下2018-05-05
使用Linux的read和write系統(tǒng)函數(shù)操作文件的方法詳解
在Linux系統(tǒng)編程中,文件操作是非?;A(chǔ)且重要的部分,Linux提供了多個(gè)系統(tǒng)調(diào)用來(lái)實(shí)現(xiàn)文件的讀寫(xiě)操作,其中read和write是最常用的兩個(gè)函數(shù),本文將詳細(xì)介紹這兩個(gè)系統(tǒng)調(diào)用的功能、使用方法以及實(shí)際應(yīng)用中的注意事項(xiàng),需要的朋友可以參考下2025-10-10
你需要知道的16個(gè)Linux服務(wù)器監(jiān)控命令
如果你想知道你的服務(wù)器正在做干什么,你就需要了解一些基本的命令,一旦你精通了這些命令,那你就是一個(gè) 專業(yè)的 Linux 系統(tǒng)管理員2012-03-03

