位置:首页 > C++ > C++11实战:深入解析与实现CachedThreadPool缓存线程池

C++11实战:深入解析与实现CachedThreadPool缓存线程池

时间:2026-08-27  |  作者:清风无痕  |  阅读:0

目录

  1. 一、线程池概述
  2. 1.1 线程池概念
  3. 1.2 按应用场景分类
  4. 1.3 线程池模式
  5. 1.4 半同步/半异步模式分析
  6. 1.5 线程池实现的关键技术分析

前言

本文基于C++11标准,深入探讨线程池的核心机制与最佳实践。线程池通过复用线程、控制并发及任务排队,显著提升系统性能。文章首先概述FixedThreadPool、SingleThreadPool等五种常见模式,重点剖析半同步/半异步架构下的关键实现技术。随后,结合VS2019与Linux g++环境,详细解析CachedThreadPool的动态线程调整需求,展示如何根据任务负载灵活创建或复用线程,为高并发场景提供高效解决方案。

C++11实战:深入解析与实现CachedThreadPool缓存线程池 的核心流程信息图
C++11实战:深入解析与实现CachedT用简体中文信息图概括C++11实战:深入解析与实现CachedT的核心流程、关键规则与实践要点。

一、线程池概述

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

4 半同步/半异步模式分析 对应的技术说明图
4 半同步/半异步模式分析展示半同步/半异步模式分析涉及的工具作用、执行步骤与结果判断。

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

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

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

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

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

二、CachedThreadPool的实现

2.1 需求

动态调整线程数量:CachedThreadPool的线程数量是动态调整的。当有新任务提交时,如果线程池中有空闲的线程,则会立即使用空闲线程执行任务;如果线程池中没有空闲线程,则会创建一个新的线程来执行任务。当线程空闲一段时间后,超过一定的时间,会被回收销毁。

2.2 SyncQueue同步队列的设计和实现

双条件变量分离生产 / 消费唤醒;m_needStop停机标志实现安全退出;Stop 接口先等待队列清空再广播通知,保证任务不丢失;同时提供单条 Take 与批量 Take 两种消费模式。

#ifndef SYNCQUEUE_HPP
#define SYNCQUEUE_HPP
#include
#include
#include
#include
using namespace::std;
template
class SynQueue
{
private:
    std::list m_queue;                     // 任务存储容器
    mutable std::mutex m_mutex;               // 全局互斥锁,保护队列所有读写
    std::condition_variable m_notEmpty;       // 条件变量:队列非空,唤醒消费者
    std::condition_variable m_notFull;        // 条件变量:队列未满,唤醒生产者
    size_t m_waitTime;                        // 阻塞等待超时时间(秒)
    int m_maxSize;                            // 有界队列最大容量
    bool m_needStop;                          // 队列停止标记,用于优雅退出
    // 判断队列是否已满
    bool IsFull()const;
    // 判断队列是否为空
    bool IsEmpty()const;
    // 底层入队通用模板,完美转发统一处理左值、右值
    template
    int Add(F&& x);
public:
    // 构造:指定队列上限、等待超时时间
    SynQueue(int maxsize = 100, int waittime = 1);
    // 左值版本入队
    int Put(const T& x);
    // 右值版本入队
    int Put(T&& x);
    // 阻塞等待任务,仅做等待探测,不取出数据
    int notTask();
    // 批量取出全部任务,移动语义转移list
    void Take(std::list& list);
    // 阻塞获取单个任务,支持超时
    int Take(T& t);
    // 优雅停止队列:等待任务消费完毕,唤醒所有阻塞线程
    void Stop();
    // 查询队列是否为空
    bool Empty()const;
    // 查询队列是否已满
    bool Full()const;
    // 获取当前任务数量
    size_t Size()const;
};
#endif // !SYNCQUEUE_HPP

Add函数

上锁后使用带超时的wait_for等待队列腾出空间,谓词同时判断停止标记与队列是否未满,规避虚假唤醒。等待超时直接返回失败;若触发停止标志则终止入队。条件满足时通过完美转发将元素加入队列,随后通知消费者队列已有任务,释放锁。

template
int Add(F&& x)
{
	std::unique_lock locker(m_mutex);
	// 等待队列有空位
	// 谓词返回 true 时停止等待(即:需要停止 OR 队列未满),wait_for 返回 false 表示超时
	if (!m_notFull.wait_for(locker, std::chrono::seconds(m_waitTime),
		[this] { return m_needStop || !IsFull(); })) 
	{
		std::cout << "task queue full, timeout return 1" << std::endl;
		return 1;		// 超时失败
	}
	// 检查是否收到停止信号
	if (m_needStop)
	{
		std::cout << "同步队列停止工作..." << std::endl;
		return 2;		// 停止状态
	}
	// 完美转发并入队
	m_queue.push_back(std::forward(x));
	// 通知消费者
	m_notEmpty.notify_one();
	return 0;			// 成功
}

Take函数

单元素 Take (T& t):上锁阻塞等待非空信号,支持超时。被唤醒后校验停止标识,正常则取出队首元素并弹出,通知生产者队列腾出位置;超时或停止对应返回不同状态码。

三、CachedThreadPool的测试

3.1 测试1

批量 Take (list&):持续等待直到队列不为空或收到停止信号;正常情况下通过移动语义一次性迁移整个队列,O (1) 批量消费,随后唤醒生产者。

void Take(std::list& list)
{
	// 获取互斥锁,保护共享队列 m_queue
	std::unique_lock locker(m_mutex);
	// 循环等待:当系统未停止且队列为空时,阻塞当前线程
	while (!m_needStop && IsEmpty())
	{
		m_notEmpty.wait(locker);
	}
	// 退出检查:若是因为收到停止信号而跳出循环,则直接返回
	if (m_needStop) 
	{
		std::cout << "同步队列停止工作..." << std::endl;
		return;
	}
	// 批量提取:使用移动语义将整个队列内容转移给外部 list,效率极高(O(1))
	list = std::move(m_queue);
	// 通知生产者:队列已清空,唤醒一个正在等待“非满”条件的生产线程
	m_notFull.notify_one();
}
int Take(T& t)
{
	std::unique_lock locker(m_mutex);
	// 带超时的条件等待
	// 谓词 [this] { return !m_needStop && !IsEmpty(); } ,若返回 false (超时),进入 if 分支
	if (!m_notEmpty.wait_for(locker, std::chrono::seconds(m_waitTime),
		[this] { return !m_needStop && !IsEmpty(); }))
	{
		return 1;       // 等待超时,队列仍为空
	}
	if (m_needStop)
	{
		std::cout << "同步队列停止工作..." << std::endl;
		return 2;       // 系统已请求停止
	}
	// 提取队首元素
	t = m_queue.front();
	m_queue.pop_front();
	m_notFull.notify_one();
	return 0;           // 成功获取任务
}

Stop函数

先上锁,循环等待队列所有任务消费完成;等待清空后置停止标志,调用notify_all广播唤醒所有阻塞在两个条件变量上的生产者、消费者线程,避免线程永久挂起,保证剩余任务处理完毕再退出。

void Stop()
{
	std::unique_lock locker(m_mutex);
	// 等待队列清空:防止在还有任务未处理时强行停止,确保数据完整性
	while (!IsEmpty())
	{
		m_notFull.wait(locker);
	}
	// 设置停止标志:通知所有工作线程准备退出
	m_needStop = true;
	// 广播通知:唤醒所有阻塞在 m_notEmpty(等待任务)和 m_notFull(等待空间)的线程
	m_notEmpty.notify_all();
	m_notFull.notify_all();
}

2.3 CachedThreadPool线程池的设计和实现

CachedThreadPool 为可缓存动态线程池,依托 SyncQueue 同步队列构建;具备线程动态扩容、空闲线程超时回收机制,常驻核心线程保障基础任务吞吐,突发任务负载下自动新建线程;任务队列满载时内置调用者运行拒绝策略,避免任务丢失。

#ifndef CACHEDTHREADPOOL_HPP
#define CACHEDTHREADPOOL_HPP
#include"SynQueue2.hpp"
#include
#include
#include
#include
using namespace std;
int MaxTaskCount = 2;
const int KeepAliveTime = 10;
class CachedThreadPool
{
public:
	using Task = std::function;		
private:
	std::unordered_map> m_threadgroup;	//线程组
	int m_coreThreadSize;					//核心线程下限
	int m_maxThreadSize;					//最大线程上限
	std::atomic_int m_idleThreadSize;		//空闲线程计数
	std::atomic_int m_curThreadSize;		//当前总线程数量
	mutable std::mutex m_mutex;				//容器互斥锁
	SynQueue m_queue;					//任务同步队列
	std::atomic_bool m_running;				//线程池运行标记
	std::once_flag m_flag;					//保证Stop仅执行一次
	void Start(int numthreads);
	void RunInThread();
	void StopThreadGroup();
public:
	CachedThreadPool(int initNumThreads=8,int taskPoolSize=MaxTaskCount);
	~CachedThreadPool();
	void Stop();
	//提交无返回值任务
	template
	void execute(Func&& func, Args&&... args);
	//提交带返回值任务
	template
	auto submit(Func&& func, Args&&... args)
		-> std::future ;
};
#endif

RunInThread 工作线程主循环

工作线程持续循环从同步队列阻塞获取任务执行;持续无任务时统计空闲时长。空闲时长超过KeepAliveTime,且总线程数大于核心线程阈值,则销毁临时线程,仅保留核心线程。单次任务执行完成后重置空闲计时起点。

void RunInThread()
{
	auto tid = std::this_thread::get_id();
	//记录线程空闲起始时间
	auto startTime = std::chrono::high_resolution_clock().now();
	while (m_running)
	{
		Task task;
		if (m_queue.Size() == 0 && m_queue.notTask())
		{
			auto now = std::chrono::high_resolution_clock().now();
			auto intervalTime = std::chrono::duration_cast(now - startTime);
			std::lock_guard lock(m_mutex);
			//空闲超时 && 当前线程数大于核心线程,销毁临时线程
			if (intervalTime.count() >= KeepAliveTime && m_curThreadSize > m_coreThreadSize)
			{
				m_threadgroup.find(tid)->second->detach();
				m_threadgroup.erase(tid);
				m_curThreadSize--;
				m_idleThreadSize--;
				cout << "空闲线程销毁 " << m_curThreadSize << " " << m_coreThreadSize << endl;
				return;
			}
			//阻塞获取任务并执行
			if (!m_queue.Take(task) && m_running)
			{
				m_idleThreadSize--;
				task();
				m_idleThreadSize++;
				//任务执行完毕,重置空闲计时
				startTime = std::chrono::high_resolution_clock().now();
			}
		}
	}
}

submit 提交带返回值任务

通过std::packaged_task封装任务,返回std::future供客户端获取返回结果。入队失败触发调用者运行策略;无空闲线程、线程总数未达上限时,动态新建工作线程扩容。

template
auto submit(Func&& func, Args&&... args)
	-> std::future 
{
	using RetType = decltype(func(args...));
	//封装任务,支持获取返回值
	auto task = std::make_shared>(
		std::bind(std::forward(func),std::forward(args)...)
	);
	std::future result = task->get_future();
	//尝试将任务放入同步队列
	if (m_queue.Put([task]() {(*task)();})!=0)
	{
		std::cout << "调用者运行策略" << std::endl;
		(*task)();
	}
	//无空闲线程,且未达到最大线程限制,新建线程扩容
	if (m_idleThreadSize<=0&&m_curThreadSize lock(m_mutex);
		auto tha = std::make_shared(
			std::thread(&CachedThreadPool::RunInThread,this)
		);
		std::thread::id tid = tha->get_id();
		tha->detach();
		m_threadgroup.emplace(tid, std::move(tha));
		m_idleThreadSize++;
		m_curThreadSize++;
	}
	return result;
}

StopThreadGroup 线程池停止逻辑

借助同步队列Stop接口封锁任务入队,关闭线程池运行标记;等待所有工作线程执行完成,统一回收线程资源,保证队列剩余任务处理完毕再优雅退出。

void StopThreadGroup()
{
	m_queue.Stop();		//停止同步队列,禁止新任务入队
	m_running = false;	//修改线程池运行标志
	//等待所有工作线程执行结束
	for (auto thread : m_threadgroup) 
	{
		thread.second->join();
	}
	m_threadgroup.clear();
}

3.2 测试2

四、线程池进阶拓展

4.1 FixedThreadPool与CachedThreadPool特性对比

4.2 最佳实践

FixedThreadPool和CachedThreadPool两者对高负载的应用都不是特别友好。

CachedThreadPool要比FixedThreadPool危险很多。

如果应用要求高负载、低延迟、最好不要选择以上两种线程池:

  • 任务队列的无边界:会导致内存溢出以及高延迟
  • 长时间运行会导致CachedThreadPool在线程创建上失控

因为两者都不是特别友好,所以推荐使用ThreadPoolExecutor,它提供了很多参数模版可以进行细粒度的控制。

  • 将任务队列设置成有边界的队列
  • 使用合适的RejectionHandler拒绝处理程序
  • 如果在任务完成前后需要执行某些操作,可以重载
  • 重载ThreadFactory,如果有线程定制化的需要
  • 在运行时动态控制线程池的大小(Dynamic Thread Pool)

4.3 使用场景

适用于以下场景:

CachedThreadPool 适用场景分析

大量短期任务处理

CachedThreadPool 专为处理大量短期任务而设计。当新任务到达时,线程池会尽可能创建新线程以执行任务;若存在空闲线程,则优先复用现有线程,避免线程闲置。这种机制有效避免了频繁创建和销毁线程所带来的额外系统开销,提升了资源利用率。

快速响应需求

该线程池适用于对响应速度要求较高的场景。它能够根据任务的到达情况快速创建并启动新线程,从而显著减少任务在队列中的等待时间,确保系统能够迅速响应用户请求或外部事件。

无限制线程数量

CachedThreadPool 的最大线程数不设上限。只要系统内存空间充足,它可以根据任务的动态到来情况,无限制地创建新线程。这种灵活性使其能够应对突发的流量高峰,但同时也要求开发者对系统资源有清晰的认知。

高并发短期任务

由于支持动态创建线程,该线程池特别适合处理具有高并发特性的短期任务。当任务处理完毕后,线程池会保留一定数量的空闲线程,以便立即响应下一批任务的到来,从而保持系统的高吞吐量和低延迟。

注意事项与总结

需要特别注意的是,CachedThreadPool 的线程数量不受限制。如果任务量过大,可能导致线程数量激增,进而造成系统资源过度消耗甚至崩溃。因此,在实际使用中,必须根据具体业务场景灵活调整线程数量,或在必要时选用其他类型的线程池以更好地控制资源使用。

总结

利用 C++11 的线程相关特性,我们可以编写出简洁高效的并发程序。例如,通过结合线程、条件变量和互斥量,可以构建一个轻量级的线程池,从而避免频繁创建线程带来的性能损耗。在使用线程池时,还需注意以下关键问题:首先,必须确保线程池中的任务不会挂死,否则会导致线程耗尽,引发假死现象;其次,应避免长时间执行单个任务,这会导致后续任务大量堆积而无法得到及时处理。对于耗时较长的任务,建议采用单独的线程进行处理,以保障系统的稳定性和响应能力。

//无返回值测试任务
void func(int index)
{
    static int num = 0;
    //静态变量验证多线程并发访问
    cout << "func_" << index << " num: " << ++num << endl;
}
//带返回值测试任务
int add(int a, int b)
{
    return a + b;
}
int main()
{
    //创建可缓存动态线程池
    CachedThreadPool mypool;
    //循环批量提交1000个任务
    for (int i = 0; i < 1000; ++i)
    {
        if (i % 2 == 0)
        {
            //偶数:提交带返回值任务,通过future接收结果
            //auto pa = mypool.submit(add, i, i + 1);
            auto pa = mypool.submit([=]() { return add(i, i + 1); });
            //阻塞等待任务完成,获取返回值并打印
            cout << pa.get() << endl;
        }
        else
        {
            //奇数:提交无需返回值的任务
            mypool.execute(func, i);
        }
    }
    return 0;
}
//创建核心线程数量为2的动态线程池
CachedThreadPool pool(2);
//耗时任务,模拟长时间业务处理
int add(int a, int b, int s)
{
    //任务休眠s秒,占用工作线程
    std::this_thread::sleep_for(std::chrono::seconds(s));
    int c = a + b;
    cout << "add begin ..." << endl;
    return c;
}
//在线程池提交任务
void add_a()
{
    //向线程池提交耗时4s任务
    auto r = pool.submit(add, 10, 20, 4);
    cout << "add_a: " << r.get() << endl;
}
void add_b()
{
    //向线程池提交耗时6s任务
    auto r = pool.submit(add, 20, 30, 6);
    cout << "add_b: " << r.get() << endl;
}
void add_c()
{
    //向线程池提交耗时1s任务
    auto r = pool.submit(add, 30, 40, 1);
    cout << "add_c: " << r.get() << endl;
}
void add_d()
{
    //向线程池提交耗时9s任务
    auto r = pool.submit(add, 10, 40, 9);
    cout << "add_d: " << r.get() << endl;
}
int main()
{
    //开启4个独立线程并发向线程池投递任务,制造任务突发压力
    std::thread tha(add_a);
    std::thread thb(add_b);
    std::thread thc(add_c);
    std::thread thd(add_d);
    tha.join();
    thb.join();
    thc.join();
    thd.join();
    //休眠20s,等待空闲线程触发超时回收逻辑
    std::this_thread::sleep_for(std::chrono::seconds(20));
    //再次发起新一轮任务
    std::thread the(add_a);
    std::thread thf(add_b);
    the.join();
    thf.join();
    return 0;
}
特性FixedThreadPoolCachedThreadPool
重用FixedThreadPool 与 CacheThreadPool 差不多,也是能 reuse 就用,但不能随时建新的线程缓存型池子,先查看池中有没有以前建立的线程,如果有,就 reuse;如果没有,就建一个新的线程加入池中
池大小可指定 nThreads,固定数量可增长,最大值 Integer.MAX_VALUE
队列大小无限制无限制
超时无 IDLE默认 60 秒 IDLE
使用场景FixedThreadPool 多数针对一些很稳定很固定的正规并发线程,多用于服务器。定长线程池;适用于执行负载重,cpu 使用频率高的任务;这个主要是为了防止太多线程进行大量的线程频繁切换,得不偿失。大量短生命周期的异步任务。适用于执行大量 (并发) 短期异步的任务;注意,任务量的负载要轻。
结束不会自动销毁注意,放入 CachedThreadPool 的线程不必担心其结束,超过 TIMEOUT 不活动,其会自动被终止。

免责声明:文中图文均来自网络,如有侵权请联系删除,心愿游戏发布此文仅为传递信息,不代表心愿游戏认同其观点或证实其描述。

相关文章

更多

精选合集

更多

大家都在玩

热门话题

大家都在看

更多