文章目录
- 概要&序論
- 一、基于环形缓冲的生产者消费者模型原理介绍
- 1.1 环形缓冲区与核心逻辑约定
- 1.2 并发访问条件与同步互斥机制
- 1.2.1 什么时候可以并发?
- 1.2.2 什么时候需要同步与互斥?
- 1.3 基于 POSIX 信号量(P/V 操作)的原理实现
- 1.3.1 信号量的资源抽象
- 1.3.2 生产者与消费者的执行逻辑
- 1.4 特殊场景:二元信号量与单缓冲区(N = 1)
- 二、POSIX信号量接口的使用
- 2.1 初始化信号量
- 2.2 销毁信号量
- 2.3 等待信号量(P操作)
- 2.4 发布信号量(V操作)
- 三、基于环形缓冲的生产者消费者模型源码解析
- 3.1RingQueue.hpp
- 3.2Sem.hpp
- 3.3Task.hpp
- 3.4Mutex.hpp
- 3.5Main.cc
- 四、以上代码中出现的相关问题
- 4.1 加锁与申请信号量的先后顺序问题
- 4.1.1 正确性保证
- 4.1.2 高效性对比(买电影票比喻)
概要&序論
Hello大家好,我是此方。本文围绕环形缓冲区实现生产者消费者模型展开,首先介绍环形缓冲区的核心逻辑与生产者、消费者之间的并发访问关系,随后深入讲解 POSIX 信号量及 P/V 操作,分析并发控制、同步互斥的实现原理。在此基础上,完成模型的代码实现,并结合特殊场景进一步理解信号量在多线程协作中的实际应用。
部分内容在前面已经讲解过了,这里不再赘述
生产消费模型概述见上一篇:基于阻塞队列的生产者消费者模型与条件变量
信号量概述见:信号量与并发编程初步认识
环形缓冲详解见:手把手教你实现环形缓冲
一、基于环形缓冲的生产者消费者模型原理介绍
1.1 环形缓冲区与核心逻辑约定
在并发编程中,基于环形缓冲(Ring Buffer)的生产者消费者模型是一种极其高效的数据同步结构。与基于单队列加锁的模式不同,环形缓冲区借助于固定大小的数组(容量为N)和POSIX 信号量,能够在满足特定条件时实现生产与消费的并发执行。
在 POSIX 标准下,信号量相比传统的 System V 信号量更加轻量化与高效。为了保证环形缓冲区在多线程环境下的数据安全与逻辑正确,我们需要建立以下四个核心约束约定:
- 约定 1:当缓冲区为空时,生产者必须先运行。
- 约定 2:当缓冲区为满时,消费者必须先运行。
- 约定 3:生产者不能将消费者“套圈”(超越一圈以上)。其本质是为了防止生产者覆盖上一轮未被消费的数据。
- 约定 4:消费者不能超过生产者。其本质是为了防止消费者读取到未生产的无效数据。
我们可以将环形缓冲区形象地比作一张大圆桌,每一个可以放置数据的槽位就是桌上的盘子(空格子)。生产者向盘子里放数据,消费者从盘子里取数据。
1.2 并发访问条件与同步互斥机制
理解环形缓冲区的精髓,在于厘清生产者与消费者什么时候可以并发运行,什么时候必须进行同步与互斥。
1.2.1 什么时候可以并发?
只要生产者与消费者不访问同一个槽位,两者就可以同时进行。
当环形队列满足不为空 && 不为满的条件时,生产者指针(tail / p_step)与消费者指针(head / c_step)指向不同的格子。此时生产与消费互不干扰,能够实现真正的并发处理,大幅提升系统吞吐量。
1.2.2 什么时候需要同步与互斥?
当生产者与消费者指向同一个槽位时,线程之间会发生资源竞争,此时必须引入互斥与同步机制:
- 缓冲区为空时:两个指针重合。由于没有可消费的数据,必须执行互斥与同步,强制生产者先运行,填充数据后唤醒消费者。
- 缓冲区为满时:两个指针再次重合(生产者套圈)。由于没有多余的空格子,必须执行互斥与同步,强制消费者先运行,释放空间后唤醒生产者。
1.3 基于 POSIX 信号量(P/V 操作)的原理实现
为了在代码层面上完美落实上述四个约定,我们需要借助POSIX 信号量提供的原子P/V 操作来管理计数资源。
1.3.1 信号量的资源抽象
我们将环形缓冲区的资源划分为两类:
- 生产者关注的资源:空槽位数量,定义为信号量sem_blank,初始值为N。
- 消费者关注的资源:数据槽位数量,定义为信号量sem_data,初始值为0。
信号量的P 操作(申请资源)具有原子性:若资源计数大于 0 则递减并成功继续;若资源计数为 0,申请线程将被阻塞挂起。信号量的V 操作(释放资源)则会递增计数并唤醒等待线程。
1.3.2 生产者与消费者的执行逻辑
生产者逻辑(Producer Process):
- 申请空位资源:P(sem_blank)(即sem_blank–)
- 在当前下标p_step位置写入数据。
- 更新索引指针:p_step = (p_step + 1) % N。
- 释放数据资源:V(sem_data)(即sem_data++,唤醒可能阻塞的消费者)。
消费者逻辑(Consumer Process):
- 申请数据资源:P(sem_data)(即sem_data–)
- 在当前下标c_step位置读取并消费数据。
- 更新索引指针:c_step = (c_step + 1) % N。
- 释放空位资源:V(sem_blank)(即sem_blank++,唤醒可能阻塞的生产者)。
生产者通过V(sem_data)激活消费者的P(sem_data),消费者通过V(sem_blank)激活生产者的P(sem_blank),二者形成精密的交替唤醒闭环。
1.4 特殊场景:二元信号量与单缓冲区(N = 1)
当环形缓冲区的容量N = 1时,环形队列退化为只有一个格子的单缓冲区。
此时,sem_blank初始为 1,sem_data初始为 0。系统变成了一种全新的同步互斥实现方式——通过这两个二元信号量,天然实现了生产者与消费者对单一临界资源轮流且互斥的严格同步访问。
二、POSIX信号量接口的使用
在使用POSIX信号量之前,需要包含头文件<semaphore.h>。
2.1 初始化信号量
#include<semaphore.h>intsem_init(sem_t*sem,intpshared,unsignedintvalue);参数说明:
- sem:指向要初始化的信号量对象的指针。
- pshared:0表示线程间共享;非零表示进程间共享。
- value:信号量的初始值(表示可用资源的数量)。
2.2 销毁信号量
intsem_destroy(sem_t*sem);用于释放信号量占用的系统资源。在销毁信号量之前,应确保没有线程正在等待该信号量。
2.3 等待信号量(P操作)
intsem_wait(sem_t*sem);// P操作功能说明:
- 等待信号量。如果信号量的值大于0,则将信号量的值减1并立即返回。
- 如果信号量的值为0,则调用线程将被阻塞,直到信号量的值大于0(即有其他线程发布了信号量)。
2.4 发布信号量(V操作)
intsem_post(sem_t*sem);// V操作功能说明:
- 发布信号量,表示资源使用完毕,可以归还资源了。将信号量值加1。
三、基于环形缓冲的生产者消费者模型源码解析
手搓代码,如有错误,还请指出私信。
3.1RingQueue.hpp
#pragmaonce#include<unistd.h>#include<cstdio>#include<vector>#include"Mutex.hpp"#include"Sem.hpp"usingnamespaceMySem;usingnamespaceMyMutex;constsize_t DEFULT_SIZE=5;namespaceProducerAndConsumerProblemByRingQueue{template<typenameT>classRingQueue{public:RingQueue(size_t N=DEFULT_SIZE):_capacity(N),_blank_sem(N),_data_sem(0),_c_step(0),_p_step(0){_RingQueue.resize(_capacity);}voidEqueue(constT&args){//Producer_blank_sem.P();{_p_mutex.Lock();_RingQueue[_p_step]=args;_p_step++;_p_step%=_capacity;_data_sem.V();_p_mutex.UnLock();}}TPop(){//ConsumerT data;_data_sem.P();{_c_mutex.Lock();data=_RingQueue[_c_step];_c_step++;_c_step%=_capacity;_blank_sem.V();_c_mutex.UnLock();}returndata;}~RingQueue(){}private:std::vector<T>_RingQueue;size_t _capacity;Sem _blank_sem;Sem _data_sem;size_t _c_step;size_t _p_step;Mutex _c_mutex;Mutex _p_mutex;};}3.2Sem.hpp
#pragmaonce#include<semaphore.h>namespaceMySem{classSem{public:Sem(size_t size){sem_init(&_sem,0,size);}voidP(){sem_wait(&_sem);}voidV(){sem_post(&_sem);}~Sem(){sem_destroy(&_sem);}private:sem_t _sem;};}3.3Task.hpp
#include<functional>#include<iostream>#include<vector>usingtask_t=std::function<void(void)>;constsize_t TASK_NUM=3;voidMemaryProblem(){std::cout<<"This is a Memary Problem"<<std::endl;}voidSQLProblem(){std::cout<<"This is a SQL Problem"<<std::endl;}voidInternetProblem(){std::cout<<"This is a Internet Problem"<<std::endl;}classTaskManager{public:TaskManager()=default;~TaskManager(){}voidRegister(task_t task){_TaskCollection.push_back(task);}task_toperator[](size_t i){return_TaskCollection[i];}private:std::vector<task_t>_TaskCollection;};3.4Mutex.hpp
#pragmaonce#include<pthread.h>namespaceMyMutex{classMutex{public:Mutex(){pthread_mutex_init(&_mutex,nullptr);}voidLock(){pthread_mutex_lock(&_mutex);}voidUnLock(){pthread_mutex_unlock(&_mutex);}~Mutex(){pthread_mutex_destroy(&_mutex);}private:pthread_mutex_t _mutex;};}3.5Main.cc
#include"RingQueue.hpp"#include"Task.hpp"#include<ctime>usingnamespaceProducerAndConsumerProblemByRingQueue;constsize_t THREAD_NUM=5;classThreadData{public:ThreadData(RingQueue<task_t>*ringqueue,char*name):_ringqueue(ringqueue),_name(name){}RingQueue<task_t>*_ringqueue;char*_name;};task_tRandTask(){TaskManager tmang;tmang.Register(MemaryProblem);tmang.Register(SQLProblem);tmang.Register(InternetProblem);returntmang[rand()%TASK_NUM];}void*Producer(void*args){char*name=static_cast<ThreadData*>(args)->_name;RingQueue<task_t>*ringqueue=static_cast<ThreadData*>(args)->_ringqueue;while(true){std::cout<<name<<"生产一个任务 "<<std::endl;ringqueue->Equeue(RandTask());}delete[](static_cast<ThreadData*>(args)->_name);}void*Consumer(void*args){char*name=static_cast<ThreadData*>(args)->_name;RingQueue<task_t>*ringqueue=static_cast<ThreadData*>(args)->_ringqueue;while(true){std::cout<<name<<"消费一个任务 "<<std::endl;task_t task=ringqueue->Pop();task();}delete[](static_cast<ThreadData*>(args)->_name);}intmain(){srand((unsignedint)time(NULL));std::vector<pthread_t>p_thread;std::vector<pthread_t>c_thread;RingQueue<task_t>*ringqueue=newRingQueue<task_t>();//生产者们for(inti=0;i<THREAD_NUM;i++){char*name=newchar[64];intn=snprintf(name,64,"ProducerThread-%d",i);(void)n;ThreadData*data=newThreadData(ringqueue,name);pthread_t tid;pthread_create(&tid,nullptr,Producer,data);p_thread.push_back(tid);}//消费者们for(inti=0;i<THREAD_NUM;i++){char*name=newchar[64];intn=snprintf(name,64,"ComsumerThread-%d",i);(void)n;ThreadData*data=newThreadData(ringqueue,name);pthread_t tid;pthread_create(&tid,nullptr,Consumer,data);c_thread.push_back(tid);}for(autoe:p_thread)pthread_join(e,nullptr);for(autoe:c_thread)pthread_join(e,nullptr);return0;}四、以上代码中出现的相关问题
在基于信号量实现环形队列的生产者-消费者模型中,涉及临界资源访问顺序、并发效率优化以及信号量特性的几个核心问题。
4.1 加锁与申请信号量的先后顺序问题
在代码实现时,关于“申请信号量(P操作)”与“申请互斥锁(Lock)”的先后顺序存在两种写法:
- 顺序一(先加锁,再申请信号量):线程先获取互斥锁进入临界区,再进行
sem_wait申请资源。 - 顺序二(先申请信号量,再加锁):线程先调用
sem_wait申请资源,成功拿到资源后再获取互斥锁。
两种方式在功能逻辑上均能正常运行,但先申请信号量,再加锁在多线程环境下的执行效率明显更高。
4.1.1 正确性保证
信号量的 P/V 操作由操作系统底层保证其原子性,无需互斥锁对其进行二次保护。
4.1.2 高效性对比(买电影票比喻)
可以通过“购买电影票”的例子直观解释两者的效率差异:
- 先加锁,再申请信号量:类似于所有人排成一条单列长队,只有排到队伍最前面的人才能拿出手机尝试买票。如果买票失败,该线程挂起等待,导致身后排队的所有人均被阻塞,整体效率低下。
- 先申请信号量,再加锁:类似于所有人先在网络上各自并发抢票,抢到票的人再去影院门口排队核验入场。
在并发场景下,若采用先申请信号量的逻辑,当某一个线程拿到资源并获取锁在临界区内更新队列下标时,其他线程完全可以并发地执行 P 操作去预分配资源,从而最大化利用多线程并发优势。