最近工作涉及到后台批量处理文件数据,而且需要极快响应结果,于是就想到了使用线程池。刚开始构思线程池的生产者在提交工作任务时,可以传入函数指针,这样就能让线程去批量执行相同的文件数据处理逻辑来处理不同的文件。从而达到利用并发计算来减少计算时间的目的。

  写完之后,效果还不错。但后续发现传入函数指针只适合于调用同一逻辑的批量任务,就比如这次工作中涉及到的场景。后续扩展性很差,想传入不同的工作函数都难以实现,并且还需要手动强转类型,安全性也得不到保证。

  就去向组内的大佬请教,看了他之前在项目架构中写的一个可复用的模板线程池。看完也是感叹大佬不愧是大佬。就趁着这次学习的机会,从头复盘下线程池的基本知识点。

 

线程池(ThreadPool):基本的线程池主要由工作线程池,任务队列,保护任务队列的互斥锁,唤醒/阻塞工作线程和提交函数的条件变量,线程池停止标志位组成。

    核心价值:线程池主要是为了解决频繁创建/销毁线程的性能开销,并统一管理并发数而创建的。

      1.线程复用: 创建/销毁线程的开销(系统调用,内核态/用户态切换,栈空间分配)远大于执行简单任务,线程池提前创建固定数量的线程并复用,可避免频繁创建销毁的开销。

     2.并发控制:限制最大线程数,避免创建过多线程,导致CPU上下文切换过载。

    3.任务调度:统一管理任务队列,可手动实现任务的优先级,排队,延迟等待。

  经典应用场景:

 1.提交的任务函数主体相同+参数不同

     例如这次的批量处理文件,处理数据的逻辑是相同的,但传入的参数(即文件存储路径)不同。

  2.提交的函数主体不同+参数不同

     例如:后台服务器中,线程池同时处理“用户登录”,“订单生产”,“日志写入“等功能,需要执行不同的业务逻辑,传入的参数也不相同。

这时候就需要用到模板+std::function<void*>的实现方式,来擦除函数类型。从而达到可以传入任意返回值,任意函数名,任意参数的效果。

这里直接附上代码实现,有我逐帧学习时的过程,233......

#include <vector>
#include <queue>
#include <thread>
#include <mutex>
#include <condition_variable>
#include <functional>
#include <future>
#include <stdexcept>

/*未实现线程数量动态扩容功能。1、可以在高负载时新增线程,低负载时回收线程(需要安全同步,复杂)。
 2、也可以使用有界队列 + 阻塞/超时/拒绝策略:当队列满时阻塞生产者或返回错误/丢弃任务。*/

class ThreadPool {

    private:
    std::vector<std::thread> workers;       // 工作线程集合
    std::queue<std::function<void()>> tasks; // 任务队列(存储无参函数对象)
    std::mutex queue_mutex;                 // 保护任务队列的互斥锁
    std::condition_variable condition;      // 用于线程等待/唤醒的条件变量
    bool stop;                              // 线程池停止标记


public:
    // 构造函数:指定线程数量(默认使用CPU核心数)
    explicit ThreadPool(size_t threads = std::thread::hardware_concurrency())
        : stop(false) {
        if (threads == 0) {
            throw std::invalid_argument("线程数量不能为0");
        }
        // 创建工作线程   
        for (size_t i = 0; i < threads; ++i) {
           workers.emplace_back([this] {
                // 工作线程循环:不断从任务队列取任务执行
                while (true) {
                    std::function<void()> task;
                    // 加锁取任务
                    {
                        std::unique_lock<std::mutex> lock(this->queue_mutex);
                        // 等待条件:队列非空 或 线程池停止  wai()的第二个参数为退出等待所需条件,返回false继续等待,返回true退出等待
                        //核心要求:工作线程需要在任务队列非空或线程池被停止时唤醒,需要在任务队列为空或者线程池未停止时阻塞
                        //在线程池停止时唤醒工作线程是为了让其监测到停止信号并退出,防止出现极端情况,线程标志位设置后,线程刚好未收到退出信号,导致一直阻塞
                        this->condition.wait(lock,
                            [this] { return this->stop || !this->tasks.empty(); });
                        // 若停止且队列空,退出线程 二次判断,防止虚假唤醒
                        if (this->stop && this->tasks.empty()) {
                            return;
                        }
                        // 取出任务 
                        task = std::move(this->tasks.front()); //移动语义,避免多余拷贝
                        this->tasks.pop();
                    }
                    // 执行任务(解锁后执行,避免阻塞其他线程取任务)
                    task();
                }
                });
        }
    }

    // 禁止拷贝构造和赋值
    ThreadPool(const ThreadPool&) = delete;
    ThreadPool& operator=(const ThreadPool&) = delete;

    // 析构函数:停止线程池
    ~ThreadPool() {
        {
            std::unique_lock<std::mutex> lock(queue_mutex);
            stop = true; // 标记停止状态
        }
        condition.notify_all(); // 唤醒所有等待的线程
        // 等待所有工作线程结束
        for (std::thread& worker : workers) {
            if (worker.joinable()) {
                worker.join();
            }
        }
    }

    //用 Lambda 替代 std::bind,避免类型推导问题
    template<class F, class... Args>
    auto enqueue(F&& f, Args&&... args)
        -> std::future<typename std::invoke_result<F, Args...>::type> {
        using return_type = typename std::invoke_result<F, Args...>::type;

        // 包装任务:用 Lambda 完美转发参数(替代 std::bind)
        auto task = std::make_shared<std::packaged_task<return_type()>>(
            [f = std::forward<F>(f), args = std::make_tuple(std::forward<Args>(args)...)]() mutable {
                return std::apply(std::move(f), std::move(args));
            }
        );

        std::future<return_type> res = task->get_future(); // 获取future
        {
            std::unique_lock<std::mutex> lock(queue_mutex);
            // 线程池已停止则不能提交任务
            if (stop) {
                throw std::runtime_error("向已停止的线程池提交任务");
            }
            tasks.emplace([task]() { (*task)(); }); // 放入任务队列
        }
        condition.notify_one(); // 唤醒一个等待的线程
        return res;
    }
};

   但是这里没实现优雅退出机制,若设置线程标志位为退出时,任务队列中还有剩余任务,则会直接舍弃,不符合标准。正确逻辑应该时线程退出时,若任务队列不为空,则消费完所有任务再退出。目前就先这样吧,后续有时间再继续优化

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐