一、概念讲解

线程池使用了池化技术,以空间换时间。线程池是线程的一种使用模式。线程过多会带来调度开销,影响整体性能。而线程池维护者多个线程,等待着被派发可并发执行的任务。

这避免了在处理短时间任务时创建与销毁线程的代价,还能保证内核被充分利用,防止过度调度

应用场景:

(1)需要大量线程来完成任务,且完成任务的时间较短。例如:web服务器完成网页请求

(2)对性能要求苛刻的应用。例如:要求服务器快速响应用户请求

(3)接收突发性的大量请求,但不至于使服务器因此产生大量线程的应用

线程池运作过程:

(1)线程池创建固定数量的线程,循环从任务队列中获取任务对象
(2)获取到任务对象之后,执行任务对象的任务接口

如上就是线程池的原理了,那么线程池首先会在自己内部创建出一批线程,那么线程被创建出来之后会进行检测任务队列,如果任务队列没有任务,那么线程此时就没有任务要执行了,所以线程就去条件变量下去等待了,那么多个线程在最开始任务队列没有任务,所以就会依次全部都去条件变量下进行等待

当任务队列来任务的时候,此时线程池的线程都已经在条件变量下进行等待了,所以此时当任务入队列之后,就要唤醒在条件变量下进行等待的一个线程,让这个线程执行任务,当这个线程执行完成之后又会去检测任务队列有没有任务,如果没有它就继续会去条件变量下进行等待,以此循环往复下运行

并且任务队列只有一份,所以任务队列是一个共享资源,既然是共享资源,那么被多线程同时访问的时候就会有问题,那么就要干什么?对的通过互斥锁保护共享资源,所以当任务队列入任务的时候,应该申请锁,当任务队列入完任务的时候,应该释放锁。当条件变量检测任务队列执行任务队列的任务的时候应该申请锁,当执行完毕应该释放锁

并且当最开始任务队列没有任务的时候,线程池的线程只能在条件变量下进行等待,等待线程将任务入任务队列。那么这就表现出一定的顺序性了,即同步\n所以这不就符合生产者消费者模型的321原则了

3种关系:(1)生产者与生产者:互斥(2)消费者与消费者:互斥(3)生产者与消费者:互斥,同步。

2种角色:(1)生产者(2)消费者。

1个场所:(1)特定结构的内存空间\n其中的1个场所是任务队列,2种角色,线程池内的线程是消费者,向任务队列入任务的线程是生产者,3种关系,由于只有1个生产者,不存在多个生产者,所以(1)自然而然不需要满足,(2)多个消费者之间的关系通过互斥锁保证了互斥(3)生产者与消费者的互斥通过互斥锁保证,同步通过条件变量来保证,所以线程池的实现就是基于生产者消费者模型实现的

代码实现:

Task.hpp

这个任务模块和之前的生产者消费者的任务是一样的

#pragma once
#include<iostream>
#include<string>
using namespace std;
string opers="+-*/%";
enum{
    DivZero=1,
    ModZero,
    Unknown
};
class Task
{
public:
    Task(int x,char op,int y)
    :data1_(x)
    ,oper_(op)
    ,data2_(y)
    ,result_(0)
    ,exitcode_()
    {}
    void run()
    {
        switch(oper_)
        {
            case '+':
            result_=data1_+data2_;
            break;
            case '-':
            result_=data1_-data2_; 
            break;
            case '*':
            result_=data1_*data2_;
            break;
            case '/':
            {
                if(data2_==0) exitcode_=ModZero;
                result_=data1_/data2_;
            }
            break;
            case '%':
            {
                if(data2_==0) exitcode_=DivZero;
                else result_=data1_%data2_;
            }
            break;
            default:
            exitcode_=Unknown;
            break;
        }
    }
    void operator()()
    {
        run();
    }
    string GetResult()
    {
        string r=to_string(data1_);
        r+=oper_;
        r+=to_string(data2_);
        r+="=";
        r+=to_string(result_);
        r+="[code";
        r+=to_string(exitcode_);
        r+="]";
        return r;
    }
    string GetTask()
    {
        string r=to_string(data1_);
        r+=oper_;
        r+=to_string(data2_);
        r+="=?";
        return r;
    }
    ~Task()
    {}
private:
    int data1_;
    int data2_;
    char oper_;
    int result_;
    int exitcode_;
};

ThreadPool.hpp

定义一个类来保存线程的信息

struct threadinfo
{
    pthread_t tid;
    string threadname;
};

我们需要实现一个任务队列,一把锁,一个条件变量,一个vector数组来存储线程信息

    queue<T> _tasks;
    vector<threadinfo> _threads;
    pthread_mutex_t _mutex;
    pthread_cond_t _cond;

基本功能的封装(基本功能的使用希望在类内使用,设置为私有)

void Lock()
    {
        pthread_mutex_lock(&_mutex);
    }
    void unlock()
    {
        pthread_mutex_unlock(&_mutex);
    }
    void WakeUp()
    {
        pthread_cond_signal(&_cond);
    }
    void ThreadSleep()
    {
        pthread_cond_wait(&_cond,&_mutex);
    }
    bool is_empty()
    {
        return _tasks.empty();
    }
    string GetThreadname(pthread_t tid)
    {
        for(auto it:_threads)
        {
            if(it.tid==tid)
            {
                return it.threadname;
            }
        }
        return "None";
    }

添加任务和取出任务

因为并不知道调用者是在多线程还是单线程的情况来调用push接口的,所有需要加锁保护。

pop接口后续是在线程安全的情况下执行的,所有不需要加锁。

void push(const T&in)
    {
        Lock();
        _tasks.push(in);
        WakeUp();
        unlock();
    }
    T pop()
    {
        T out=_tasks.front();
        _tasks.pop();
        return out;
    }

Start接口(创建多线程,来处理任务池中的任务)

void Start()
    {
        for(int i=0;i<defaultnum;i++)
        {
            _threads[i].threadname="thread - "+to_string(i);
            pthread_create(&(_threads[i].tid),nullptr,handletask,this);
        }
    }

handletask(线程执行函数)

因为是在类内,所以会有一个隐含的this指针,但是不符合线程函数的格式,所以设置为静态成员函数,就没有this指针,后续调用成员函数就需要通过传进来的this指针来访问。

所以接下来要判断访问共享资源任务队列中是否有任务了,所以我们要先申请锁,然后再使用while循环判断任务队列是否为空,如果为空那么则让线程去在条件变量下等待即可,这里采用while循环可以有效避免线程伪唤醒的情况。

那么走到下一步一定是当前任务队列不为空,或者线程由于任务队列加入了任务被唤醒的情况,所以线程就可以出Pop出任务了,Pop的返回值就是任务,所以我们接收任务即可,然后直接释放锁,没错,释放锁即可,因为此时线程已经拿到了任务了,已经访问完成了任务队列拿到任务了,当前线程已经在独立栈上有了应该执行什么样的任务如何执行任务了,目前已经不会对共享资源访问了,所以释放锁即可,然后执行任务,打印对应的线程名称name以及处理结果即可。

static void*handletask(void*args)
    {
        ThreadPool*tp=static_cast<ThreadPool<T>*>(args);
        string threadname=tp->GetThreadname(pthread_self());
        while(true)
        {
            tp->Lock();
            while(tp->is_empty())
            {
                tp->ThreadSleep();
            }
            T t=tp->pop();
            tp->unlock();
            t();
            cout << threadname << " run " << " result :" << t.GetResult() << endl;
        }
    }

main.cpp

创建线程池类,执行线程函数,主线程while循环一直生产任务,线程池会不断检测任务,有任务就会去执行。

这里我每次生产一个任务就会sleep1秒,所以最后呈现出来的结果是每次生产1一个任务,就会被马上执行。

#include <iostream>
#include <ctime>
#include <unistd.h>
#include "Task.hpp"
#include "ThreadPool.hpp"

int main()
{
    srand(time(nullptr));
    int len = opers.size();

    ThreadPool<Task>* tp = new ThreadPool<Task>();

    tp->Start();
    
    while(true)
    {
        int data1 = rand() % 10;
        usleep(10);
        int data2 = rand() % 5;
        char op = opers[rand() % len];

        Task t(data1,op,data2);

        tp->push(t);

        std::cout << "main thread make task: " << t.GetTask() << std::endl;
        sleep(1);
    }

    delete tp;

    return 0;
}

运行结果:

代码源码:

Task.hpp上面已经有了。

ThreadPool.hpp

#include <iostream>
#include <vector>
#include <pthread.h>
#include <string>
#include <unistd.h>
#include <queue>
using namespace std;
const int defaultnum=5;
struct threadinfo
{
    pthread_t tid;
    string threadname;
};
template<class T>
class ThreadPool
{
private:
    void Lock()
    {
        pthread_mutex_lock(&_mutex);
    }
    void unlock()
    {
        pthread_mutex_unlock(&_mutex);
    }
    void WakeUp()
    {
        pthread_cond_signal(&_cond);
    }
    void ThreadSleep()
    {
        pthread_cond_wait(&_cond,&_mutex);
    }
    bool is_empty()
    {
        return _tasks.empty();
    }
    string GetThreadname(pthread_t tid)
    {
        for(auto it:_threads)
        {
            if(it.tid==tid)
            {
                return it.threadname;
            }
        }
        return "None";
    }
    //设置静态,没有this指针,里面接口都需要使用传进来的tp指针来使用
    static void*handletask(void*args)
    {
        ThreadPool*tp=static_cast<ThreadPool<T>*>(args);
        string threadname=tp->GetThreadname(pthread_self());
        while(true)
        {
            tp->Lock();
            while(tp->is_empty())
            {
                tp->ThreadSleep();
            }
            T t=tp->pop();
            tp->unlock();
            t();
            cout << threadname << " run " << " result :" << t.GetResult() << endl;
        }
    }
public:
    ThreadPool(int num=defaultnum)
    :_threads(num)
    {
        pthread_mutex_init(&_mutex,nullptr);
        pthread_cond_init(&_cond,nullptr);
    }
    void Start()
    {
        for(int i=0;i<defaultnum;i++)
        {
            _threads[i].threadname="thread - "+to_string(i);
            pthread_create(&(_threads[i].tid),nullptr,handletask,this);
        }
    }
    void push(const T&in)
    {
        Lock();
        _tasks.push(in);
        WakeUp();
        unlock();
    }
    T pop()
    {
        T out=_tasks.front();
        _tasks.pop();
        return out;
    }
    ~ThreadPool()
    {
        pthread_mutex_destroy(&_mutex);
        pthread_cond_destroy(&_cond);
    }
private:
    queue<T> _tasks;
    vector<threadinfo> _threads;
    pthread_mutex_t _mutex;
    pthread_cond_t _cond;
};

main.cpp

#include <iostream>
#include <ctime>
#include <unistd.h>
#include "Task.hpp"
#include "ThreadPool.hpp"

int main()
{
    srand(time(nullptr));
    int len = opers.size();

    ThreadPool<Task>* tp = new ThreadPool<Task>();

    tp->Start();
    
    while(true)
    {
        int data1 = rand() % 10;
        usleep(10);
        int data2 = rand() % 5;
        char op = opers[rand() % len];

        Task t(data1,op,data2);

        tp->push(t);

        std::cout << "main thread make task: " << t.GetTask() << std::endl;
        sleep(1);
    }

    delete tp;

    return 0;
}

单例模式

单例模式是一种常用的,经典的设计模式之一,是IT界的很多大佬针对一些常见的,经典的场景设计出来的解决方案,即设计模式

单例模式的特点就是一个类只能实例化出一个对象。例如在服务器开发的场景中,需要一个服务器将很多数据加载到内存,但是只需要一个服务器,所以就可以使用单例的类来构建服务器管理数据

单例模式有两种实现方式分别为懒汉模式和饿汉模式。

懒汉模式最核心的就是延时加载,从而可以优化服务器的启动速度

单例模式优化

基本思路:

1.公有静态获取接口:来给外部提供获取实例化对象的接口。

static ThreadPool<T> *GetInstance() // 懒汉模式
    {
        pthread_mutex_lock(&lock_); // 避免出现多个线程抢夺同一份资源
        if (tp_ == nullptr)
        {
            tp_ = new ThreadPool<T>();
        }
        pthread_mutex_unlock(&lock_);
    }

如果使用者在使用这个接口的情况是在多线程的情况下,是会有线程安全的问题的,多线程同时进入 tp_==nullptr 的判断分支,因 new 对象和给 tp_赋值 的指令未原子执行,线程切换后后续线程仍判定 tp_ 为空并重复创建对象;后续线程的赋值操作会覆盖 tp_ 原有地址,导致先创建的对象无指针指向,造成内存泄漏,彻底破坏单例唯一性。

所以我们需要进行加锁保护,每次只有一次线程可以进行访问_tp变量。

但是这样子设计还是存在缺陷,因为多线程情况下面,如果_tp已经被实例化了,但是线程并不知道,还是回去申请锁,和释放锁,锁的申请和释放都是需要时间的,会非常影响效率。

要注意:这个锁需要是静态的,因为静态成员函数只可以调用静态成员对象。

static ThreadPool<T> *GetInstance() // 懒汉模式
    {
        if (tp_ == nullptr)//减少申请锁和释放锁的次数
        {
            pthread_mutex_lock(&lock_); // 避免出现多个线程抢夺同一份资源
            while (tp_ == nullptr)
            {
                tp_ = new ThreadPool<T>();
            }
            pthread_mutex_unlock(&lock_);
        }
        return tp_;
    }

可以在外面再加一次判断,如果被实例化了,就直接返回。

2.私有静态成员变量:静态函数中只能调用静态成员对象,同时静态成员对象,在整个程序中只有一份也符合单例模式中一个类只能实例化出一个对象。

private:
    vector<ThreadInfo> threads_;
    queue<T> tasks_;
    pthread_mutex_t mutex_;
    pthread_cond_t cond_;
    static pthread_mutex_t lock_;
    static ThreadPool<T> *tp_; // 懒汉模式
};
template <class T>
ThreadPool<T> *ThreadPool<T>::tp_ = nullptr;
template <class T>
pthread_mutex_t ThreadPool<T>::lock_ = PTHREAD_MUTEX_INITIALIZER;

3.私有构造函数:禁止外部通过new创建对象,保证单例的唯一性。

private:
    ThreadPool(int num = defaultnum)
        : threads_(num)
    {
        pthread_mutex_init(&mutex_, nullptr);
        pthread_cond_init(&cond_, nullptr);
    }
    ThreadPool(const ThreadPool<T>&)=delete;
    const ThreadPool<T>&operator=(const ThreadPool<T>&)=delete;
    ~ThreadPool()
    {
        pthread_mutex_destroy(&mutex_);
        pthread_cond_destroy(&cond_);
    }

单例模式优化源码

main.cpp

#include"ThreadPool.hpp"
#include"Task.hpp"
#include<iostream>
using namespace std;
int main()
{
    cout<<"process run >>>>"<<endl;
    sleep(3);
    ThreadPool<Task>::GetInstance()->Start();
    srand(time(nullptr)^getpid());
    while(true)
    {
        int x=rand()%10+1;
        usleep(10);
        int y=rand()%10;
        char op=opers[rand()%opers.size()];
        Task t(x,op,y);
        ThreadPool<Task>::GetInstance()->Push(t);
        cout<<"main thread make task: "<<t.GetResult()<<endl;
        sleep(1);
    }
    return 0;
}

只有我们使用到了ThreadPool里面的成员接口,GetInstance才会去实例化。

ThreadPool.hpp

#include <iostream>
#include <vector>
#include <pthread.h>
#include <string>
#include <unistd.h>
#include <queue>
using namespace std;
struct ThreadInfo
{
    pthread_t tid;
    string name;
};
static const int defaultnum = 5;
template <class T>
class ThreadPool
{
public:
    void Lock()
    {
        pthread_mutex_lock(&mutex_);
    }
    void UnLock()
    {
        pthread_mutex_unlock(&mutex_);
    }
    void WakeUp()
    {
        pthread_cond_signal(&cond_);
    }
    void ThreadSleep()
    {
        pthread_cond_wait(&cond_, &mutex_);
    }
    bool IsQueueEmpty()
    {
        return tasks_.empty();
    }
    string GetThreadName(pthread_t tid)
    {
        for (auto &t : threads_)
        {
            if (t.tid == tid)
            {
                return t.name;
            }
        }
        return "None";
    }

public:
    static void *HandlerTask(void *args)
    {
        ThreadPool *tp = static_cast<ThreadPool<T>*>(args);
        string name = tp->GetThreadName(pthread_self());
        while (true)
        {
            tp->Lock();
            while (tp->IsQueueEmpty())
            {
                tp->ThreadSleep();
            }
            T t = tp->Pop();
            tp->UnLock();
            t();
            cout << name << " run " << " result :" << t.GetResult() << endl;
        }
    }
    void Start()
    {
        int len = threads_.size();
        for (int i = 0; i < len; i++)
        {
            threads_[i].name = "thread -" + to_string(i);
            pthread_create(&(threads_[i].tid), nullptr, HandlerTask, this); // this指针指向的是调用成员函数的对象(ThreadPool)
        }
    }
    T Pop()
    {
        T t = tasks_.front();
        tasks_.pop();
        return t;
    }
    void Push(const T &in)
    {
        Lock();
        tasks_.push(in);
        WakeUp();
        UnLock();
    }
    static ThreadPool<T> *GetInstance() // 懒汉模式
    {
        if (tp_ == nullptr)//减少申请锁和释放锁的次数
        {
            pthread_mutex_lock(&lock_); // 避免出现多个线程抢夺同一份资源
            while (tp_ == nullptr)
            {
                tp_ = new ThreadPool<T>();
            }
            pthread_mutex_unlock(&lock_);
        }
        return tp_;
    }

private:
    ThreadPool(int num = defaultnum)
        : threads_(num)
    {
        pthread_mutex_init(&mutex_, nullptr);
        pthread_cond_init(&cond_, nullptr);
    }
    ThreadPool(const ThreadPool<T>&)=delete;
    const ThreadPool<T>&operator=(const ThreadPool<T>&)=delete;
    ~ThreadPool()
    {
        pthread_mutex_destroy(&mutex_);
        pthread_cond_destroy(&cond_);
    }

private:
    vector<ThreadInfo> threads_;
    queue<T> tasks_;
    pthread_mutex_t mutex_;
    pthread_cond_t cond_;
    static pthread_mutex_t lock_;
    static ThreadPool<T> *tp_; // 懒汉模式
};
template <class T>
ThreadPool<T> *ThreadPool<T>::tp_ = nullptr;
template <class T>
pthread_mutex_t ThreadPool<T>::lock_ = PTHREAD_MUTEX_INITIALIZER;

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐