ARTICLE DETAIL

资讯详情

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

C++11实现工作窃取线程池

C++11实现工作窃取线程池

一、线程池概述

1.1 线程池概念

线程池技术通过在系统中预先创建一定数量的线程,当任务请求到来时从线程池中分配一个预先创建的线程去处理,线程在处理完任务之后并不会销毁,而是把线程还到线程池中,继续为后续的任务提供服务。

线程池的特点:

线程复用:线程池会在内部维护一定数量的线程,并在需要时重复使用这些线程来执行任务,避免频繁地创建和销毁线程,从而提高性能和效率。

控制并发性:对于多核处理器,由于多线程被分配到多个处理器中,提高并行处理效率。

任务队列:当线程池中的线程已经全部被占用时,新提交的任务会被放入一个任务队列中进行排队等待执行,排队机制可以根据具体线程池实现,选择不同的队列类型,如有界队列或无界队列。

开发环境:

window: vs2019

Linux: g++ 要求g++版本能够支持C++11以上

1.2 按应用场景分类

1. FixedThreadPool

固定线程池:线程池中的线程数量固定,这些线程一直存在,不会随任务的增加或减少而动态调整,超出的任务会在队列中等待。

使用场景:任务量比较固定但耗时较长的任务。

2. CachedThreadPool

缓存线程池:可根据需要创建新线程的线程池,如果新任务到达,但线程池中没有可用线程,则创建一个新线程并添加到池中,如果有被使用完但是还没有销毁的线程,就复用该线程。

使用场景:任务量大但耗时少的任务。

3. SingleThreadPool

单线程池:使用唯一的工作线程来执行任务,保证所有任务按照指定顺序(FIFO,LIFO,优先级)执行。

使用场景:多个任务顺序执行(FIFO,优先级)。

4. WorkStealingPool

工作窃取线程池:创建一个拥有多个任务队列(以便减少连接数)的线程池。

使用场景:高并发下的负载均衡。

5. ScheduledThreadPool

计划线程池(定时线程池,调度线程池)

使用场景:定时以及周期性执行任务。

1.3 线程池模式

线程池模式一般分为两种:L/F领导者与跟随者模式,HS/HA半同步/半异步模式。

1.4 半同步/半异步模式分析

1. 同步服务层,它处理来自上层的任务请求,上层的请求可能是并发的,这些请求不是马上就会被处理,而是将这些任务放到一个同步队列中,等待处理。

2. 同步排队层,来自上层的任务请求都会加到排队层中等待处理。

3. 异步服务层:这一层会有多个线程同时处理排队层中的任务,异步服务层从同步排队层中取出任务并行的处理。

1.5 线程池实现的关键技术分析

线程池有两个活动过程,一个是往同步队列中添加任务的过程,另一个是从同步队列中取任务的过程。

半同步半异步线程池活动图

二、WorkStealingPool的实现

2.1 需求

工作窃取算法:WorkStealingPool采用了工作窃取算法,具体来说就是当某个线程执行完自己队列中的任务后,会从其他线程中“偷取”任务来执行。这种算法可以提高线程利用率,减少线程之间的竞争,以及减少线程的等待时间。

WorkStealingPool可以设定多个工作线程,每个工作线程都有一个自己的任务队列,每个线程在执行任务时会首先从自己的队列中获取任务,如果自己队列为空,则从其他线程的队列中获取任务。这种设计可以充分发挥多核处理器的并行能力,提高整体的任务处理效率。

2.2 SyncQueue的设计和实现

该分段分桶阻塞同步队列专为工作窃取线程池设计,通过vector<list<T>>将任务分散至 N 个独立子队列,使每个线程绑定专属队列以消除全局锁竞争,并利用双条件变量与移动语义保障高吞吐下的低延迟。

其内置的超时等待、安全停机机制及标准化返回码确保了系统健壮性与优雅退出,而多队列架构天然支持本地优先消费与跨队列窃取逻辑,在维持缓存局部性的同时实现动态负载均衡,完美解决高并发场景下的性能瓶颈与任务倾斜问题。

template<class T> class SyncQueue { private: // 分桶队列数组,实现多队列打散锁竞争 std::vector<std::list<T>> m_taskQueues; size_t m_bucketSize; // 桶数量(队列个数,对应线程池线程数) size_t m_maxSize; // 单个队列最大容量 mutable std::mutex m_mutex; std::condition_variable m_notEmpty; // 消费者等待:队列非空 std::condition_variable m_notFull; // 生产者等待:队列未满 size_t m_waitTime; // 条件变量超时等待时长(秒) bool m_needStop; // 队列停止标记,用于析构/关闭 // 判断指定编号的队列是否已满 bool IsFull(const int index) const; // 判断指定编号的队列是否为空 bool IsEmpty(const int index) const; // 底层添加任务模板实现,支持左值/右值引用 template<class F> int Add(F&& task,const int index); public: // 构造函数:bucketSize=分桶数量,maxSize单队列上限,timeout等待超时 SyncQueue(int bucketsize,int maxsize = 200, size_t timeout = 1); // 析构函数 ~SyncQueue(); // 左值版本入队 int Put(const T& task,const int index); // 右值版本入队(移动语义优化) int Put(T&& task,const int index); // 批量取出整个队列数据,移动取出 int Take(std::list<T>& list,const int index); // 单个取出队首元素 int Take(T& task,const int index); // 停止队列,唤醒所有阻塞的生产者、消费者线程 void Stop(); };

Add函数

该队列通过全局互斥锁保障线程安全,在Add核心逻辑中集成停止标记检测,若处于停机状态则直接返回 2 并丢弃任务。针对生产者阻塞,采用wait_for实现超时等待,满队时休眠并在超时后返回 1,避免永久卡死。

入队过程利用完美转发技术区分左值拷贝与右值移动,显著降低内存复制开销,并在成功插入后唤醒所有消费者条件变量以通知新任务到达。上层Put接口通过两个重载分别处理左值与右值,内部统一调用Add,从而对外提供简洁且高效的差异化入队体验。

template<class F> int SyncQueue<T>::Add(F&& task,const int index) { std::unique_lock<std::mutex> locker(m_mutex); // 队列满则限时等待,谓词:未停止 且 队列已满继续阻塞 bool waitret = m_notFull.wait_for(locker, std::chrono::seconds(m_waitTime), [=] { return m_needStop || !IsFull(index); }); // 等待超时,入队失败 if (!waitret) { return 1; } // 队列已经停止,拒绝入队 if (m_needStop) { return 2; } // 任务转发入队,右值直接移动,左值拷贝 m_taskQueues[index].push_back(std::forward<F>(task)); // 唤醒消费者:有新任务了 m_notEmpty.notify_all(); return 0; }

Take函数

消费者在获取锁后,若队列为空则通过m_notEmpty条件变量进入限时阻塞等待;超时返回 1 以触发上层工作窃取逻辑,检测至停机标记则返回 2 促使线程优雅退出。

一旦成功唤醒,即取出队列头部任务并弹出节点,随后通知生产者条件变量释放空间,确保生产端能继续入队,从而在保证线程安全的同时实现高效的负载均衡与资源流转。

template<class T> int SyncQueue<T>::Take(T& task,const int index) { std::unique_lock<std::mutex> locker(m_mutex); // 队列为空限时阻塞,停机或者有任务则解除阻塞 bool waitret = m_notEmpty.wait_for(locker, std::chrono::seconds(m_waitTime), [this, index]()->bool { return m_needStop || !IsEmpty(index); }); // 等待超时无任务 if (!waitret) { return 1; } // 队列停止,退出消费 if (m_needStop) { return 2; } // 取出队首元素 task = m_taskQueues[index].front(); m_taskQueues[index].pop_front(); // 唤醒生产者,队列腾出空位 m_notFull.notify_all(); return 0; }

Stop函数

关闭流程首先上锁遍历所有分桶队列,阻塞等待直至所有任务被完全消费,确保不丢失任何未执行作业。随后设置全局停止标志m_needStop=true并广播唤醒全部条件变量:阻塞的生产者检测到该标志后立即终止入队操作,而阻塞的消费者则据此退出消费循环。

这一机制通过有序的状态同步与全量唤醒,实现了线程池的平滑关闭,既保证了数据完整性,又避免了残留阻塞线程导致的资源泄漏。

template<class T> void SyncQueue<T>::Stop() { std::unique_lock<std::mutex> locker(m_mutex); // 等待所有桶队列任务消费完毕,再标记停止 for(int i = 0;i<m_bucketSize;++i) { while (!m_needStop && !IsEmpty(i)) { m_notFull.wait(locker); } } // 置位停止标志 m_needStop = true; // 唤醒所有等待的消费者、生产者线程 m_notEmpty.notify_all(); m_notFull.notify_all(); }

2.3 WorkStealingPool的设计和实现

该工作窃取线程池依托前文SyncQueue分段分桶阻塞队列搭建,将工作线程与队列一一绑定,依托本地优先消费、跨队列任务窃取实现负载均衡,解决传统固定队列线程池任务倾斜、长尾阻塞、全局锁竞争严重的痛点。

template<class T> class WorkStealingPool { private: SyncQueue<T> m_taskQueue; // 底层分桶同步队列 std::vector<std::thread> m_workers;// 工作线程数组 size_t m_threadNum; // 线程总数,等价于分桶数量 std::atomic<bool> m_poolStop; // 线程池停止原子标记 // 线程主循环:自身队列消费+空闲窃取逻辑 void WorkerRun(int threadIdx); // 空闲线程从其他队列窃取任务 bool TryStealTask(T& task, int selfIdx); public: // 构造:指定线程数、单队列最大容量、条件变量超时时间 WorkStealingPool(size_t threadCount, size_t queueMaxSize = 200, size_t waitSec = 1); // 析构触发池销毁 ~WorkStealingPool(); // 任务提交,投递至当前线程绑定队列 bool SubmitTask(const T& task, int threadIdx); // 右值版本任务提交,移动语义优化 bool SubmitTask(T&& task, int threadIdx); // 整体关闭线程池,同步等待任务执行完毕 void Shutdown(); };

WorkerRun 线程主循环函数

线程以自身编号绑定专属队列持续循环消费,本地队列拉取任务超时后,触发跨队列窃取逻辑;检测到线程池停机标记直接结束线程,是整个窃取调度的核心驱动逻辑。

线程优先调用SyncQueue::Take读取自身绑定队列任务,成功则直接执行业务逻辑;若队列等待超时(返回码 1),代表本地无任务,调用TryStealTask尝试窃取其他队列任务;当线程池触发关闭(返回码 2),跳出循环完成线程退出。循环结构持续复用线程,规避频繁创建销毁线程的开销,同时依靠队列超时阻塞避免 CPU 空轮询。

template<class T> void WorkStealingPool<T>::WorkerRun(int threadIdx) { T task; while (!m_poolStop) { int ret = m_taskQueue.Take(task, threadIdx); if (ret == 0) { // 正常取出本地任务执行 task(); } else if (ret == 1) { // 本地队列为空超时,尝试窃取任务 if (TryStealTask(task, threadIdx)) { task(); } } else if (ret == 2) { // 队列停机,终止线程循环 break; } } }

TryStealTask 任务窃取函数

针对本地队列空闲场景,轮询遍历其余所有分桶队列,从远端队列拉取任务实现负载均衡,规避单队列任务堆积导致的线程算力闲置。

窃取逻辑会跳过自身队列,依次尝试从其他队列调用单任务Take接口拉取任务,一旦窃取成功立刻返回执行;全部队列遍历完毕均无可用任务则返回false,线程回到阻塞等待状态。该逻辑保证忙碌队列的任务被空闲线程分摊,抹平各线程任务量差值,相比统一任务分发大幅降低锁冲突。

template<class T> bool WorkStealingPool<T>::TryStealTask(T& task, int selfIdx) { for (int i = 0; i < m_threadNum; ++i) { if (i == selfIdx) continue; int ret = m_taskQueue.Take(task, i); if (ret == 0) { return true; } } return false; }

Shutdown 线程池停机销毁函数

串联SyncQueue的停机能力与线程等待回收,先标记线程池停止状态,调用队列Stop阻塞等待全部存量任务执行完成,再join所有工作线程,保证任务不丢失、线程资源完整回收,无残留阻塞线程与僵尸线程。

首先设置原子停机标识,调用底层队列Stop阻塞等待所有子队列任务消费完毕,之后遍历所有工作线程执行join阻塞回收,确保主线程等待所有工作线程正常退出后再释放资源,实现业务任务无损关闭,适配服务平滑下线场景。

template<class T> void WorkStealingPool<T>::Shutdown() { m_poolStop = true; // 同步停止队列,等待存量任务执行完成 m_taskQueue.Stop(); // 回收全部工作线程 for (auto& th : m_workers) { if (th.joinable()) { th.join(); } } }

三、WorkStealingPool的测试

本次测试旨在全面验证工作窃取线程池在高并发场景下的稳定性与负载均衡能力,核心聚焦于多生产者无阻塞提交、任务窃取机制生效及资源安全回收三大维度。

通过设计包含休眠逻辑的加法任务add()模拟 CPU 密集型计算,并利用随机时长制造任务执行差异,人为构建负载倾斜环境以触发工作窃取行为。生产者线程add_a()负责并发提交任务并阻塞等待结果,以此检验线程池在海量请求涌入时的响应效率及内部调度算法是否能有效缓解单队列堆积,确保忙碌线程的任务能被空闲线程及时“窃取”执行。主函数逻辑构建了高强度的压力测试场景,创建包含 2000 个任务的线程池实例,并启动同等数量的生产者线程进行并发投递,模拟极端高负载下的系统表现。

#include "workStealingPool.hpp" #include <thread> #include <iostream> #include <vector> #include <chrono> #include <climits> #include <system_error> #include <cstdlib> #include <ctime> using namespace std; // 线程池执行任务:加法计算,可开启sleep模拟耗时任务 int add(int a, int b, int s) { clog << "add begin ..." << endl; int c = a + b; // 取消注释开启任务耗时差异,用于验证工作窃取负载均衡 // std::this_thread::sleep_for(std::chrono::seconds(s)); clog << "add end .. " << endl; return c; } // 生产者线程函数:向线程池提交异步任务 void add_a(int ch, int x, int y, WorkStealingPool<int>& mypool) { // 随机0~9秒任务耗时,制造任务长短差异,触发工作窃取 auto r = mypool.submit(add, x, y, rand() % 10); // get()阻塞等待任务执行完成,获取返回值 cout << "add_" << ch << " " << r.get() << endl; } int main() { srand((unsigned int)time(nullptr)); // 随机数种子初始化 const int task_producer_num = 2000; // 生产者线程总数,并发提交任务 WorkStealingPool<int> mypool; // 初始化工作窃取线程池 std::vector<std::thread> producer_threads(task_producer_num); int create_cnt = 0; // 成功创建的生产者线程计数 // 批量创建生产者线程,并发投递任务 for (int i = 0; i < task_producer_num; ++i) { try { producer_threads[i] = std::thread(add_a, i, i + 20, i + 10, std::ref(mypool)); create_cnt = i; } catch (std::system_error& e) { // 线程资源不足异常捕获,避免程序直接崩溃 cout << "Thread create error: " << e.what() << endl; break; } } // 阻塞等待所有生产者线程执行完毕(全部任务提交完成) for (int i = 0; i <= create_cnt; ++i) { if (producer_threads[i].joinable()) { producer_threads[i].join(); } } cout << "Success create producer thread num: " << create_cnt + 1 << endl; // 线程池优雅关闭,等待池内全部任务执行完成、工作线程回收 mypool.Shutdown(); return 0; }

四、线程池进阶拓展

4.1 三种线程池对比

对比维度固定式线程池 (FixedThreadPool)缓存式线程池 (CachedThreadPool)工作窃取线程池 (WorkStealingPool)
线程数量固定初始化,全程线程数不变动态伸缩,任务暴涨新建线程,空闲超时销毁固定线程数,线程绑定独立任务队列
任务队列结构全局单阻塞队列无界同步队列多分段分桶队列
负载均衡方式主线程统一分发,容易出现队列任务堆积、长尾阻塞按需扩容线程扛峰值,高峰过后大量空闲线程回收线程本地优先消费,空闲线程主动窃取繁忙队列任务,均衡负载
锁竞争程度全局队列锁,高并发入队锁冲突激烈单队列锁竞争同样明显分桶打散锁,锁粒度小,并发吞吐高
适用场景任务量平稳、负载稳定的长驻后台任务短时脉冲峰值任务、任务执行极快的场景CPU 密集型并行计算、任务耗时不均、存在长尾任务的场景

4.2 使用场景

WorkStealingPool 适用于以下场景:

  1. 任务分解型应用:当一个任务需要被分解成多个子任务进行并行处理时,WorkStealingPool 可以自动管理任务的分配和调度,充分利用多核处理器的并行能力,提高任务处理效率。例如,图像处理、数据处理、并行排序等。
  2. 递归型任务:对于递归型的任务,WorkStealingPool 能够适应任务的动态变化,根据需要创建和调度子任务,以实现更高效的递归执行。例如,斐波那契数列计算、归并排序等。
  3. 高吞吐量任务:WorkStealingPool 的工作窃取算法可以减小线程之间的竞争,并且能够在任务队列为空时从其他线程窃取任务,从而减少线程的等待时间,提高整体的任务处理吞吐量。适用于需要高吞吐量的任务场景。
  4. CPU 密集型任务:对于需要大量的 CPU 计算而没有 I/O 阻塞的任务,使用 WorkStealingPool 可以更好地充分利用 CPU 核心,并且可以根据需要增加或减少线程数量,以适应任务的计算量。

需要注意的是,WorkStealingPool 在任务数较少或任务之间存在 I/O 等阻塞时可能不如其他类型的线程池效果好,因为工作窃取算法适用于 CPU 密集型任务。在实际应用中,根据具体情况选择合适的线程池类型和参数才能达到最佳的性能和效果。

返回列表