使用C++11 STL线程库实现一个线程池。处理机制是抢占式的,即所有线程从一个队列(std::queue)中获取任务执行(计算字符串简单HASH值),使用std::mutex和std::conditional_variable实现队列访问并发协调。
1#include <iostream> 2#include <iomanip> 3#include <thread> 4#include <mutex> 5#include <string> 6#include <queue> 7#include <condition_variable> 8#include <algorithm> 9#include <sstream> 10 11using namespace std; 12 13static std::mutex G_lockPrint; 14 15void print_message(int value, const string& str) { 16 lock_guard<mutex> lock(G_lockPrint); 17 cout<<setw(8)<<right<<this_thread::get_id(); 18 cout<<setw(12)<<right<<value<<" "<<str<<endl; 19} 20 21#define THREAD_COUNT 10 22int main() 23{ 24 thread thpool[THREAD_COUNT]; 25 mutex quelock; 26 condition_variable quecv; 27 queue<string> strqueue; 28 29 volatile bool stop = false; 30 31 for(int i = 0; i < THREAD_COUNT; ++i ) { 32 thpool[i] = thread([&quelock, &quecv, &strqueue, &stop]() 33 { 34 string str; 35 while ( !stop ) { 36 { 37 unique_lock<mutex> lock(quelock); 38 if ( strqueue.empty() ) { 39 auto ret = quecv.wait_for(lock, chrono::seconds(1)); 40 if ( ret == cv_status::timeout) continue; 41 } 42 43 if ( !strqueue.empty() ) { 44 str = strqueue.front(); 45 strqueue.pop(); 46 } else { 47 continue; 48 } 49 } 50 51 int hash = 0; 52 for(size_t i = 0; i < str.length(); ++i) { 53 hash = (hash << 5) - i + str[i]; 54 } 55 print_message(hash, str); 56 } // end while 57 } 58 ); 59 } 60 61 for(int i = 0; i < 100000; ++i) { 62 stringstream ss; 63 ss<<"aaaaa_"<<i; 64 lock_guard<mutex> lock(quelock); 65 strqueue.push(ss.str()); 66 quecv.notify_one(); 67 } 68 69 while (1) { 70 this_thread::sleep_for(chrono::seconds(1)); 71 lock_guard<mutex> lock(quelock); 72 if ( strqueue.empty()) break; 73 } 74 stop = true; 75 for(int i = 0; i < THREAD_COUNT; ++i ) thpool[i].join(); 76 cout<<"program exit"<<endl; 77 return 0; 78}