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

基于Mutex与Condition Variable的生产者消费者线程同步实现咨询

双向通知的生产者消费者实现方案

嘿,这个生产者和消费者互相通知等待的问题,确实是多线程编程里容易卡壳的点。其实核心就是用条件变量+状态标记来实现双向的等待与唤醒,我给你写个完整的可运行示例,再拆解关键逻辑。

核心思路

咱们的需求是:

  • 消费者处理完当前数据 → 通知生产者生成新数据
  • 生产者生成好新数据 → 通知消费者来处理

所以需要两个关键的同步点:一个让生产者等消费者处理完,另一个让消费者等生产者准备好。用两个条件变量分别对应这两个场景,再配合一个状态标记来避免虚假唤醒。

完整代码示例

#include <iostream>
#include <thread>
#include <mutex>
#include <condition_variable>
#include <vector>
#include <atomic>
#include <chrono>

// 全局同步变量(实际项目里建议封装成类,这里为了演示简化)
std::mutex m;
std::condition_variable cv_producer;  // 用于通知生产者:可以生产新数据了
std::condition_variable cv_consumer;  // 用于通知消费者:有新数据可以处理了
std::vector<int> task_queue;          // 任务队列
bool is_processed = true;             // 标记:当前队列数据是否已被处理
std::atomic<bool> stop_running = false;  // 线程终止标记,原子变量保证线程可见性

// 生产者线程函数
void producer() {
    int task_id = 0;
    while (!stop_running) {
        std::unique_lock<std::mutex> lock(m);
        // 等待消费者处理完上一批数据,或者收到终止信号
        cv_producer.wait(lock, []{ return is_processed || stop_running; });

        if (stop_running) break;

        // 模拟生产数据:这里生成一个递增的任务ID
        task_queue.clear();
        task_queue.push_back(++task_id);
        std::cout << "[生产者] 生成新任务: " << task_id << std::endl;

        // 标记数据未处理,通知消费者可以开始处理
        is_processed = false;
        cv_consumer.notify_one();
    }
}

// 消费者线程函数
void consumer() {
    while (!stop_running) {
        std::unique_lock<std::mutex> lock(m);
        // 等待生产者准备好新数据,或者收到终止信号
        cv_consumer.wait(lock, []{ return !is_processed || stop_running; });

        if (stop_running) break;

        // 模拟处理数据:打印任务ID
        for (int id : task_queue) {
            std::cout << "[消费者] 处理任务: " << id << std::endl;
        }
        task_queue.clear();

        // 标记数据已处理,通知生产者可以生成新数据
        is_processed = true;
        cv_producer.notify_one();

        // 模拟处理耗时(可选,实际根据业务调整)
        std::this_thread::sleep_for(std::chrono::milliseconds(600));
    }
}

int main() {
    // 启动生产者和消费者线程
    std::thread prod_thread(producer);
    std::thread cons_thread(consumer);

    // 让线程运行5秒后终止
    std::this_thread::sleep_for(std::chrono::seconds(5));
    stop_running = true;

    // 唤醒所有等待的线程,确保它们能正常退出循环
    cv_producer.notify_one();
    cv_consumer.notify_one();

    // 等待线程结束
    prod_thread.join();
    cons_thread.join();

    std::cout << "所有线程已终止" << std::endl;
    return 0;
}

关键细节解释

  • 双条件变量:cv_producer负责唤醒生产者,cv_consumer负责唤醒消费者,分工明确,避免混淆。
  • 状态标记+wait谓词:is_processed标记当前数据的处理状态,配合wait的第二个谓词参数,能有效避免虚假唤醒(系统会自发唤醒条件变量的情况),确保线程只在真正满足条件时才继续执行。
  • 原子终止标记:stop_running用std::atomic修饰,保证在多线程环境下的可见性,避免线程一直阻塞在wait里无法退出。
  • 锁的范围:用std::unique_lock而不是std::lock_guard,因为wait会自动释放锁,唤醒后重新获取锁,这是条件变量的必备操作。

扩展建议

如果你的场景是批量生产/消费,或者有多个生产者/消费者,只需要调整状态标记(比如用计数代替布尔值),或者在通知时用notify_all()代替notify_one()即可,核心逻辑是通用的。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 03:55:21