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

资讯详情

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

Re:Linux系统篇(五十七)线程篇 · 十:基于环形缓冲的生产者消费者模型与信号量

Re:Linux系统篇(五十七)线程篇 · 十:基于环形缓冲的生产者消费者模型与信号量 ◆ 博主名称 小此方-CSDN博客大家好欢迎来到小此方的博客。⭐️Linux系列个人专栏 【主题曲】Linux⭐️此方的GitHub github_此方⭐️Re系列专栏我们思考 (Rethink) · 我们重建 (Rebuild) · 我们记录 (Record)文章目录概要序論一、基于环形缓冲的生产者消费者模型原理介绍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.hpp3.2Sem.hpp3.3Task.hpp3.4Mutex.hpp3.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初始为 1sem_data初始为 0。系统变成了一种全新的同步互斥实现方式——通过这两个二元信号量天然实现了生产者与消费者对单一临界资源轮流且互斥的严格同步访问。二、POSIX信号量接口的使用在使用POSIX信号量之前需要包含头文件semaphore.h。2.1 初始化信号量#includesemaphore.hintsem_init(sem_t*sem,intpshared,unsignedintvalue);参数说明sem指向要初始化的信号量对象的指针。pshared0表示线程间共享非零表示进程间共享。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#includeunistd.h#includecstdio#includevector#includeMutex.hpp#includeSem.hppusingnamespaceMySem;usingnamespaceMyMutex;constsize_t DEFULT_SIZE5;namespaceProducerAndConsumerProblemByRingQueue{templatetypenameTclassRingQueue{public:RingQueue(size_t NDEFULT_SIZE):_capacity(N),_blank_sem(N),_data_sem(0),_c_step(0),_p_step(0){_RingQueue.resize(_capacity);}voidEqueue(constTargs){//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::vectorT_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#includesemaphore.hnamespaceMySem{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#includefunctional#includeiostream#includevectorusingtask_tstd::functionvoid(void);constsize_t TASK_NUM3;voidMemaryProblem(){std::coutThis is a Memary Problemstd::endl;}voidSQLProblem(){std::coutThis is a SQL Problemstd::endl;}voidInternetProblem(){std::coutThis is a Internet Problemstd::endl;}classTaskManager{public:TaskManager()default;~TaskManager(){}voidRegister(task_t task){_TaskCollection.push_back(task);}task_toperator[](size_t i){return_TaskCollection[i];}private:std::vectortask_t_TaskCollection;};3.4Mutex.hpp#pragmaonce#includepthread.hnamespaceMyMutex{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#includeRingQueue.hpp#includeTask.hpp#includectimeusingnamespaceProducerAndConsumerProblemByRingQueue;constsize_t THREAD_NUM5;classThreadData{public:ThreadData(RingQueuetask_t*ringqueue,char*name):_ringqueue(ringqueue),_name(name){}RingQueuetask_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*namestatic_castThreadData*(args)-_name;RingQueuetask_t*ringqueuestatic_castThreadData*(args)-_ringqueue;while(true){std::coutname生产一个任务 std::endl;ringqueue-Equeue(RandTask());}delete[](static_castThreadData*(args)-_name);}void*Consumer(void*args){char*namestatic_castThreadData*(args)-_name;RingQueuetask_t*ringqueuestatic_castThreadData*(args)-_ringqueue;while(true){std::coutname消费一个任务 std::endl;task_t taskringqueue-Pop();task();}delete[](static_castThreadData*(args)-_name);}intmain(){srand((unsignedint)time(NULL));std::vectorpthread_tp_thread;std::vectorpthread_tc_thread;RingQueuetask_t*ringqueuenewRingQueuetask_t();//生产者们for(inti0;iTHREAD_NUM;i){char*namenewchar[64];intnsnprintf(name,64,ProducerThread-%d,i);(void)n;ThreadData*datanewThreadData(ringqueue,name);pthread_t tid;pthread_create(tid,nullptr,Producer,data);p_thread.push_back(tid);}//消费者们for(inti0;iTHREAD_NUM;i){char*namenewchar[64];intnsnprintf(name,64,ComsumerThread-%d,i);(void)n;ThreadData*datanewThreadData(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 操作去预分配资源从而最大化利用多线程并发优势。好的本期内容就到这里如果对你有帮助还不要忘记点赞三联支持。我是此方我们下期再见。bye! Linux、C、算法持续连载中欢迎关注WeChat Official Account 【此方的技术栈】。
返回列表