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

如何利用Barrier同步生产者线程与三个消费者线程?

我来帮你搞定这个生产者-消费者的同步逻辑,确保生产者能等三个消费者都处理完任务后,再去处理结果。先给你梳理清楚实现思路,再上完整代码:

核心同步思路

这里用boost::barrier是非常合适的——你已经初始化了barrier(4),正好对应1个生产者+3个消费者的总线程数。每次生产者生成完任务,消费者处理完任务后,所有线程都在barrier处等待,等全员到齐后,barrier会自动放行,此时生产者就可以安全地处理上一轮的任务结果,然后开启下一轮的任务生产。

完整代码实现
#include <condition_variable>
#include <mutex>
#include <boost/thread/barrier.hpp>
#include <vector>
#include <queue>
#include <iostream>
#include <thread>

// 全局同步对象(实际项目里建议封装成类,这里为了简化用全局)
std::condition_variable cond;
std::mutex mutex;
boost::barrier barrier(4);
std::queue<int> task_queue; // 任务队列,这里用int代表任务,你可以换成自定义任务类型
std::vector<int> results;   // 存储消费者处理后的结果
bool stop_flag = false;     // 停止线程的标志

// 消费者线程函数
void consumer(int id) {
    while (true) {
        std::unique_lock<std::mutex> lock(mutex);
        // 等待生产者发布任务或者停止信号
        cond.wait(lock, []{ return !task_queue.empty() || stop_flag; });
        
        if (stop_flag && task_queue.empty()) {
            // 收到停止信号且任务为空,退出线程
            lock.unlock();
            barrier.wait(); // 最后一次等待,让生产者知道所有消费者都退出了
            break;
        }

        // 取出任务
        int task = task_queue.front();
        task_queue.pop();
        lock.unlock();

        // 模拟任务处理(这里换成你的实际业务逻辑)
        std::cout << "Consumer " << id << " processing task " << task << std::endl;
        int result = task * 2; // 示例处理逻辑:任务值翻倍
        lock.lock();
        results.push_back(result);
        lock.unlock();

        // 处理完任务,等待所有线程到达barrier
        barrier.wait();
    }
}

// 生产者线程函数
void producer() {
    for (int round = 1; round <= 3; ++round) { // 模拟生产3轮任务
        std::unique_lock<std::mutex> lock(mutex);
        // 生成任务(这里每次生成3个任务,对应3个消费者)
        for (int i = 0; i < 3; ++i) {
            int task = round * 10 + i;
            task_queue.push(task);
            std::cout << "Producer generated task " << task << std::endl;
        }
        lock.unlock();

        // 通知所有消费者可以取任务了
        cond.notify_all();

        // 等待所有消费者处理完当前轮的任务
        barrier.wait();

        // 所有消费者都处理完了,现在处理结果
        std::unique_lock<std::mutex> result_lock(mutex);
        std::cout << "\nRound " << round << " results: ";
        for (int res : results) {
            std::cout << res << " ";
        }
        std::cout << "\n-------------------------\n";
        results.clear(); // 清空结果,准备下一轮
        result_lock.unlock();
    }

    // 所有任务轮次结束,通知消费者停止
    std::unique_lock<std::mutex> lock(mutex);
    stop_flag = true;
    lock.unlock();
    cond.notify_all();

    // 等待所有消费者退出
    barrier.wait();
    std::cout << "Producer finished, all consumers exited.\n";
}

int main() {
    std::vector<std::thread> workers;

    // 创建3个消费者线程
    for (int i = 0; i < 3; ++i) {
        workers.emplace_back(consumer, i+1);
    }
    // 创建生产者线程
    std::thread producer_thread(producer);

    // 等待所有线程结束
    producer_thread.join();
    for (auto& t : workers) {
        t.join();
    }

    return 0;
}
关键细节解释
  • barrier的作用:每一轮任务处理完成后,生产者和消费者都会调用barrier.wait(),只有当4个线程都调用这个方法时,所有线程才会继续往下执行。这就保证了生产者必须等所有消费者处理完当前任务,才能去处理结果。
  • 条件变量的使用:消费者通过cond.wait()阻塞等待任务,生产者生成任务后用cond.notify_all()唤醒所有消费者,避免消费者空轮询浪费CPU。
  • 线程安全:所有对任务队列、结果容器的访问都被std::mutex保护,避免多线程竞态条件导致的数据错乱。
  • 优雅退出:通过stop_flag让生产者可以通知消费者停止,最后一轮barrier等待确保生产者能确认所有消费者都已退出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 03:38:04