#include "threadPool.h"
#include "Task.h"
#include "TaskQueue.h"
#include "TaskQueue.cpp" //éè¦å
嫿ºæä»¶ï¼å¦å模æ¿ç±» ç彿°å°±ä¼ æ¥é âæªå®ä¹çå¼ç¨â
#include
#include
#include
#include
#include
using namespace std;
template
ThreadPool::ThreadPool(int min, int max)
{
//å¯ä»¥ä½¿ç¨break æ¿ä»£ returnï¼å¨å½æ°å
é¨éåº
do
{
//å®ä¾åä»»å¡éå
taskQ=new TaskQueue;
if(taskQ==nullptr)
{
cout<<"malloc taskQ fail...\n";
break;
}
//åå§åç»æä½æå
threadIDs= new pthread_t[max];
if(threadIDs==nullptr)
{
//å
ååé
失败
cout<<"malloc pthreadIDs fail...\n";
break;
}
//åå§å线ç¨
memset(threadIDs,0,sizeof(pthread_t)*max);
minNum=min;
maxNum=max;
busyNum=0;
liveNum=min;
exitNum=0;
//åå§å线ç¨
if(pthread_mutex_init(&mutexPool,NULL)!=0 ||
pthread_cond_init(¬Empty,NULL)!=0)
{
//åå§å线ç¨å¤±è´¥
cout<<"mutex or condition init fail...\n";
//return NULL;
break;
}
//åå§æ¶ä¸éæ¯
shutDown=false;
//å建线ç¨
//è¿é manageréè¦æ¯éæå彿°ï¼æè
ä½ä¸º éææåï¼ä¸åå±äºå¯¹è±¡ï¼èæ¯å±äºç±»ï¼å°±æäºå°åï¼å°±å¯ä»¥ä¼ é彿°å°åäº
pthread_create(&managerID,NULL,manager,this); //管çè
线ç¨, 使ç¨thisä¼ ééæå®ä¾å¯¹è±¡ï¼managerå°±å¯ä»¥ç®¡çééææåäº
for(int i=0;i
void * ThreadPool::worker(void* arg)
{
//ç±»å转æ¢
ThreadPool* pool=static_cast(arg);
while(true)
{
//线ç¨ä½¿ç¨ä¹åå é
pthread_mutex_lock(&pool->mutexPool);
//夿å½åä»»å¡éå
while (pool->taskQ->taskNumber()==0&& !pool->shutDown)
{
//é»å¡å·¥ä½çº¿ç¨
pthread_cond_wait(&pool->notEmpty,&pool->mutexPool);
//鿝 é»å¡çå·¥ä½çº¿ç¨
if(pool->exitNum>0)
{
pool->exitNum--;
if (pool->liveNum>pool->minNum)
{
//åå°çº¿ç¨æ°
pool->liveNum--;
//è§£å¼äºæ¥é
pthread_mutex_unlock(&pool->mutexPool);
//éæ¯çº¿ç¨
pool->threadExit();
}
}
}
//å¤æçº¿ç¨æ± æ¯å¦è¢«å
³éäº
if(pool->shutDown)
{
//å
è§£é, é¿å
æ»é
pthread_mutex_unlock(&pool->mutexPool);
//éåºçº¿ç¨
pool->threadExit();
}
//ä»ä»»å¡éåä¸ååºä¸ä¸ªä»»å¡
Task task=pool->taskQ->takeTask();
//è§£é
pool->busyNum++;
//线ç¨ä½¿ç¨å®ä¹åè§£é
pthread_mutex_unlock(&pool->mutexPool);
cout<<"thread "<mutexPool); //线ç¨ä¼å¤è®¿é® å é
pool->busyNum--;
pthread_mutex_unlock(&pool->mutexPool); //è§£é
}
return NULL;
}
//管çè
线ç¨
template
void * ThreadPool::manager(void* arg)
{
//强å¶ç±»å转æ¢
ThreadPool* pool=static_cast(arg);
//æç
§ä¸å®çé¢çï¼æ£æµ å è°æ´çº¿ç¨ä¸ªæ°
while(!pool->shutDown) //çº¿ç¨æ± æªå
³éå°±æ£æµ
{
//æ¯é3S æ£æµä¸æ¬¡
sleep(3);
//ååºçº¿ç¨æ± ä¸ä»»å¡çæ°éåå½åçº¿ç¨æ± çæ°é, 鲿¢æå
¶ä»çº¿ç¨å¨åå
¥æ°æ®ï¼æä»¥éè¦ é
pthread_mutex_lock(&pool->mutexPool);
int queueSize=pool->taskQ->taskNumber();
int liveNum=pool->liveNum;
int busyNum=pool->busyNum;
pthread_mutex_unlock(&pool->mutexPool);
//æ·»å 线ç¨
//ä»»å¡ä¸ªæ°>åæ´»ç线ç¨ä¸ªæ° å¹¶ä¸ åæ´»ç线ç¨ä¸ªæ°<æå¤§çº¿ç¨æ° æ¶ï¼æä¼ç»§ç»å¢å æ°çº¿ç¨
if(queueSize>liveNum&&liveNummaxNum)
{
pthread_mutex_lock(&pool->mutexPool);
int counter=0;
//ä¸ä»
è¦å¨ä¸é¢å¤æï¼ä¹éè¦å¨ è¿é忬¡å¤æï¼é²æ¢å¨è¿æé´æ°éåéååï¼å¯¼è´ 个æ°åºç°é®é¢
for(int i=0; i < pool->maxNum && counter < NUMBER && pool->liveNum < pool->maxNum;i++)
{
//æ¾å°æªä½¿ç¨ç线ç¨ID(è¢«éæ¯ç线ç¨)
if(pool->threadIDs[i]==0)
{
//å建线ç¨, ç´æ¥ä½¿ç¨è¯¥çº¿ç¨ID
pthread_create(&pool->threadIDs[i],NULL,worker,pool);
counter++; //æ»çº¿ç¨ä¸ªæ°
pool->liveNum++; //åæ´»çº¿ç¨ä¸ªæ°
}
}
pthread_mutex_unlock(&pool->mutexPool);
}
//éæ¯çº¿ç¨
//éæ¯çº¿ç¨çæ¡ä»¶ï¼å¿çº¿ç¨*2<åæ´»çº¿ç¨ å¹¶ä¸ åæ´»ç线ç¨>æå°çº¿ç¨æ°
if(busyNum*2pool->minNum)
{
//æä½çº¿ç¨æ± ä¸çå¼é½è¦å é
pthread_mutex_lock(&pool->mutexPool);
//æ¯æ¬¡éæ¯åå¼ä¿ç两个线ç¨
pool->exitNum=NUMBER;
pthread_mutex_unlock(&pool->mutexPool);
//让线ç¨èªå·±ç»æ
for(int i=0;inotEmpty);
}
}
}
return NULL;
}
template
void ThreadPool::threadExit()
{
//è·åå½åç线ç¨ID
pthread_t tid=pthread_self();
for(int i=0;i
void ThreadPool::addTask(Task task)
{
//æ·»å ä»»å¡ä¸åéè¦éäºï¼å 为ææçæä½é½æ¯ç± TaskQueueæ¥æ§å¶äºï¼ä¸ä¼åçå²çª
if(shutDown)
{
return;
}
//æ·»å ä»»å¡
taskQ->addTask(task);
//å¤éé»å¡ç线ç¨
pthread_cond_signal(¬Empty); //å¤éæ¶è´¹è
}
//å½åçº¿ç¨æ± ä¸çº¿ç¨ä¸ªæ°
template
int ThreadPool::getBusyNum()
{
//è¿éæå¯è½ä¼è¢«åæ¶è®¿é®ï¼æä»¥éè¦å é
pthread_mutex_lock(&mutexPool);
int busyNum=this->busyNum;
pthread_mutex_unlock(&mutexPool);
return busyNum;
}
//è·åå½åå建ç线ç¨ä¸ªæ°
template
int ThreadPool::getAliveNum()
{
pthread_mutex_lock(&mutexPool);
int aliveNum=this->liveNum;
pthread_mutex_unlock(&mutexPool);
return aliveNum;
}
template
ThreadPool::~ThreadPool()
{
//å
³éçº¿ç¨æ±
shutDown=true;
//åæ¶é»å¡ç®¡çè
线ç¨
pthread_join(managerID,NULL);
//å¤éé»å¡çæ¶è´¹è
线ç¨
for(int i=0;i