NOTE / 10/29/2019
C++ Concurrent Programming: Part 2
In most systems, assigning every task to a particular thread is impractical. Existing concurrency can instead be used for concurrent execution. A thread pool provides this: submitted tasks are queued, worker threads take tasks from the queue, and each worker takes another task after completing one. — C++ Concurrency in Action
The following is a simple thread-pool implementation:
#include <atomic>
#include <condition_variable>
#include <functional>
#include <iostream>
#include <mutex>
#include <queue>
#include <thread>
#include <vector>
class ThreadPool {
public:
ThreadPool(const int thread_num) {
threads_.reserve(thread_num);
for (size_t i = 0; i < thread_num; ++i) {
threads_.emplace_back(&ThreadPool::Thread, this);
}
}
~ThreadPool() {
// Stop all threads in the pool.
done_.store(true);
thread_cond_var_.notify_all();
for (size_t i = 0; i < threads_.size(); ++i) {
threads_.at(i).join();
}
std::cout << "[Threads stoped]\n";
}
void Submit(const std::function<void()>& work) {
{
// Add work to the queue.
std::lock_guard<std::mutex> lg(mutex_);
works_.push(work);
}
// Tell all threads to work.
thread_cond_var_.notify_all();
}
private:
void Thread() {
done_.store(false);
while (true) {
// Wait for work.
std::function<void()> work;
{
std::unique_lock<std::mutex> ul(mutex_);
thread_cond_var_.wait(ul, [this] { return !works_.empty() || done_.load(); });
if (done_.load()) {
break;
}
work = works_.front();
works_.pop();
}
// Do work.
work();
}
}
std::atomic_bool done_;
std::mutex mutex_;
std::condition_variable thread_cond_var_;
std::queue<std::function<void()>> works_;
std::vector<std::thread> threads_;
};
class Task {
public:
void PrintText(const std::string& text) {
std::this_thread::sleep_for(std::chrono::milliseconds(100));
{
std::lock_guard<std::mutex> lg(mutex_);
std::cout << "Thread " << std::this_thread::get_id() << " do task " << text << std::endl;
}
}
std::mutex mutex_;
};
int main (int argc, char** argv) {
ThreadPool thread_pool(2);
Task task;
for (size_t i = 0; i < 10; ++i) {
thread_pool.Submit(std::bind(&Task::PrintText, &task, std::to_string(i)));
std::cout << "Submit task: " << i << std::endl;
}
std::this_thread::sleep_for(std::chrono::seconds(5));
return 1;
}