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 线程池实现的关键技术分析
线程池有两个活动过程,一个是往同步队列中添加任务的过程,另一个是从同步队列中取任务的过程。
半同步半异步线程池活动图
二、FixedThreadPool的实现
2.1 需求
FixedThreadPool 是一个固定大小的线程池,在创建时会指定线程池中线程的数量、每当有任务提交到线程池时,线程会启动一个线程来执行任务,直到达到线程池的最大线程数。
2.2 SyncQueue同步队列的设计
同步队列为线程池中三层结构中的中间层,主要作用是保证任务队列中,共享数据的线程安全,还为上一层同步服务层提供添加新任务的接口,以及为下一层异步服务层提供获取任务的接口。同时,还要限制任务数的上限,避免任务过多导致内存暴涨问题。
同步队列的实现,我们会用到C++11的互斥锁、条件变量、右值引用、std::move以及std::forward。std::move是为了实现移动语义,std::forward是为了实现完美转发。同步队列的锁是用来线程同步的,条件变量是用来实现线程通信的,即线程池空了就要等待,不空就通知一个线程去处理;线程池满了就等待,直到没有满的时候才通知上层添加新任务。
2.3 SyncQueue类型的代码实现
1 同步队列类型:生产者消费者队列
2 同步队列代码实现
下面具体介绍同步队列的3个函数Take、Add、Stop的实现
1. Take函数
先创建一个unique_lock获取mutex,然后通过条件变量m_notEmpty来等待判断式,判断式由两个条件组成,一个是停止的标志,另一个是不为空的条件,当不满足任何一个条件时,条件变量会释放mutex并将线程置与waiting状态,等待其他线程调用notify_one/notify_all将其唤醒;当满足任何一个条件时,则继续往下执行后面的逻辑。将队列中的任务取出,并唤醒一个正处于等待状态的添加任务的线程去添加任务。
//批量消费 void Take(std::list<T>& list) { std::unique_lock<std::mutex> locker(m_mutex); // 等待队列非空或停止信号 m_notEmpty.wait(locker, [this] { return m_needStop || !IsEmpty(); }); if (m_needStop) { return; } // 移动语义:将整个内部队列的所有权转移给外部 list list = std::move(m_queue); // 通知生产者:因为内部队列清空了,肯定“不满”了 m_notFull.notify_one(); } //单件消费 void Take(T& take) { std::unique_lock<std::mutex> locker(m_mutex); m_notEmpty.wait(locker, [this] { return m_needStop || !IsEmpty(); }); if (m_needStop) { return; } // 拷贝/移动前端元素 take = m_queue.front(); // 移除前端元素 m_queue.pop_front(); m_notFull.notify_one(); }2. Add函数
Add的过程和Take的过程是类似的,也是先获取mutex,然后检查条件是否满足,不满足条件时,释放mutex继续等待,如果满足条件,则将新的任务插入到队列中,并唤醒取任务的线程去取数据。
template<class F> void Add(F&& task) { // 加独占互斥锁,管控队列并发访问 std::unique_lock<std::mutex> locker(m_mutex); // 队列满则阻塞生产者,线程池停止也直接放行跳出等待 m_notFull.wait(locker, [this] { return m_needStop || !IsFull(); }); // 线程池已停止,放弃添加任务,直接返回 if (m_needStop) { return; } // 完美转发任务,推入任务队列,避免拷贝 m_queue.push_back(std::forward<F>(task)); // 唤醒一个等待任务的工作线程 m_notEmpty.notify_one(); }3. Stop函数
Stop函数先获取mutex,然后将停止标志置为true。由于线程m_needStop为true时会退出,所有所有的等待线程会相继退出。将m_notFull.notify_all()放到了lock_guard保护范围之外,被唤醒的线程获取锁的时候不需要等待lock_guard释放锁,性能会好一点。
void Stop() { { std::unique_lock<std::mutex> locker(m_mutex); m_needStop = true; } m_notFull.notify_all(); m_notEmpty.notify_all(); }2.4 FixedThreadPool线程池的设计
一个完整的线程池包括三层:同步服务层、排队层和异步服务层,其实这也是一种生产者-消费者模式,同步层是生产者,不断将新任务添加到排队层,因此,线程需要提供一个添加新任务的接口供生产者使用;消费者是异步层,具体由线程中预先创建的线程去处理排队层中的任务,排队层是一个同步队列,它内部保证了上下层对共享数据的安全访问,同时还要保证不会被无限制地添加任务导致内存暴涨。另外,线程还要一个停止的接口,让用户能够在需要时候停止线程池的运行。
2.5 FixedThreadPool代码实现
class FixedThreadPool { public: using Task = std::function<void(void)>; private: std::list<std::shared_ptr<std::thread>> m_threadgroup; //线程组 SyncQueue<Task> m_queue; //同步队列 std::atomic_bool m_running; //停止线程池 std::once_flag m_flag; void Start(int numthreads) { m_running = true; for (int i = 0;i < numthreads;i++) { m_threadgroup.push_back( std::make_shared<std::thread>( &FixedThreadPool::RunInThread, this ) ); } } void RunInThread() { while (m_running.load(std::memory_order_acquire)) { Task task; m_queue.Take(task); if (task && m_running.load(std::memory_order_acquire)) { task(); } } } void StopThreadGroup() { m_queue.Stop(); m_running.store(false, std::memory_order_release); for (auto& thread : m_threadgroup) { if (thread) { thread->join(); } } m_threadgroup.clear(); } public: // 仅保留带参构造,默认参数实现无参创建,删除空声明 FixedThreadPool(int numThreads = std::thread::hardware_concurrency()) :m_queue(MaxTaskCount), m_running(false) { Start(numThreads); } ~FixedThreadPool() { Stop(); } void Stop() { std::call_once(m_flag, [this] {StopThreadGroup();}); } void AddTask(Task&& task) { m_queue.Put(std::forward<Task>(task)); } void AddTask(const Task& task) { m_queue.Put(task); } };三、FixedThreadPool的测试
3.1 测试1
//测试任务:加法计算,通过promise返回结果 void Add(int a, int b, std::promise<int>& c_promise) { std::cout << "add begin ..." << std::endl; std::this_thread::sleep_for(std::chrono::milliseconds(2000)); int c = a + b; c_promise.set_value(c); std::this_thread::sleep_for(std::chrono::milliseconds(1000)); std::cout << "add end ..." << std::endl; } //线程a:提交10+20任务 void add_a() { std::promise<int> c_promise; std::future<int> a_future = c_promise.get_future(); std::function<void(void)> f = std::bind(Add, 10, 20, std::ref(c_promise)); pool.AddTask(f); std::cout << "add_a: " << a_future.get() << std::endl; } //线程b:提交20+30任务 void add_b() { std::promise<int> c_promise; std::future<int> a_future = c_promise.get_future(); std::function<void(void)> f = std::bind(Add, 20, 30, std::ref(c_promise)); pool.AddTask(f); std::cout << "add_b: " << a_future.get() << std::endl; } //线程c:提交30+40任务 void add_c() { std::promise<int> c_promise; std::future<int> a_future = c_promise.get_future(); std::function<void(void)> f = std::bind(Add, 30, 40, std::ref(c_promise)); pool.AddTask(f); std::cout << "add_c: " << a_future.get() << std::endl; } int main() { //创建3个独立生产者线程,分别提交三组不同任务 std::thread tha(add_a); std::thread thb(add_b); std::thread thc(add_c); //主线程阻塞等待所有生产者线程执行完毕 tha.join(); thb.join(); thc.join(); return 0; }3.2 测试2
//线程池内执行:堆内存分配任务 void my_malloc(int size, std::promise<int*>& c_promise) { //在堆上分配指定大小内存,返回堆起始地址 int* p = (int*)malloc(size); //将分配得到的堆指针存入promise,唤醒外部阻塞的future.get() c_promise.set_value(p); } //线程池内执行:释放堆内存任务 void my_free(int* p) { free(p); } //生产者线程a:申请10个int大小堆内存,使用后投递释放任务 void my_a() { // 单次分配10个int的字节长度 int n = 10; // 用于存储异步分配返回的堆指针 std::promise<int*> c_promise; // 绑定promise,用于阻塞等待任务执行结果 std::future<int*> a_future = c_promise.get_future(); // 绑定my_malloc任务,promise不能拷贝,必须用std::ref传递引用 std::function<void(void)> f = std::bind(my_malloc, sizeof(int) * n, std::ref(c_promise)); // 将内存分配任务提交至线程池执行 pool.AddTask(f); // 阻塞当前线程,直到线程池内malloc任务完成、指针写入promise int* p = a_future.get(); // 判断内存分配是否失败 if (p == nullptr) { std::cout << "失败 ... " << std::endl; exit(1); } // 对堆内存写入数据,验证指针可用 *p = 5; // 打印堆地址与存储的数据 std::cout << "p: " << p << " *p: " << *p << std::endl; // 再次向线程池投递释放任务,归还堆内存 pool.AddTask(std::bind(my_free, p)); } //生产者线程b:申请5个int大小堆内存,使用后投递释放任务 void my_b() { int n = 5; std::promise<int*> c_promise; std::future<int*> a_future = c_promise.get_future(); std::function<void(void)> f = std::bind(my_malloc, sizeof(int) * n, std::ref(c_promise)); pool.AddTask(f); int* p = a_future.get(); if (p == nullptr) { std::cout << "失败 ... " << std::endl; exit(1); } *p = 10; std::cout << "p: " << p << " *p: " << *p << std::endl; pool.AddTask(std::bind(my_free, p)); } //生产者线程c:申请1个int大小堆内存,使用后投递释放任务 void my_c() { int n = 1; std::promise<int*> c_promise; std::future<int*> a_future = c_promise.get_future(); std::function<void(void)> f = std::bind(my_malloc, sizeof(int) * n, std::ref(c_promise)); pool.AddTask(f); int* p = a_future.get(); if (p == nullptr) { std::cout << "失败 ... " << std::endl; exit(1); } *p = 15; std::cout << "p: " << p << " *p: " << *p << std::endl; pool.AddTask(std::bind(my_free, p)); } int main() { // 创建3个独立生产者线程,并发向线程池提交内存分配任务 std::thread tha(my_a); std::thread thb(my_b); std::thread thc(my_c); // 主线程阻塞,等待3个生产者线程全部执行完毕 tha.join(); thb.join(); thc.join(); return 0; }四、线程池进阶拓展
4.1 线程池的拒绝策略
AbortPolicy: 中止策略。默认的拒绝策略,直接抛出RejectedExecutionException。调用者可以捕获这个异常,然后根据需求编写自己的处理代码。
DiscardPolicy: 抛弃(丢弃)策略。什么都不做,直接抛弃被拒绝的任务。
DiscardOlderPolicy: 抛弃最老策略。抛弃阻塞队列中最老的任务,相当于就是队列中下一个将要被执行的任务,然后重新提交被拒绝的任务。
CallerRunsPolicy: 调用者运行策略。在调用者线程中执行该任务。该策略实现了一种调节机制,该策略既不会抛弃任务,也不会抛出异常,而是将任务回退到调用者(调用线程执行任务的主线程),由于执行任务需要一定时间,因此主线程至少在一段时间内不能提交任务,从而使得线程池有时间来处理完正在执行的任务。
4.2 实现调用者运行策略的代码实现
SyncQueue
//调用者运行策略 //0 任务添加成功 1 任务队列已达上限 2 任务队列停止工作 template<class F> int Add(F&& task) { // 加独占互斥锁,管控队列并发访问 std::unique_lock<std::mutex> locker(m_mutex); // 等待队列非满或停止信号。若1秒内未满足条件则超时返回1 if (!m_notFull.wait_for(locker, std::chrono::seconds(1), [this] {return m_needStop || IsFull();})) { return 1; } // 检查是否已收到停止指令 if (m_needStop) { return 2; } // 完美转发任务,推入任务队列,避免拷贝 m_queue.push_back(std::forward<F>(task)); // 唤醒一个等待任务的工作线程 m_notEmpty.notify_one(); return 0 } int Put(const T& task) { return Add(task); } int Put(T&& task) { return Add(std::forward<T>(task)); }FixedThreadPool
//调用者运行策略 void AddTask(Task&& task) { if (m_queue.Put(std::forward<Task>(task)) != 0) { std::cerr << "task queue is full,Add task fail." << std::endl; task(); } } void AddTask(const Task& task) { if (m_queue.Put(task) != 0) { std::cerr << "task queue is full,Add task fail." << std::endl; task(); } }4.3 改写AddTask函数降低客户端得到返回值的难度
版本1:基础实现(存在拷贝开销)
template<class Func, class... Args> auto AddTask(Func&& func, Args&&... args) -> std::future<decltype(std::forward<Func>(func)(std::forward<Args>(args)...))> { // 推导返回类型 using RetType = decltype(func(args...)); // 包装任务:bind 会拷贝 func 和 args,可能存在性能损耗 std::packaged_task<RetType()> task(std::bind(func, args...)); // 获取与任务绑定的 future 对象,用于后续获取结果 std::future<RetType> result = task.get_future(); // 同步执行:直接在当前线程运行任务 task(); return result; }版本2:优化实现(支持完美转发)
template<class Func, class... Args> auto AddTask(Func&& func, Args&&... args) -> std::future<decltype(std::forward<Func>(func)(std::forward<Args>(args)...))> { // 推导返回类型 using RetType = decltype(func(args...)); // 包装任务:使用 std::forward 实现完美转发,避免不必要的拷贝,支持移动语义 std::packaged_task<RetType()> task( std::bind(std::forward<Func>(func), std::forward<Args>(args)...) ); // 获取 future 对象 std::future<RetType> result = task.get_future(); // 同步执行:直接在当前线程运行任务 task(); return result; }版本3(推荐)
//版本3(推荐) template<class Func, class... Args> auto AddTask(Func&& func, Args&&... args) -> std::future<decltype(std::forward<Func>(func)(std::forward<Args>(args)...))> { // 推导返回类型 using RetType = decltype(std::forward<Func>(func)(std::forward<Args>(args)...)); // 使用 shared_ptr 包装 packaged_task,便于在 Lambda 中捕获并异步执行 auto task = std::make_shared<std::packaged_task<RetType()>>( std::bind(std::forward<Func>(func), std::forward<Args>(args)...) ); // 获取 future 对象 std::future<RetType> result = task->get_future(); // 尝试放入队列:若成功则异步执行,若失败则当前线程同步执行(回退策略) if (m_queue.Put([task]() { (*task)(); }) != 0) { (*task)(); // 队列满时,由调用者同步执行 } return result; }4.4 FixedThreadPool的使用场景
1.并发限制:当有大量任务需要执行,但希望限制并发线程数时,可以使用FixedThreadPool。可以控制并发执行的线程数量,避免系统资源过度占用和线程竞争导致性能下降。
2.稳定且可控的任务执行:当任务量稳定且任务的执行时间较短时,FixedThreadPool是一个合适的选择。由于线程池中的线程数量固定,可以提供稳定的执行环境,避免频繁地创建和销毁线程的开销。
3.服务器应用:FixedThreadPool适用于服务器中需要处理大量请求的场景。
4.批量任务处理:当需要对一批任务进行并发处理时,FixedThreadPool可以提供线程池管理和调度的支持。
