C++多线程学习之生产者-消费者模型和线程池实现
前面几篇学习了线程中的互斥体mutex, lock_guard, unique_lock, scoped_lock(C17), shared_lock(C14),condition_variable等相关示例。本篇结合这些知识实现生产者-消费者模型示例以及线程池的实现。生产者-消费者源码//producer_consumer.cpp 生产者消费者模型 example //#include iostream // std::cout #include thread // std::thread #include mutex // std::mutex #include vector //std::vector #include condition_variable std::mutex mtx; std::condition_variable cv_msg_empty; //条件变量消息为空 std::condition_variable cv_msg_full; //条件变量消息满了 volatile int msg_count 10; std::vectorint vec; void consumer(int id) //消费者 { while (true) { std::unique_lockstd::mutex lock(mtx); while (msg_count 0) { printf( consumer cv_msg_full.waitstart%d\n, msg_count); cv_msg_full.wait(lock); printf( consumer cv_msg_full.waitend%d\n, msg_count); } printf(consumer1id%d\n, id); //读取10条消息 std::this_thread::sleep_for(std::chrono::seconds(1)); //暂停2秒 msg_count 10; //消费10条消息 for (int msg : vec) { printf(consumermsg_id%d\n, msg); } vec.clear(); cv_msg_empty.notify_one(); } printf(threadid%d\n, id); } //生产者 void producer() { int number{ 1000 }; while (true) { std::unique_lockstd::mutex lock(mtx); while (msg_count 0) { printf( producer cv_msg_empty.waitstart%d\n, msg_count); cv_msg_empty.wait(lock); printf( producer cv_msg_empty.waitend%d\n, msg_count); } vec.push_back(number); //生产10条消息 msg_count--; std::this_thread::sleep_for(std::chrono::seconds(1)); //暂停2秒 printf(producermsg_count%d\n, msg_count); if (msg_count 0) //消息满了 { cv_msg_full.notify_all();//通知所有线程 } } } int main() { //1.记录开始时间点 auto start std::chrono::steady_clock::now(); printf(main1\n); std::thread consume(consumer, 1); std::thread produce(producer); consume.join(); produce.join(); printf(main2\n); //2.记录结束时间点 auto end std::chrono::steady_clock::now(); //3.计算时间差 auto elapsed end - start; printf(耗时: %lld 纳秒\n, std::chrono::duration_caststd::chrono::nanoseconds(elapsed).count()); printf(耗时: %lld 微秒\n, std::chrono::duration_caststd::chrono::microseconds(elapsed).count()); printf(耗时: %lld 毫秒\n, std::chrono::duration_caststd::chrono::milliseconds(elapsed).count()); printf(耗时: %lld 秒\n, std::chrono::duration_caststd::chrono::seconds(elapsed).count()); printf(hello learn producer_consumer\n); return 0; }编译运行2.线程池实现ThreadPool.h文件//ThreadPool.h文件 #ifndef THREADPOOL_H #define THREADPOOL_H #include vector #include queue #include memory #include thread #include mutex #include condition_variable #include future #include functional #include stdexcept #include type_traits class ThreadPool { public: ThreadPool(size_t); //templateclass F, class... Args //auto enqueue(F f, Args... args) // -std::futuretypename std::result_ofF(Args...)::type; //c14 templateclass F, class... Args auto enqueue(F f, Args... args) -std::futurestd::invoke_result_tF, Args...; // C17 ~ThreadPool(); private: // need to keep track of threads so we can join them std::vector std::thread workers; // the task queue std::queue std::functionvoid() tasks; // synchronization std::mutex queue_mutex; std::condition_variable condition; bool stop; }; // the constructor just launches some amount of workers inline ThreadPool::ThreadPool(size_t threads) : stop(false) { for (size_t i 0; i threads; i) workers.emplace_back( [this] { for (;;) { std::functionvoid() task; { std::unique_lockstd::mutex lock(this-queue_mutex); 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(); } } ); } // add new work item to the pool (C14) //templateclass F, class... Args //auto ThreadPool::enqueue(F f, Args... args) - std::futuretypename std::result_ofF(Args...)::type //{ // using return_type typename std::result_ofF(Args...)::type; // // auto task std::make_shared std::packaged_taskreturn_type() ( // std::bind(std::forwardF(f), std::forwardArgs(args)...) // ); // // std::futurereturn_type res task-get_future(); // { // std::unique_lockstd::mutex lock(queue_mutex); // // // dont allow enqueueing after stopping the pool // if (stop) // throw std::runtime_error(enqueue on stopped ThreadPool); // // tasks.emplace([task]() { (*task)(); }); // } // condition.notify_one(); // return res; //} //C 17标准 templateclass F, class... Args auto ThreadPool::enqueue(F f, Args... args) - std::futurestd::invoke_result_tF, Args... // 正确F, Args... { using return_type std::invoke_result_tF, Args...; // 正确F, Args... auto task std::make_sharedstd::packaged_taskreturn_type()( std::bind(std::forwardF(f), std::forwardArgs(args)...) ); std::futurereturn_type res task-get_future(); { std::unique_lockstd::mutex lock(queue_mutex); if (stop) throw std::runtime_error(enqueue on stopped ThreadPool); tasks.emplace([task]() { (*task)(); }); } condition.notify_one(); return res; } // the destructor joins all threads inline ThreadPool::~ThreadPool() { { std::unique_lockstd::mutex lock(queue_mutex); stop true; } condition.notify_all(); for (std::thread worker : workers) worker.join(); } #endifThreadPool.cpp文件实现// ThreadPool.cpp #include iostream #include vector #include chrono #include ThreadPool.h int main() { ThreadPool pool(4); std::vector std::futureint results; for (int i 0; i 8; i) { results.emplace_back( pool.enqueue([i] { std::cout hello i std::endl; std::this_thread::sleep_for(std::chrono::seconds(1)); std::cout world i std::endl; return i * i; }) ); } for (auto result : results) std::cout result.get() ; std::cout std::endl; return 0; }编译运行参考【C多线程】RAII 锁家族lock_guard / unique_lock / scoped_lock 的设计哲学与源码剖析100行代码使用C实现一个线程池C TheadPool 线程池实现
上一篇/下一篇内容由系统自动关联
返回资讯列表 →