ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

【C++ 标准项目】发布订阅消息队列 (篇五):C++11 异步编程,手写线程池,async/promise/packaged_task/future 实战

【C++ 标准项目】发布订阅消息队列 (篇五):C++11 异步编程,手写线程池,async/promise/packaged_task/future 实战 C11 异步操作实现线程池1.std::futurea.介绍链接如下cplusplus.com/reference/future/future/std::future 是 C11 标准库中的一个模板类它表示一个异步操作的结果。当我们在多线程编程中使用异步任务时std::future 可以帮助我们在需要的时候获取任务的执行结果。std::future 的一个重要特性是能够阻塞当前线程直到异步操作完成从而确保我们在获取结果时不会遇到未完成的操作。b.应用场景异步任务当我们需要在后台执行一些耗时操作时如网络请求或计算密集型任务等std::future 可以用来表示这些异步任务的结果。通过将任务与主线程分离我们可以实现任务的并行处理从而提高程序的执行效率并发控制在多线程编程中我们可能需要等待某些任务完成后才能继续执行其他操作。通过使用 std::future我们可以实现线程之间的同步确保任务完成后再获取结果并继续执行后续操作结果获取std::future 提供了一种安全的方式来获取异步任务的结果。我们可以使用std::future::get()函数来获取任务的结果此函数会阻塞当前线程直到异步操作完成。这样在调用 get()函数时我们可以确保已经获取到了所需的结果c.用法示例使用 std::async 关联异步任务std::async 是一种将任务与 std::future 关联的简单方法。它创建并运行一个异步任务并返回一个与该任务结果关联的 std::future 对象。默认情况下std::async 是否启动一个新线程或者在等待 future 时任务是否同步运行都取决于你给的 参数。这个参数为 std::launch 类型 std::launch::deferred 表明该函数会被延迟调用直到在 future 上调用 get()或者 wait()才会开始执行任务 std::launch::async 表明函数会在自己创建的线程上运行 std::launch::deferred | std::launch::async 内部通过系统等条件自动选择策略来看代码运行结果使用 std::packaged_task 和 std::future 配合std::packaged_task 就是将任务和 std::future 绑定在一起的模板是一种对任务的封装。我们可以通过 std::packaged_task 对象获取任务相关联的 std::future 对象通过调用 get_future()方法获得。std::packaged_task 的模板参数是函数签名。可以把 std::future 和 std::async 看成是分开的 而 std::packaged_task 则是一个整体。来看代码运行结果使用 std::promise 和 std::future 配合std::promise 提供了一种设置值的方式它可以在设置之后通过相关联的 std::future 对象进行读取。换种说法就是之前说过 std::future 可以读取一个异步函数的返回值了 但是要等待就绪而 std::promise 就提供一种 方式手动让 std::future 就绪来看代码2.c11 线程池实现来看代码运行结果代码解析如下a.整体结构简答说就是push负责交任务工作线程负责取任务执行future负责传结果。b.核心流程push()的四步templatetypename F, typename ...Args auto push(F func, Args ...args) - std::futuredecltype(func(args...)) { // 1. 把函数 参数绑成一个整体并推断返回类型 using return_type decltype(func(args...)); auto tmp_func std::bind(std::forwardF(func), std::forwardArgs(args)...); // 2. 用 packaged_task 包装它并用 shared_ptr 持有 auto task std::make_sharedstd::packaged_taskreturn_type()(tmp_func); std::futurereturn_type fu task-get_future(); // 取出配对凭证 // 3. 构造一个只捕获 shared_ptr 的 lambda塞进队列 { std::unique_lockstd::mutex lock(_mutex); _taskpool.push_back([task](){ (*task)(); }); _cv.notify_one(); // 叫醒一个工作线程 } return fu; // 4. 把凭证交给调用方 }entry()的三步void entry() { while(!_stop) { std::vectorfunctor tmp_taskpool; { std::unique_lockstd::mutex lock(_mutex); // 1. 等任务池非空 或 要停了 _cv.wait(lock, [this](){ return _stop || !_taskpool.empty(); }); // 2. 把整批任务 swap 到本地 tmp_taskpool.swap(_taskpool); } // 3. 【锁已经释放】逐个执行 for (auto task : tmp_taskpool) task(); } }c.四个关键设计1. 为什么必须用shared_ptr包packaged_task这是整段代码最关键的一处很多人写线程池就卡在这里std::vectorfunctor _taskpool; // functor std::functionvoid()std::function要求存进去的可调用对象是可拷贝的。但std::packaged_task是move-only只能移动不能拷贝std::packaged_taskint() t1(...); std::packaged_taskint() t2 t1; // 错误 拷贝构造被删除 _taskpool.push_back(task); // 错误 packaged_task 无法放进 std::function解法用shared_ptr包一层。shared_ptr是可拷贝的lambda 捕获它也只是增加引用计数auto task std::make_sharedstd::packaged_taskreturn_type()(tmp_func); _taskpool.push_back([task](){ (*task)(); }); // 正确 shared_ptr 可拷贝同时它还顺手解决了生命周期问题task被 lambda 捕获 -- 引用计数 1任务执行完lambda 被销毁 -- 引用计数 -1 -- 自动释放不用手动delete也不用担心任务还在队列里对象就没了小结shared_ptr 在这里干了两件事——让 move-only 的对象能被 std::function 持有以及自动管理生命周期。2.为什么要把任务swap到本地再执行{ std::unique_lockstd::mutex lock(_mutex); tmp_taskpool.swap(_taskpool); // - 只做这一件事然后立刻解锁 } for (auto task : tmp_taskpool) task(); // - 在锁【外面】执行如果直接在锁里执行任务会怎样// 反例 { std::unique_lockstd::mutex lock(_mutex); while(!_taskpool.empty()) { _taskpool.front()(); // 执行任务可能耗时很久 _taskpool.pop_front(); } }后果某个任务跑了 5 秒这 5 秒内其他线程一个都 push 不进来锁被占着—— 线程池变成了串行池。swap 的妙处加锁的时间只有交换两个 vector 指针那么短O(1)任务的实际执行完全不持锁。3. 为什么wait必须带谓词_cv.wait(lock, [this](){ return _stop || !_taskpool.empty(); });等价于while (!(_stop || !_taskpool.empty())) { // --- 必须用 while不是 if _cv.wait(lock); }两个原因原因说明虚假唤醒条件变量可能无缘无故醒来while会重新检查条件if会带着错误的状态往下跑丢失唤醒如果不持锁检查条件可能错过notify_one—— 带谓词的wait会自动处理这件事不带谓词的写法是错的_cv.wait(lock); // 错误的 醒来后不检查条件可能任务池是空的4. 为什么_stop要用std::atomicbool_stop被多个线程同时读写工作线程读、stop()写必须保证可见性和原子性std::atomicbool _stop; // 正确 bool _stop; // 错误 数据竞争工作线程可能永远看不到变化d.五个有待改进得问题问题 1main里push完立刻get等于没有并发for(int i 0; i 10; i) { std::futureint fu pool.push(Add, 11, i); std::cout fu.get() std::endl; // -- 立刻阻塞等待 }每个任务都是提交 -- 等它做完 -- 再提交下一个线程池退化成单线程顺序执行。 想真正看到并发应该先把任务都提交完再统一取结果std::vectorstd::futureint fus; for (int i 0; i 10; i) fus.push_back(pool.push(Add, 11, i)); // 先全部提交 for (auto f : fus) std::cout f.get() std::endl; // 再统一取这样 10 个任务能同时被多个工作线程并行处理。问题 2std::bind会拷贝参数丢失移动语义auto tmp_func std::bind(std::forwardF(func), std::forwardArgs(args)...);std::bind存参数时会decay退化——std::forward白写了传进来的右值会被拷贝而不是移动move-only 的参数比如unique_ptr、packaged_task根本传不进去直接编译报错。C14 起有更好的写法lambda 初始化捕获auto task std::make_sharedstd::packaged_taskreturn_type()( [func std::forwardF(func), args std::make_tuple(std::forwardArgs(args)...)]() mutable { return std::apply(func, std::move(args)); });std::apply 是 C17。你编译用 -stdc11所以 std::bind 是这个版本的标准做法不算错但要知道它的局限问题 3stop()可能丢弃剩余任务while(!_stop) // -- 工作线程在这里检查 { ... _cv.wait(lock, [this](){ return _stop || !_taskpool.empty(); }); tmp_taskpool.swap(_taskpool); ... }如果 _stop 在工作线程执行完一批任务、准备回到while判断时被置位而队列里还有新任务 ——它会直接退出循环那些任务永远不会被执行。这个丢任务是静默的不会报错。想保证停之前把活干完需要在退出前再清一次队列void entry() { for (;;) { std::vectorfunctor tmp; { std::unique_lockstd::mutex lock(_mutex); _cv.wait(lock, [this]{ return _stop || !_taskpool.empty(); }); tmp.swap(_taskpool); // 只有在要停 且 队列已空时才退出 if (tmp.empty() _stop) return; } for (auto t : tmp) t(); } }不过这里得 main 里是push 一个 get 一个退出时队列本来就是空的所以实际不会触发。问题 4stop()的 check-then-set 不是原子的if(_stop true) return; // -- 检查 _stop true; // -- 设置两个线程同时调stop()可能都通过检查。更好的写法是if (_stop.exchange(true)) return; // 一步完成检查并设置问题 5tmp_taskpool每轮循环都重建while(!_stop) { std::vectorfunctor tmp_taskpool; // --- 每次都重新分配内存把它移到循环外面就能复用已分配的内存swap 之后 vector 的容量会保留在另一方std::vectorfunctor tmp_taskpool; // 放外面 while(!_stop) { { ... tmp_taskpool.swap(_taskpool); } for (auto task : tmp_taskpool) task(); tmp_taskpool.clear(); }e.运行结果因为main是push 一个、get 一个所以顺序输出——这反而证明了future.get()的同步作用每次 get 都会阻塞到对应任务执行完。想验证多线程真的在跑可以在Add里打印线程 IDint Add(int a, int b) { std::cout 线程 std::this_thread::get_id() 处理任务\n; return a b; }然后用上面问题 1里那个先全部提交、再统一取结果的写法——你会看到任务被不同线程接手。f.小结设计价值shared_ptrpackaged_task让 move-only 对象能被std::function持有 自动管理生命周期swap到本地再执行锁只保护队列操作不保护任务执行避免串行化wait带谓词防虚假唤醒、防丢失唤醒atomicbool _stop多线程可见性析构调stop()RAII忘了手动停也不会泄漏线程
返回列表