C语言线程池详解1. 线程池核心部分的分析(1) 线程同步、生产者和消费者(2) 任务和任务队列(3)工作线程(4) 管理者线程(5) 线程池2. cthreadpool.h3. cthreadpool.c3. 测试程序main.c
任务队列:存储任务。
工作线程:子线程,根据和任务数对比来动态增减数量,作用是取出任务并处理任务。
管理者线程:对比任务数和线程数,执行创建与销毁线程。
线程池:管理任务、线程和整个线程池的运行。
主线程生产任务,子线程消费任务,需使用锁进行线程同步。
生产者:主线程往线程池的任务队列添加任务,需加锁,并使用条件变量判断任务队列是否已满:
任务队列已满,解锁并阻塞线程;
任务队列未满,往队尾添加任务,解锁并pthread_cond_signal通知任务队列不为空。
消费者:子线程循环从线程池的任务队列取出任务,需加锁,并使用条件变量盘对任务队列是否为空:
任务队列不为空:于对首取出任务,解锁并pthread_cond_signal通知任务队列未满;
任务队列已空:解锁并阻塞子线程。
任务:工作任务函数,执行任务的过程就是调用函数的过程。
任务队列:包括参数:当前总的任务数,最大任务数量,提供接口添加和删除任务。
xxxxxxxxxx//1. 任务和任务队列//定义一个表示线程池任务的结构体,包含一个函数指针,和一个参数 typedef struct PoolTask { void (*taskFunc)(void*); void* taskArg; }PoolTask; PoolTask* taskQueue; //这是一个任务队列,用数组表示 int tasksMax; //限制任务队列的任务最大数 int tasksCurr; //当前有多少任务 int queueFront; //队首位置 int queueBack; //队尾位置//2. 任务队列初始化 pool->taskQueue = (PoolTask*)malloc(sizeof(PoolTask) * maxNumTasks); if(NULL == pool->taskQueue) { break; //跳出do-while循环 } pool->tasksMax = maxNumTasks; pool->tasksCurr = 0; pool->queueFront = 0; pool->queueBack = 0; //3. 添加任务 pool->taskQueue[pool->queueBack].taskFunc = task->taskFunc; pool->taskQueue[pool->queueBack].taskArg = task->taskArg; pool->queueBack = (pool->queueBack + 1) % pool->tasksMax; pool->tasksCurr++;//4. 获取并执行任务 PoolTask nowTask; nowTask.taskFunc = pool->taskQueue[pool->queueFront].taskFunc; nowTask.taskArg = pool->taskQueue[pool->queueFront].taskArg; pool->queueFront = (pool->queueFront + 1) % pool->tasksMax; pool->tasksCurr--; nowTask.taskFunc(nowTask.taskArg);! 每个线程都在对任务队列的操作,故需注意同步。
工作线程,主要负责处理任务,线程池维护着一定数量的线程数,这些线程是动态的,根据任务数量而变化。
C语言:
xxxxxxxxxx//1. 定义 pthread_t* threadsWorker; //用于表示线程数组 int threadsMax; //最大线程数量 int threadsMin; //最小线程数量 int threadsAll; //当前线程池的所有线程数量,包括工作中和闲的线程数 int threadsWorking; //当前线程池的工作中的线程数量 int threadsClean; //用于存储需要销毁的线程数量//2. 线程初始化 pool->threadsMax = maxNumThreads; pool->threadsMin = minNumThreads; pool->threadsAll = minNumThreads; //初始化时,当前线程数 = 最小线程数 pool->threadsWorking = 0; pool->threadsClean = 0; pool->threadsWorker = (pthread_t*)calloc(maxNumThreads, sizeof(pthread_t)); if(NULL == pool->threadsWorker) { break; } for(int i = 0; i < minNumThreads; ++i) { int ret = pthread_create(&pool->threadsWorker[i], NULL, comsumeTask, pool); //初始化时,新建minNumThreads个线程 if(ret != 0) { // perror("pthread_create"); break; } }管理者并不需要处理任务,其根据线程数量和任务数量,动态创建和销毁线程。
xxxxxxxxxx //管理者线程初始化 int ret = pthread_create(&pool->threadManager, NULL, manager, pool); if(ret != 0) { break; }/* 管理者工作函数*/void *manager(void *arg){ //(1)while循环,每隔几秒检测需不需要创建和销毁线程 //(2)获得当前总线程数和总任务数 //(3)获得正在工作的线程数量 //(4)添加线程,当任务数 > 线程数,且线程数不大于线程最大数 if(countTask > countThreadAll && countThreadAll < pool->threadsMax) // (5)销毁线程:当工作中线程数 * 2 < 总线程数, 且线程数 > 最小线程数 if(countThreadWorking * 2 > countThreadAll && countThreadAll > pool->threadsMin)}线程池管理着整个线程池的运转,包括:任务队列及其参数(标志位)、工作中的线程及其参数、管理者线程、线程同步资源等。
xxxxxxxxxx//1. 定义线程池typedef struct ThreadPool{ //任务 //线程:工作线程 pthread_t* threadsWorker; //用于表示线程数组 int threadsMax; //最大线程数量 int threadsMin; //最小线程数量 int threadsAll; //当前线程池的所有线程数量,包括工作中和闲的线程数 int threadsWorking; //当前线程池的工作中的线程数量 int threadsClean; //用于存储需要销毁的线程数量 //线程:管理者线程 pthread_t threadManager; //线程同步 pthread_mutex_t mutexPool; //锁,整个线程池 pthread_mutex_t mutexWorking; //锁,工作中的线程数 pthread_cond_t condTaskFull; //条件变量,任务数满了 pthread_cond_t condTaskEmpty; //条件变量,任务数空了 bool offPool; //关闭线程池}ThreadPool;
//2. 线程池函数 //创建线程池并初始化ThreadPool *createThreadPool(int minNumThreads, int maxNumThreads, int maxNumTasks); //销毁线程池void clean ThreadPool(ThreadPool *pool);//将任务添加至线程池void produceTask(ThreadPool *pool, PoolTask* task);//工作子线程函数void *comsumeTask(void *arg);//管理者子线程函数void *manager(void *arg);//销毁单个线程函数void cleanOneThread(ThreadPool *pool);xxxxxxxxxx
typedef struct PoolTask{ void (*taskFunc)(void*); void* taskArg;}PoolTask;
typedef struct ThreadPool{ //任务 PoolTask* taskQueue; int tasksMax; int tasksCurr; int queueFront; int queueBack; //线程 //工作线程 pthread_t* threadsWorker; int threadsMax; int threadsMin; int threadsAll; int threadsWorking; int threadsClean; //管理者线程 pthread_t threadManager; //线程同步 pthread_mutex_t mutexPool; pthread_mutex_t mutexWorking; pthread_cond_t condTaskFull; pthread_cond_t condTaskEmpty; bool offPool; }ThreadPool;ThreadPool *createThreadPool(int minNumThreads, int maxNumThreads, int maxNumTasks); void cleanThreadPool(ThreadPool *pool);void produceTask(ThreadPool *pool, PoolTask* task);void *comsumeTask(void *arg);void *manager(void *arg);void cleanOneThread(ThreadPool *pool);
// !_CTHREADPOOL_Hxxxxxxxxxx
ThreadPool *createThreadPool(int minNumThreads, int maxNumThreads, int maxNumTasks){ ThreadPool* pool = (ThreadPool*)malloc(sizeof(ThreadPool)); do { if(NULL == pool) { // perror("malloc"); break; } //任务队列初始化 pool->taskQueue = (PoolTask*)malloc(sizeof(PoolTask) * maxNumTasks); if(NULL == pool->taskQueue) { // perror("malloc"); break; } pool->tasksMax = maxNumTasks; pool->tasksCurr = 0; pool->queueFront = 0; pool->queueBack = 0; //工作线程 pool->threadsMax = maxNumThreads; pool->threadsMin = minNumThreads; pool->threadsAll = minNumThreads; pool->threadsWorking = 0; pool->threadsClean = 0; pool->threadsWorker = (pthread_t*)calloc(maxNumThreads, sizeof(pthread_t)); if(NULL == pool->threadsWorker) { break; } for(int i = 0; i < minNumThreads; ++i) { int ret = pthread_create(&pool->threadsWorker[i], NULL, comsumeTask, pool); if(ret != 0) { // perror("pthread_create"); break; } } //管理者线程 int ret = pthread_create(&pool->threadManager, NULL, manager, pool); if(ret != 0) { break; } //线程同步初始化 if( pthread_mutex_init(&pool->mutexPool, NULL) != 0 || pthread_mutex_init(&pool->mutexWorking, NULL) != 0 || pthread_cond_init(&pool->condTaskFull, NULL) != 0 || pthread_cond_init(&pool->condTaskEmpty, NULL) != 0 ) { break; } pool->offPool = false;
return pool; //初始化没问题,可以return了 } while(0); if(pool->threadsWorker != NULL) { free(pool->threadsWorker); pool->threadsWorker = NULL; } if(pool->taskQueue != NULL) { free(pool->taskQueue); pool->taskQueue = NULL; } if(pool != NULL) { free(pool); pool = NULL; } return NULL; //初始化有问题}void cleanThreadPool(ThreadPool *pool){ if(NULL == pool) { return; } pool->offPool = true; //线程回收 pthread_join(pool->threadManager, NULL); for(int i = 0; i < pool->threadsAll; ++i) { pthread_cond_signal(&pool->condTaskEmpty); } sleep(20); //内存回收 if(pool->taskQueue != NULL) { free(pool->taskQueue); pool->taskQueue = NULL; } if(pool->threadsWorker != NULL) { free(pool->threadsWorker); pool->threadsWorker = NULL; }
pthread_mutex_destroy(&pool->mutexPool); pthread_mutex_destroy(&pool->mutexWorking); pthread_cond_destroy(&pool->condTaskEmpty); pthread_cond_destroy(&pool->condTaskFull);
if(pool != NULL) { free(pool); pool = NULL; }}//生产者:主线程往队列添加任务void produceTask(ThreadPool *pool, PoolTask* task){ if(NULL == pool) { return; } pthread_mutex_lock(&pool->mutexPool); while(pool->tasksMax == pool->tasksCurr && !pool->offPool) { pthread_cond_wait(&pool->condTaskFull, &pool->mutexPool); } if(pool->offPool) { pthread_mutex_unlock(&pool->mutexPool); return; } pool->taskQueue[pool->queueBack].taskFunc = task->taskFunc; pool->taskQueue[pool->queueBack].taskArg = task->taskArg; pool->queueBack = (pool->queueBack + 1) % pool->tasksMax; pool->tasksCurr++;
pthread_cond_signal(&pool->condTaskEmpty); pthread_mutex_unlock(&pool->mutexPool);}void *comsumeTask(void *arg){ ThreadPool* pool = (ThreadPool*)arg; while (1) { pthread_mutex_lock(&pool->mutexPool); while (pool->tasksCurr == 0 && !pool->offPool) { pthread_cond_wait(&pool->condTaskEmpty, &pool->mutexPool); if(pool->threadsClean > 0)//是否需要销毁一个子线程 { pool->threadsClean--; if(pool->threadsAll > pool->threadsMin) { pool->threadsAll--; pthread_mutex_unlock(&pool->mutexPool); cleanOneThread(pool); } } } if(pool->offPool) { pthread_mutex_unlock(&pool->mutexPool); cleanOneThread(pool); } //获取任务 PoolTask nowTask; nowTask.taskFunc = pool->taskQueue[pool->queueFront].taskFunc; nowTask.taskArg = pool->taskQueue[pool->queueFront].taskArg; pool->queueFront = (pool->queueFront + 1) % pool->tasksMax; pool->tasksCurr--; pthread_cond_signal(&pool->condTaskFull); //唤醒被阻塞的生产者线程 pthread_mutex_unlock(&pool->mutexPool);
//执行任务 pthread_mutex_lock(&pool->mutexWorking); pool->threadsWorking++; pthread_mutex_unlock(&pool->mutexWorking);
// printf("*****child thread ID: %ld, comsuming task!*****\n", pthread_self()); nowTask.taskFunc(nowTask.taskArg); pthread_mutex_lock(&pool->mutexWorking); pool->threadsWorking--; pthread_mutex_unlock(&pool->mutexWorking);
} return NULL;}//管理者:根据任务数和当前正在工作的线程数对比,动态增减线程void *manager(void *arg){ ThreadPool* pool = (ThreadPool*)arg; //每隔3秒检测一次 while(!pool->offPool) { sleep(3); //取出任务数,线程数 pthread_mutex_lock(&pool->mutexPool); int countTask = pool->tasksCurr; int countThreadAll = pool->threadsAll; pthread_mutex_unlock(&pool->mutexPool);
pthread_mutex_lock(&pool->mutexWorking); int countThreadWorking = pool->threadsWorking; pthread_mutex_unlock(&pool->mutexWorking); //增加线程数 if(countTask > countThreadAll && countThreadAll < pool->threadsMax) { pthread_mutex_lock(&pool->mutexPool); int countAdd = 0; for(int i = 0; i < pool->threadsMax && countAdd < 2 && pool->threadsAll < pool->threadsMax; ++i) //每次增加2个线程 { if(0 == pool->threadsWorker[i]) { pthread_create(&pool->threadsWorker[i], NULL, comsumeTask, pool); countAdd++; pool->threadsAll++; } } pthread_mutex_unlock(&pool->mutexPool); } //减少线程数 if(countThreadWorking * 2 > countThreadAll && countThreadAll > pool->threadsMin) { pthread_mutex_lock(&pool->mutexPool); pool->threadsClean = 2; //每次销毁2个线程 pthread_mutex_unlock(&pool->mutexPool); for(int i = 0; i < 2; ++i) { pthread_cond_signal(&pool->condTaskEmpty); } } } return NULL;}void cleanOneThread(ThreadPool *pool){ if(NULL == pool) { return; } pthread_t tid = pthread_self(); for(int i = 0; i < pool->threadsMax; ++i) { if(tid == pool->threadsWorker[i]) { pool->threadsWorker[i] = 0; break; } } pthread_exit(NULL);}xxxxxxxxxx
void taskTest(void* arg){ int* p = (int*)arg; (*p)++; printf("result-----child thread %ld at taskTest function, the arg: %d-----\n", pthread_self(), *p); sleep(1);}
int testArg = 0;
int main(){ ThreadPool* pool = createThreadPool(2, 6, 100); for(int i = 0; i < 100; ++i) { PoolTask task; task.taskArg = &testArg; task.taskFunc = taskTest; produceTask(pool, &task); }
sleep(50); //主线程不要提前关闭 cleanThreadPool(pool); printf("The result: testArg = %d\n", testArg); //输出最终结果 return 0;}