仅用Semaphore实现BlockingQueue的线程同步问题求助
仅用TaggedSemaphore实现BlockingQueue的线程安全问题
要求仅通过TaggedSemaphore实现线程安全的BlockingQueue,禁止使用循环、条件变量、普通互斥锁和if-else语句。当前实现的代码在并发场景下存在线程安全问题:当多个线程同时执行Put和Take操作时,信号量操作与缓冲区修改不同步,导致临界区处理错误数据;但提前锁定互斥信号量又会在队列满/空时引发死锁。
当前代码实现
template <typename T> class BlockingQueue { using Token = typename TaggedSemaphore<T>::Token; using Guard = typename TaggedSemaphore<T>::Guard; public: explicit BlockingQueue(size_t capacity): empty_ (capacity), taken_ (0), mutex_(1), put_mutex_(1), take_mutex_(1) { } // Inserts the specified element into this queue, // waiting if necessary for space to become available. void Put(T value) { Guard put_guard (put_mutex_); Token tok (std::move (empty_.Acquire())); Guard guard (mutex_); buffer_.push_back(std::move(value)); taken_.Release(std::move(tok)); } // Retrieves and removes the head of this queue, // waiting if necessary until an element becomes available T Take() { Guard take_guard (take_mutex_); Token tok (std::move (taken_.Acquire())); Guard guard (mutex_); T ret_value = std::move(buffer_.front()); buffer_.pop_front(); empty_.Release(std::move(tok)); return std::move(ret_value); } private: TaggedSemaphore<T> empty_, taken_, mutex_, put_mutex_, take_mutex_; std::deque<T> buffer_; };
问题根源
- 跨信号量复用Token:Put操作中用
empty_的Token去Releasetaken_,Take操作中用taken_的Token去Releaseempty_,违反了TaggedSemaphore的Token语义(Token应仅用于对应信号量的Release),导致信号量计数与缓冲区实际状态完全脱节。 - 多余的锁无意义:
put_mutex_和take_mutex_并未解决核心同步问题,反而增加了不必要的阻塞点,无法保证信号量获取与缓冲区操作的原子性。 - 临界区时机错误:信号量获取与临界区进入之间无原子性保证,多个线程可能在获取信号量后乱序进入临界区,导致缓冲区操作冲突。
解决方案
核心思路是保证信号量等待、临界区操作、信号量通知的顺序性与原子性,同时严格遵循TaggedSemaphore的Token语义:
- 移除多余的
put_mutex_和take_mutex_,仅保留三个核心信号量:empty_:初始值为队列容量,标记可用空槽数量taken_:初始值为0,标记已占用槽数量mutex_:初始值为1,作为互斥信号量保护缓冲区操作
- 调整操作顺序:先等待对应信号量(空槽/元素),再进入临界区修改缓冲区,最后通知对应信号量(元素可用/空槽可用)
- 严格保证Token仅用于对应信号量的操作,跨信号量通知使用独立的Release操作(若
TaggedSemaphore支持无Token的Release,或通过合法方式生成对应信号量的Token)
修正后的代码
template <typename T> class BlockingQueue { using Token = typename TaggedSemaphore<T>::Token; using Guard = typename TaggedSemaphore<T>::Guard; public: explicit BlockingQueue(size_t capacity): empty_(capacity), taken_(0), mutex_(1) { } void Put(T value) { // 等待队列有空槽,获取empty_的Token(消耗一个空槽) Token empty_token = std::move(empty_.Acquire()); // 进入临界区,保证同一时间只有一个线程操作缓冲区 { Guard mutex_guard(mutex_); buffer_.push_back(std::move(value)); } // 离开作用域自动释放mutex_ // 通知Take线程有元素可取,增加taken_的计数 taken_.Release(); // 释放empty_的Token(若TaggedSemaphore要求必须归还Token,否则可省略) empty_.Release(std::move(empty_token)); } T Take() { // 等待队列有元素,获取taken_的Token(消耗一个已用槽) Token taken_token = std::move(taken_.Acquire()); T ret_value; // 进入临界区操作缓冲区 { Guard mutex_guard(mutex_); ret_value = std::move(buffer_.front()); buffer_.pop_front(); } // 自动释放mutex_ // 通知Put线程有空槽可用,增加empty_的计数 empty_.Release(); // 释放taken_的Token(若TaggedSemaphore要求必须归还Token,否则可省略) taken_.Release(std::move(taken_token)); return std::move(ret_value); } private: TaggedSemaphore<T> empty_, taken_, mutex_; std::deque<T> buffer_; };
关键说明
- 信号量等待操作(
empty_.Acquire()/taken_.Acquire())必须在临界区之前执行,避免队列满/空时持有互斥信号量阻塞,引发死锁。 - 临界区用局部作用域包裹
mutex_的Guard,确保缓冲区操作完成后立即释放互斥信号量,提升并发性能。 - 严格遵循
TaggedSemaphore的Token语义,每个信号量的Acquire与Release对应,跨信号量的通知使用独立的Release操作,保证信号量计数与缓冲区状态完全一致。
内容的提问来源于stack exchange,提问作者Amir Mokhammed-Ali
相关产品推荐
相关产品推荐

