无锁环 C实现

发布时间:2026/8/9 12:17:37
无锁环 C实现 无锁环MPMC目标1、实现多生产者、多消费者无锁队列2、不用互斥锁3、数组循环复用4、FIFO5、不能丢消息、不能重复读、不能读到脏数据底层载体1、固定数组大小 slots[]2、capacity 必须是 2 的幂次#includestdio.h#includestdlib.h#includepthread.h#includestdatomic.h#includestring.h/* PPT无锁环关键点 1. 环形数组缓存元素定长这里用uint32_t4字节符合PPT 4字节整数倍 2. 多生产者、多消费者MPMC 3. 不加锁全部依靠原子CAS 4. 全局共享prod_head/prod_tailcons_head/cons_tail线程本地私有副本做缓存减少共享变量访问 */typedefuint32_telem_t;// 定长元素4字节和PPT要求匹配typedefstruct{atomic_uint seq;// 槽位序列号无锁同步核心elem_tdata;}slot_t;typedefstruct{slot_t*slots;size_tcapacity;// 必须是2^Nsize_tmask;// capacity‑1代替%取模// 全局共享变量所有线程可见atomic_size_tprod_head;// 生产者headatomic_size_tprod_tail;// 生产者tailatomic_size_tcons_head;// 消费者headatomic_size_tcons_tail;// 消费者tail}lockfree_ring_t;// 创建无锁环capacity必须为2的幂lockfree_ring_t*ring_create(size_tcapacity){lockfree_ring_t*rmalloc(sizeof(lockfree_ring_t));r-capacitycapacity;r-maskcapacity-1;r-slotscalloc(capacity,sizeof(slot_t));atomic_init(r-prod_head,0);atomic_init(r-prod_tail,0);atomic_init(r-cons_head,0);atomic_init(r-cons_tail,0);for(size_ti0;icapacity;i){atomic_init(r-slots[i].seq,i);}returnr;}// 销毁voidring_destroy(lockfree_ring_t*r){free(r-slots);free(r);}/** * 多生产者入队 enqueue * 返回0成功‑1队列满 */intring_enqueue(lockfree_ring_t*r,elem_tval){size_tpos,next;// 自旋CAS争抢prod_head多生产者竞争while(1){posatomic_load_explicit(r-prod_head,memory_order_relaxed);nextpos1;size_tcons_tatomic_load_explicit(r-cons_tail,memory_order_acquire);if((next-cons_t)r-capacity){return-1;// 队列满}// CAS抢占位置if(atomic_compare_exchange_weak_explicit(r-prod_head,pos,next,memory_order_relaxed,memory_order_relaxed)){break;}}size_tidxposr-mask;slot_t*sr-slots[idx];// 等待槽位空闲上一轮消费完成while(atomic_load_explicit(s-seq,memory_order_acquire)!pos){;}s-dataval;atomic_store_explicit(s-seq,pos1,memory_order_release);// 更新prod_tail推进发布指针允许其他线程看到写入while(!atomic_compare_exchange_weak_explicit(r-prod_tail,pos,next,memory_order_release,memory_order_relaxed)){posatomic_load_explicit(r-prod_tail,memory_order_relaxed);}return0;}/** * 多消费者出队 dequeue * 返回0成功‑1队空 */intring_dequeue(lockfree_ring_t*r,elem_t*out){size_tpos,next;while(1){posatomic_load_explicit(r-cons_head,memory_order_relaxed);nextpos1;size_tprod_tatomic_load_explicit(r-prod_tail,memory_order_acquire);if(posprod_t){return-1;// 队列为空}if(atomic_compare_exchange_weak_explicit(r-cons_head,pos,next,memory_order_relaxed,memory_order_relaxed)){break;}}size_tidxposr-mask;slot_t*sr-slots[idx];// 等待生产者写完数据while(atomic_load_explicit(s-seq,memory_order_acquire)!pos1){;}*outs-data;atomic_store_explicit(s-seq,posr-capacity,memory_order_release);// 更新cons_tailwhile(!atomic_compare_exchange_weak_explicit(r-cons_tail,pos,next,memory_order_release,memory_order_relaxed)){posatomic_load_explicit(r-cons_tail,memory_order_relaxed);}return0;}// 测试多生产者多消费者#defineRING_CAP(14)// 162的幂#definePRODUCER_NUM3#defineCONSUMER_NUM2#definePER_PROD_CNT20lockfree_ring_t*g_ring;/* 函数功能向无锁环里塞数据的生产者线程 *//* 入参arg 把 tid 数字本身伪装成指针值没有对应内存*//* 返回值现成结束返回值可以给 pthread_join 拿到 */void*producer_thread(void*arg){longtid(long)arg;for(inti0;iPER_PROD_CNT;i){elem_tvtid*1000i;while(ring_enqueue(g_ring,v)!0){;// 满则自旋}printf([P%ld] enqueue: %u\n,tid,v);}returnNULL;}void*consumer_thread(void*arg){longtid(long)arg;elem_tval;inttotalPER_PROD_CNT*PRODUCER_NUM;staticatomic_int recv_cnt0;while(atomic_load(recv_cnt)total){if(ring_dequeue(g_ring,val)0){atomic_fetch_add(recv_cnt,1);printf( [C%ld] dequeue: %u\n,tid,val);}}returnNULL;}intmain(void){g_ringring_create(RING_CAP);pthread_tprods[PRODUCER_NUM];pthread_tcons[CONSUMER_NUM];for(inti0;iPRODUCER_NUM;i){pthread_create(prods[i],NULL,producer_thread,(void*)(long)i);}for(inti0;iCONSUMER_NUM;i){pthread_create(cons[i],NULL,consumer_thread,(void*)(long)i);}for(inti0;iPRODUCER_NUM;i)pthread_join(prods[i],NULL);for(inti0;iCONSUMER_NUM;i)pthread_join(cons[i],NULL);ring_destroy(g_ring);printf(all done\n);return0;}给坏代码打补丁/** * 多生产者入队 enqueue * 返回0成功‑1队列满 */intring_enqueue(lockfree_ring_t*r,elem_tval){size_tpos,next;// 自旋CAS争抢prod_head多生产者竞争while(1){posatomic_load_explicit(r-prod_head,memory_order_relaxed);nextpos1;size_tcons_tatomic_load_explicit(r-cons_tail,memory_order_acquire);if((next-cons_t)r-capacity){return-1;// 队列满}// CAS抢占位置if(atomic_compare_exchange_weak_explicit(r-prod_head,pos,next,memory_order_relaxed,memory_order_relaxed)){break;}}size_tidxposr-mask;slot_t*sr-slots[idx];// 等待槽位空闲上一轮消费完成while(atomic_load_explicit(s-seq,memory_order_acquire)!pos){;}s-dataval;atomic_store_explicit(s-seq,pos1,memory_order_release);// 更新prod_tail推进发布指针允许其他线程看到写入// 以下 for 循环函数是给注释代码打补丁注释代码存在问题/*while(!atomic_compare_exchange_weak_explicit( r-prod_tail, pos, next, memory_order_release, memory_order_relaxed )){ pos atomic_load_explicit(r-prod_tail, memory_order_relaxed); }*/for(;;){size_texpectedatomic_load_explicit(r-prod_tail,memory_order_relaxed);if(expected!pos){continue;//还没轮到我等别的线程推进}size_tdesiredexpected1;if(atomic_compare_exchange_weak_explicit(r-prod_tail,expected,desired,memory_order_release,memory_order_relaxed)){break;}// CAS失败expected已经更新下一轮重新判断}return0;}