You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

仅用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_;
};

问题根源

  1. 跨信号量复用Token:Put操作中用empty_的Token去Releasetaken_,Take操作中用taken_的Token去Releaseempty_,违反了TaggedSemaphore的Token语义(Token应仅用于对应信号量的Release),导致信号量计数与缓冲区实际状态完全脱节。
  2. 多余的锁无意义:put_mutex_和take_mutex_并未解决核心同步问题,反而增加了不必要的阻塞点,无法保证信号量获取与缓冲区操作的原子性。
  3. 临界区时机错误:信号量获取与临界区进入之间无原子性保证,多个线程可能在获取信号量后乱序进入临界区,导致缓冲区操作冲突。

解决方案

核心思路是保证信号量等待、临界区操作、信号量通知的顺序性与原子性,同时严格遵循TaggedSemaphore的Token语义:

  1. 移除多余的put_mutex_和take_mutex_,仅保留三个核心信号量:
    • empty_:初始值为队列容量,标记可用空槽数量
    • taken_:初始值为0,标记已占用槽数量
    • mutex_:初始值为1,作为互斥信号量保护缓冲区操作
  2. 调整操作顺序:先等待对应信号量(空槽/元素),再进入临界区修改缓冲区,最后通知对应信号量(元素可用/空槽可用)
  3. 严格保证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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.25 09:06:36