基于Boost.Fibers实现生产者/消费者模型:结合Promise/Future完成通知
基于Boost.Fibers实现带Promise/Future完成信号的生产者消费者模型
我最近在尝试用Boost.Fibers实现生产者/消费者模型,参考官方示例后觉得用channels来传递任务是个合理的方案。不过因为需要通过std::promise/std::future来发送完成信号,所以得对示例做些修改。下面是我写的一段仅用于传递完成信号的基础代码框架:
#include <thread> #include <boost/fibers/channel.hpp> #include <boost/fibers/fiber.hpp> #include <future> // 定义任务类型,可根据实际需求扩展 struct task {}; struct fiber_worker { fiber_worker() { wthread = std::thread([self{this}]() { // 创建4个worker fiber处理任务 for (int i = 0; i < 4; ++i) { boost::fibers::fiber{ [self]() { task tsk; // 循环从channel取任务,直到channel关闭 while (boost::fibers::channel_op_status::closed != self->task_channel.pop(tsk)) { // 这里添加具体任务处理逻辑 // ... // 任务完成后可更新计数,触发完成信号 if (--self->remaining_tasks == 0) { self->completion_promise.set_value(); } } }}.detach(); // 让fiber在后台运行 } }); } ~fiber_worker() { task_channel.close(); // 关闭channel,通知worker停止 if (wthread.joinable()) { wthread.join(); } } // 提供给生产者的任务推送接口 void push_task(const task& tsk) { task_channel.push(tsk); ++remaining_tasks; } // 获取用于等待完成的future std::future<void> get_completion_future() { return completion_promise.get_future(); } std::thread wthread; boost::fibers::unbuffered_channel<task> task_channel; std::promise<void> completion_promise; std::atomic<int> remaining_tasks{0}; };
这段代码的核心设计思路:
- 用
std::thread承载Boost.Fibers的调度上下文,在该线程内创建多个worker fiber并行处理任务 - 借助
boost::fibers::unbuffered_channel实现任务的生产者-消费者传递,保证线程安全的任务分发 - 通过
std::promise/std::future实现完成信号的传递,当所有任务处理完成时,触发promise通知主线程
需要注意的是,这段代码是基础框架,实际使用时还需要根据业务场景调整任务类型、任务数量的管理逻辑,比如批量任务的计数初始化等细节。
内容的提问来源于stack exchange,提问作者kreuzerkrieg
相关产品推荐
相关产品推荐

