OS80.【Linux】基于环形队列的单生产者-单消费者模型

📅 2026/8/14 12:26:53
OS80.【Linux】基于环形队列的单生产者-单消费者模型
目录1.知识回顾生产者-消费者模型信号量的概念POSIX信号量环形队列2.理论分类讨论3种情况P操作和V操作关注的资源生产者工作的伪代码消费者工作的伪代码3.单生产者、单消费者代码准备工作环形队列类框架成员变量成员函数构造函数析构函数push入队函数pop出队函数P函数V函数生产者线程函数消费者线程函数main函数完整代码运行结果4.再解LeetCode循环队列题代码错误代码正确代码1: 使用现成的函数正确代码2: 手动访问sem_t底层数据结构(偏难)查看leetcode评测机使用的glibc库的版本找到2.39稳定版的sem_t的实现sem_tsem_initnew_sem检查leetcode评测机的环境提交结果1.知识回顾生产者-消费者模型参见以下文章:OS78.【Linux】线程互斥(6) 基于阻塞队列的生产者-消费者模型(单生产者、单消费者的初步版本)OS79.【Linux】线程互斥(7) 基于阻塞队列的生产者-消费者模型(多生产者、多消费者)信号量的概念参见OS55.【Linux】理解信号量(不是信号)文章POSIX信号量参见OS80.【Linux】POSIX信号量文章环形队列环形队列也成为循环队列,之前在以下文章讲过环形队列:L35.【LeetCode题解】循环队列(数组解法)L36.【LeetCode笔记】循环队列(链表解法)LeetCode上的需要实现7个成员函数://class MyCircularQueue MyCircularQueue(int k); bool enQueue(int value); bool deQueue() int Front() int Rear() bool isEmpty() bool isFull()这里代码实现基于环形队列的生产者-消费者模型不需要判空的isEmpty和判满的isFull,可以用信号量,而且不需要像L35和L36那样开辟k1个元素2.理论分类讨论3种情况按照生产者-消费者模型的逻辑,生产者将数据入队列(push),消费者将数据出队列(pop)这里使用信号量把临界资源分成多份队列不空且不满时,生产者指针和消费者指针指向不同的位置:结论: 生产者指针和消费指针只要没有访问同一个元素,那么就能同时进行生产和消费; 队列不空且不满时,两指针指向不同的位置队列为满,生产者指针和消费者指针都指向同一个位置:队列为空,生产者指针和消费者指针都指向同一个位置:结论: 生产者指针和消费者指针指向同一个位置的时候,生产者和消费者之间是互斥关系(之前在OS78.【Linux】线程互斥(6) 基于阻塞队列的生产者-消费者模型(单生产者、单消费者的初步版本)文章讲过),它们不能访问同一个临界资源,换句话说,它们不能同时访问队列否则可能出现消费者套圈生产者(消费者指针指向的位置超过生产者指针指向的位置)或者生产者套圈消费者,由于信号量的管控,不会出现这种情况结论: 最开始队列为空,一定是生产者先执行; 如果队列满了,只能消费者执行P操作和V操作关注的资源P操作是申请信号量,对信号量的值-1,关注的是剩余空间V操作是释放信号量,对信号量的值1,关注的是还有多少剩余数据生产者工作的伪代码P(剩余空间); //剩余空间-- 执行生产,入队列(); V(剩余数据); //剩余数据消费者工作的伪代码P(剩余数据); //剩余数据-- 执行消费,出队列(); V(剩余空间); //剩余空间生产者和消费者P操作、V操作不一样3.单生产者、单消费者代码准备工作新建以下文件:producer_consumer_ring_queue/ ├── makefile └── project.cppmakefile写入:project.out:project.cpp g -o $ $^ -g -stdc11 -lpthread .PHONY:clean clean: rm -f project.out以单生产者、单消费者为例:环形队列类框架为了支持插入不同类型的数据,使用模版template class T class ring_queue { public: //...... private: //...... };成员变量这里环形队列的底层用vector数组实现: vectorT _ring_queue真实情况下的环形队列的容量是有上限的: int _ max_capacity从上面的理论知道: 需要设置生产者、消费者的指针,这里为下标:int _c_i; //consumer index int _p_i; //producer index将环形队列拆为多份临界资源,需要信号量.而且要定义2个:消费者关注的剩余数据资源(sem_t _c_data_sem;)、生产者关注的剩余空间资源(sem_t _p_space_sem)★★★单消费者、单生产者不需要锁,因为信号量保证单消费者、单生产者之间是互斥的,如果是多消费者、多生产者,需要使用锁来约束多消费者之间、多生产者之间是互斥的这里写单消费者、单生产者成员函数构造函数初始化成员变量:sem_init的pshared参数设置为0,表示线程间共享一开始vector为空,留出_max_capacity个空位修信号量: 剩余数据个数为0,剩余空间个数为队列容量_max_capacityring_queue(size_t max_capacity20) :_max_capacity(max_capacity) ,_c_i(0) ,_p_i(0) { sem_init(_c_data_sem,0,0); sem_init(_p_space_sem,0,_max_capacity); _ring_queue.resize(_max_capacity); }析构函数清空vector,销毁信号量~ring_queue() { sem_destroy(_c_data_sem); sem_destroy(_p_space_sem); _ring_queue.clear(); }push入队函数按照提供的伪代码:P(剩余空间); //剩余空间-- 执行生产,入队列(); V(剩余数据); //剩余数据那么:void push(const T obj) { P(_p_space_sem); _ring_queue[_p_i]obj; _p_i; _p_i%_max_capacity; V(_c_data_sem); }注意: _c_i可能越界,取余维持环形特性pop出队函数按照提供的伪代码:P(剩余数据); //剩余数据-- 执行消费,出队列(); V(剩余空间); //剩余空间那么:T pop() { P(_c_data_sem); T ret_ring_queue[_c_i]; _c_i; _c_i%_max_capacity; V(_p_space_sem); return ret; }Push和pop不需要判断为空和为满,信号量会控制, 生产者和消费者不会同时访问一个位置(为空,只会让生产者访问;为满,只会让消费者访问),更不会出现”套圈”的情况P函数P函数执行P操作,就是执行sem_waitvoid P(sem_t sem) { sem_wait(sem); }V函数V函数执行V操作,就是执行sem_post,阻塞等待void V(sem_t sem) { sem_post(sem); }生产者线程函数为了简单起见,这里生产随机整数:void* produce(void* args) { ring_queueint* bqstatic_castring_queueint*(args); for (;;) { int datarand()%10; bq-push(data); printf(生产者生产了数据%d\n,data); } return nullptr; }消费者线程函数void* consume(void* args) { ring_queueint* bqstatic_castring_queueint*(args); for (;;) { int databq-pop(); printf(消费者消费了数据%d\n,data); } return nullptr; }main函数int main() { srand((unsigned int)time(NULL)); ring_queueint* bqnew ring_queueint(); pthread_t producer,consumer; pthread_create(producer,nullptr,produce,bq); pthread_create(consumer,nullptr,consume,bq); pthread_join(producer,nullptr); pthread_join(consumer,nullptr); delete bq; return 0; }完整代码#include pthread.h #include unistd.h #include cstdio #include cstdlib #include vector #include semaphore.h template class T class ring_queue { public: ring_queue(size_t max_capacity20) :_max_capacity(max_capacity) ,_c_i(0) ,_p_i(0) { sem_init(_c_data_sem,0,0); sem_init(_p_space_sem,0,_max_capacity); _ring_queue.resize(_max_capacity); } void push(const T obj) { P(_p_space_sem); _ring_queue[_p_i]obj; _p_i; _p_i%_max_capacity; V(_c_data_sem); } T pop() { P(_c_data_sem); T ret_ring_queue[_c_i]; _c_i; _c_i%_max_capacity; V(_p_space_sem); return ret; } ~ring_queue() { sem_destroy(_c_data_sem); sem_destroy(_p_space_sem); _ring_queue.clear(); } private: void V(sem_t sem) { sem_post(sem); } void P(sem_t sem) { sem_wait(sem); } std::vectorT _ring_queue; int _max_capacity; int _c_i; int _p_i; sem_t _c_data_sem; sem_t _p_space_sem; }; void* produce(void* args) { ring_queueint* bqstatic_castring_queueint*(args); for (;;) { int datarand()%10; bq-push(data); printf(生产者生产了数据%d\n,data); } return nullptr; } void* consume(void* args) { ring_queueint* bqstatic_castring_queueint*(args); for (;;) { int databq-pop(); printf(消费者消费了数据%d\n,data); } return nullptr; } int main() { srand((unsigned int)time(NULL)); ring_queueint* bqnew ring_queueint(); pthread_t producer,consumer; pthread_create(producer,nullptr,produce,bq); pthread_create(consumer,nullptr,consume,bq); pthread_join(producer,nullptr); pthread_join(consumer,nullptr); delete bq; return 0; }运行结果produce和consume各sleep(1):produce进行sleep(1),consume进行sleep(3):生产和消费符合队列先进先出的特点4.再解LeetCode循环队列题https://leetcode.cn/problems/design-circular-queue/设计你的循环队列实现。 循环队列是一种线性数据结构其操作表现基于 FIFO先进先出原则并且队尾被连接在队首之后以形成一个循环。它也被称为“环形缓冲器”。循环队列的一个好处是我们可以利用这个队列之前用过的空间。在一个普通队列里一旦一个队列满了我们就不能插入下一个元素即使在队列前面仍有空间。但是使用循环队列我们能使用这些空间去存储新的值。你的实现应该支持如下操作MyCircularQueue(k): 构造器设置队列长度为 k 。Front: 从队首获取元素。如果队列为空返回 -1 。Rear: 获取队尾元素。如果队列为空返回 -1 。enQueue(value): 向循环队列插入一个元素。如果成功插入则返回真。deQueue(): 从循环队列中删除一个元素。如果成功删除则返回真。isEmpty(): 检查循环队列是否为空。isFull(): 检查循环队列是否已满。示例MyCircularQueue circularQueue new MyCircularQueue(3); // 设置长度为 3 circularQueue.enQueue(1); // 返回 true circularQueue.enQueue(2); // 返回 true circularQueue.enQueue(3); // 返回 true circularQueue.enQueue(4); // 返回 false队列已满 circularQueue.Rear(); // 返回 3 circularQueue.isFull(); // 返回 true circularQueue.deQueue(); // 返回 true circularQueue.enQueue(4); // 返回 true circularQueue.Rear(); // 返回 4提示所有的值都在 0 至 1000 的范围内操作数将在 1 至 1000 的范围内请不要使用内置的队列库。代码和上面的单生产者-单消费者模型不同,上面的代码入队如果不成功会阻塞等待,LeetCode此题必须用sem_trywait,出错立即返回,因为enQueue和deQueue都有返回值比较容易的是,LeetCode此题是单线程版本的,只有一个线程会执行MyCircularQueue里面的成员函数,而且生产者和消费者都是这个线程,不需要加锁消费者下标在队头,生产者下标在队尾错误代码class MyCircularQueue { public: MyCircularQueue(int k) :_max_capacity(k) ,_c_i(0) ,_p_i(0) { sem_init(_c_data_sem,0,0); sem_init(_p_space_sem,0,_max_capacity); _ring_queue.resize(_max_capacity); } bool enQueue(int obj) { if (!P(_p_space_sem)) return false; _ring_queue[_p_i]obj; _p_i; _p_i%_max_capacity; V(_c_data_sem); return true; } bool deQueue() { if (!P(_c_data_sem)) return false; _c_i; _c_i%_max_capacity; V(_p_space_sem); return true; } int Front() { if (isEmpty()) return -1; return _ring_queue[_c_i]; } int Rear() { if (isEmpty()) return -1; if (_p_i-10) return _ring_queue[_max_capacity-1]; return _ring_queue[_p_i-1]; } bool isEmpty() //剩余数据个数为0,剩余空间个数为满 { //借助P操作,但不能真的申请 return !P(_c_data_sem); } bool isFull() //剩余数据个数为满,剩余空间个数为0 { //借助P操作,但不能真的申请 return !P(_p_space_sem); } private: bool V(sem_t sem) { //sem_post(sem)返回值为0代表V操作成功 return sem_post(sem)0; } bool P(sem_t sem) { //sem_trywait(sem)返回值为0代表P操作成功 return sem_trywait(sem)0; } std::vectorint _ring_queue; int _max_capacity; int _c_i; int _p_i; sem_t _c_data_sem; sem_t _p_space_sem; };提交结果:第一次调用Rear返回6,第二次调用Rear返回-1,很奇怪错在isEmpty的实现上,如果队列不为空,调用isEmpty,由于return !P(_c_data_sem),那么P操作会成功! 但是isEmpty是不能修改信号量的值的正确代码1: 使用现成的函数方法: 使用sem_getvalue取得信号量的值class MyCircularQueue { public: MyCircularQueue(int k) :_max_capacity(k) ,_c_i(0) ,_p_i(0) { sem_init(_c_data_sem,0,0); sem_init(_p_space_sem,0,_max_capacity); _ring_queue.resize(_max_capacity); } bool enQueue(int obj) { if (!P(_p_space_sem)) return false; _ring_queue[_p_i]obj; _p_i; _p_i%_max_capacity; V(_c_data_sem); return true; } bool deQueue() { if (!P(_c_data_sem)) return false; _c_i; _c_i%_max_capacity; V(_p_space_sem); return true; } int Front() { if (isEmpty()) return -1; return _ring_queue[_c_i]; } int Rear() { if (isEmpty()) return -1; if (_p_i-10) return _ring_queue[_max_capacity-1]; return _ring_queue[_p_i-1]; } bool isEmpty() { int ret; sem_getvalue(_p_space_sem,ret); return ret_max_capacity; } bool isFull() { int ret; sem_getvalue(_c_data_sem,ret); return ret_max_capacity; } private: bool V(sem_t sem) { //sem_post(sem)返回值为0代表V操作成功 return sem_post(sem)0; } bool P(sem_t sem) { //sem_trywait(sem)返回值为0代表P操作成功 return sem_trywait(sem)0; } std::vectorint _ring_queue; int _max_capacity; int _c_i; int _p_i; sem_t _c_data_sem; sem_t _p_space_sem; };正确代码2: 手动访问sem_t底层数据结构(偏难)之前在OS80.【Linux】POSIX信号量文章讲过sem_t的底层实现:POSIX信号量的系统调用把sem_t联合体当结构体使用,通过强制类型转换访问内部成员,使用联合体为了对外隐藏结构查看leetcode评测机使用的glibc库的版本const char* version gnu_get_libc_version(); std::cout glibc version: version std::endl; const char* release gnu_get_libc_release(); std::cout glibc release: release std::endl;运行结果:找到2.39稳定版的sem_t的实现bootlin网托管了glibc 2.39稳定版的代码仓库: elixir.bootlin.com | glibc-2.39sem_t#if __WORDSIZE 64 # define __SIZEOF_SEM_T 32 #else # define __SIZEOF_SEM_T 16 #endif /* Value returned if sem_open failed. */ #define SEM_FAILED ((sem_t *) 0) typedef union { char __size[__SIZEOF_SEM_T]; long int __align; } sem_t;sem_initint __new_sem_init (sem_t *sem, int pshared, unsigned int value) { ASSERT_PTHREAD_INTERNAL_SIZE (sem_t, struct new_sem); /* Parameter sanity check. */ if (__glibc_unlikely (value SEM_VALUE_MAX)) { __set_errno (EINVAL); return -1; } pshared pshared ! 0 ? PTHREAD_PROCESS_SHARED : PTHREAD_PROCESS_PRIVATE; int err futex_supports_pshared (pshared); if (err ! 0) { __set_errno (err); return -1; } /* Map to the internal type. */ struct new_sem *isem (struct new_sem *) sem; /* Use the values the caller provided. */ #if __HAVE_64B_ATOMICS isem-data value; #else isem-value value SEM_VALUE_SHIFT; /* pad is used as a mutex on pre-v9 sparc and ignored otherwise. */ isem-pad 0; isem-nwaiters 0; #endif isem-private (pshared PTHREAD_PROCESS_PRIVATE ? FUTEX_PRIVATE : FUTEX_SHARED); return 0; } versioned_symbol (libc, __new_sem_init, sem_init, GLIBC_2_34);new_sem/* Semaphore variable structure. */ struct new_sem { #if __HAVE_64B_ATOMICS /* The data field holds both value (in the least-significant 32 bits) and nwaiters. */ # if __BYTE_ORDER __LITTLE_ENDIAN # define SEM_VALUE_OFFSET 0 # elif __BYTE_ORDER __BIG_ENDIAN # define SEM_VALUE_OFFSET 1 # else # error Unsupported byte order. # endif # define SEM_NWAITERS_SHIFT 32 # define SEM_VALUE_MASK (~(unsigned int)0) uint64_t data; int private; int pad; #else # define SEM_VALUE_SHIFT 1 # define SEM_NWAITERS_MASK ((unsigned int)1) unsigned int value; int private; int pad; unsigned int nwaiters; #endif };如果leetcode评测机支持__HAVE_64B_ATOMICS64为原子操作,那么new_sem定义为:/* Semaphore variable structure. */ struct new_sem { uint64_t data; int private;//和C的private冲突 int pad; };否则为:/* Semaphore variable structure. */ struct new_sem { unsigned int value; int private;//和C的private冲突 int pad; unsigned int nwaiters; };检查leetcode评测机的环境std::cout Architecture: ; #if defined(__x86_64__) || defined(__amd64__) std::cout x86_64 std::endl; #elif defined(__i386__) std::cout x86 std::endl; #elif defined(__aarch64__) std::cout ARM64 std::endl; #else std::cout unknown std::endl; #endif运行结果:64位一般支持64位原子操作,但__HAVE_64B_ATOMICS是 glibc 内部宏,用户根本无法知道__HAVE_64B_ATOMICS是否定义!__HAVE_64B_ATOMICS是否定义,信号量的值的算法不一样!//sem_init #if __HAVE_64B_ATOMICS isem-data value; #else isem-value value SEM_VALUE_SHIFT; /* pad is used as a mutex on pre-v9 sparc and ignored otherwise. */ isem-pad 0; isem-nwaiters 0; #endif可以先假设__HAVE_64B_ATOMICS定义了,如果过了说明没问题:#include semaphore.h /* Semaphore variable structure. glibc 2.39 x86_64 */ struct new_sem { uint64_t data;//在sem_init直接是isem-data value; int _private;//和C的private冲突 int pad; }; class MyCircularQueue { public: MyCircularQueue(int k) :_max_capacity(k) ,_c_i(0) ,_p_i(0) { sem_init(_c_data_sem,0,0); sem_init(_p_space_sem,0,_max_capacity); _ring_queue.resize(_max_capacity); } bool enQueue(int obj) { coutsizeof(sem_t); if (!P(_p_space_sem)) return false; _ring_queue[_p_i]obj; _p_i; _p_i%_max_capacity; V(_c_data_sem); return true; } bool deQueue() { if (!P(_c_data_sem)) return false; _c_i; _c_i%_max_capacity; V(_p_space_sem); return true; } int Front() { if (isEmpty()) return -1; return _ring_queue[_c_i]; } int Rear() { if (isEmpty()) return -1; if (_p_i-10) return _ring_queue[_max_capacity-1]; return _ring_queue[_p_i-1]; } bool isEmpty() { struct new_sem* isemreinterpret_caststruct new_sem*(_p_space_sem); return isem-data_max_capacity; } bool isFull() { struct new_sem* isemreinterpret_caststruct new_sem*(_c_data_sem); return isem-data_max_capacity; } private: bool V(sem_t sem) { //sem_post(sem)返回值为0代表V操作成功 return sem_post(sem)0; } bool P(sem_t sem) { //sem_trywait(sem)返回值为0代表P操作成功 return sem_trywait(sem)0; } std::vectorint _ring_queue; uint64_t _max_capacity; int _c_i; int _p_i; sem_t _c_data_sem; sem_t _p_space_sem; };提交通过100%提交结果正确代码1:正确代码2: