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

为何基于std::condition_variable的单消费者多生产者代码偶发挂起?

单消费者多生产者队列偶发挂起问题排查与修复

问题描述

我在C++中实现了标准的单消费者多生产者模式,对队列任务数量做了限制。Worker在独立线程运行消息队列,生产者向其发送任务;若队列已达max_num_tasks_个任务,生产者需等待。代码在VS2022的Windows x64平台用v143工具集编译。这段代码偶尔会在Worker::Send()的cv_.wait(lock, [&] {return (mq.Size() < max_num_tasks_); });语句处挂起,恳请指出问题所在。

相关代码实现

Worker类实现

#pragma once
#include <queue>
#include <mutex>
#include <functional>

class Worker {
public:
    ~Worker() {
        Send([this] {done = true; });
        thd.join();
    }

    Worker(size_t max_num_tasks) : max_num_tasks_(max_num_tasks), done(false), thd([this] {
        while (!done) {
            mq.PopFront()();
            cv_.notify_all();
        }
        })
    { }

    void Send(std::function<void()>&& m) {
        {
            std::unique_lock<std::mutex> lock(m_);
            cv_.wait(lock, [&] {return (mq.Size() < max_num_tasks_); });
        }
        mq.PushBack(std::move(m));
    }
private:
    bool done;
    size_t max_num_tasks_;
    ThreadSafeQueue<std::function<void()>> mq;
    std::thread thd;
    std::mutex m_;
    std::condition_variable cv_;
};

ThreadSafeQueue类实现

#pragma once
#include <queue>
#include <mutex>
#include <functional>

template <typename T>
class ThreadSafeQueue {
public:
    void PushBack(T&& val) {
        {
            std::unique_lock<std::mutex> lock(q_mutex);
            q.push(std::move(val));
        }
        cv.notify_one();
    }

    void PushBack(const T& val) {
        {
            std::unique_lock<std::mutex> lock(q_mutex);
            q.push(val);
        }
        cv.notify_one();
    }

    T PopFront() {
        std::unique_lock<std::mutex> lock(q_mutex);
        cv.wait(lock, [&] { return (!q.empty()); });
        T v = q.front();
        q.pop();
        return std::move(v);
    }

    bool Empty() const {
        std::unique_lock<std::mutex> lock(q_mutex);
        return q.empty();
    }

    size_t Size() const {
        std::unique_lock<std::mutex> lock(q_mutex);
        return q.size();
    }
private:
    mutable std::mutex q_mutex;
    std::condition_variable cv;
    std::queue<T> q;
};

偶发触发问题的单元测试示例

TEST(WorkerTests, Compute_PI) {

    std::atomic<double> value;
    auto multiply_by_pi_over_four = [&value] {
        int n = 750;
        double v = 0;
        for (int i = 0; i < n; i++) {
            v += std::pow(-1, i) / (2 * i + 1);
        }
        value = value * v;
    };

    auto add_pi_over_four = [&value] {
        int n = 750;
        double v = 0;
        for (int i = 0; i < n; i++) {
            v += std::pow(-1, i) / (2 * i + 1);
        }
        value = value + v;
    };

    auto add_one = [&value] {
        value = value + 1;
        std::this_thread::sleep_for(std::chrono::milliseconds(3));
    };

    auto multiply_by_three = [&value] {
        value = value * 3;
        std::this_thread::sleep_for(std::chrono::milliseconds(1));
    };

    for (int i = 0; i < 100; i++) {
        value = 0.0;
        {
            Worker worker(1);
            worker.Send(add_one);
            worker.Send(multiply_by_pi_over_four);
            worker.Send(multiply_by_three);
            worker.Send([]() {});
            worker.Send(add_pi_over_four);
        }
        EXPECT_GE(3.15, value.load());
        EXPECT_LE(3.14, value.load());
    }
}

问题根源分析

1. 竞态条件:队列大小判断与入队操作不同步

Worker::Send()中,判断队列大小mq.Size() < max_num_tasks_后立刻释放了Worker的锁m_,再执行mq.PushBack()。这导致多个生产者线程可能同时通过条件判断,然后依次执行入队操作,瞬间让队列大小超过max_num_tasks_。

比如:

  • 生产者A拿到m_锁,判断队列大小为0(小于max=1),释放锁。
  • 生产者B同时拿到m_锁,同样判断队列大小为0,释放锁。
  • 生产者A先完成入队,队列大小变为1。
  • 生产者B再完成入队,队列大小变为2,超过限制。后续生产者的cv_.wait()会因为队列一直处于“满”状态而挂起——因为消费者处理任务后,只会通知ThreadSafeQueue内部的条件变量,不会通知Worker的cv_。

2. 双重锁导致的状态不一致

mq.Size()内部使用ThreadSafeQueue自己的锁q_mutex,而Worker的条件判断使用的是m_锁。这两个锁没有关联,导致判断队列大小时,队列状态可能在判断完成后立刻变化,条件判断的结果已经失效。

3. done变量的线程可见性问题

done是普通bool变量,没有同步机制,消费者线程可能无法及时看到生产者设置的done = true,导致线程无法正常退出(虽然这不是当前挂起的直接原因,但属于潜在风险)。

修复方案

修改Worker类核心逻辑

将PushBack()操作放到Worker锁的保护范围内,确保队列大小判断与入队操作是原子性的;同时将done改为原子变量,保证线程可见性:

#pragma once
#include <queue>
#include <mutex>
#include <functional>
#include <atomic>

class Worker {
public:
    ~Worker() {
        Send([this] { done.store(true); });
        thd.join();
    }

    Worker(size_t max_num_tasks) : max_num_tasks_(max_num_tasks), done(false), thd([this] {
        while (!done.load()) {
            auto task = mq.PopFront();
            task();
            // 任务执行完成后,通知等待的生产者队列有空间了
            std::lock_guard<std::mutex> lock(m_);
            cv_.notify_all();
        }
        })
    { }

    void Send(std::function<void()>&& m) {
        std::unique_lock<std::mutex> lock(m_);
        // 等待队列有剩余空间
        cv_.wait(lock, [&] {return (mq.Size() < max_num_tasks_); });
        // 入队操作在锁保护下,确保判断与入队的原子性
        mq.PushBack(std::move(m));
    }
private:
    std::atomic<bool> done; // 原子变量保证线程间可见性
    size_t max_num_tasks_;
    ThreadSafeQueue<std::function<void()>> mq;
    std::thread thd;
    std::mutex m_;
    std::condition_variable cv_;
};

修复说明

  1. 原子性入队:将mq.PushBack()移到m_锁的保护范围内,确保队列大小判断和入队操作不会被其他线程打断,彻底避免多个生产者同时入队导致队列超限。
  2. 线程可见性:done改为std::atomic<bool>,确保消费者线程能立刻看到生产者设置的终止信号。
  3. 正确通知时机:消费者完成任务后,持有m_锁并通知cv_,让等待的生产者重新判断队列空间。

内容的提问来源于stack exchange,提问作者Jens

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 03:17:09