如何利用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
相关产品推荐
相关产品推荐

