Thread pooling in C++11
c++, c++11, multithreading, stdthread, threadpool
Solution
This is adapted from my answer to another very similar post.
Let's build a `ThreadPool` class:
class ThreadPool {
public:
void Start();
void QueueJob(const std::function<void()>& job);
void Stop();
bool busy();
private:
void ThreadLoop();
bool should_terminate = false; // Tells threads to stop looking for jobs
std::mutex queue_mutex; // Prevents data races to the job queue
std::condition_variable mutex_condition; // Allows threads to wait on new jobs or termination
std::vector<std::thread> threads;
std::queue<std::function<void()>> jobs;
};
- `ThreadPool::Start`
For an efficient threadpool implementation, once threads are created according to `num_threads`, it's better not to create new ones or destroy old ones (by joining). There will be a performance penalty, and it might even make your application go slower than the serial version. Thus, we keep a pool of threads that can be used at any time (if they aren't already running a job).
Each thread should be running its own infinite loop, constantly waiting for new tasks to grab and run.
void ThreadPool::Start() {
const uint32_t num_threads = std::thread::hardware_concurrency(); // Max # of threads the system supports
for (uint32_t ii = 0; ii < num_threads; ++ii) {
threads.emplace_back(std::thread(&ThreadPool::ThreadLoop,this))
}
}
- `ThreadPool::ThreadLoop`
The infinite loop function. This is a `while (true)` loop waiting for the task queue to open up.
void ThreadPool::ThreadLoop() {
while (true) {
std::function<void()> job;
{
std::unique_lock<std::mutex> lock(queue_mutex);
mutex_condition.wait(lock, [this] {
return !jobs.empty() || should_terminate;
});
if (should_terminate) {
return;
}
job = jobs.front();
jobs.pop();
}
job();
}
}
- `ThreadPool::QueueJob`
Add a new job to the pool; use a lock so that there isn't a data race.
void ThreadPool::QueueJob(const std::function<void()>& job) {
{
std::unique_lock<std::mutex> lock(queue_mutex);
jobs.push(job);
}
mutex_condition.notify_one();
}
To use it:
thread_pool->QueueJob([] { /* ... */ });
- `ThreadPool::busy`
bool ThreadPool::busy() {
bool poolbusy;
{
std::unique_lock<std::mutex> lock(queue_mutex);
poolbusy = !jobs.empty();
}
return poolbusy;
}
The busy() function can be used in a while loop, such that the main thread can wait the threadpool to complete all the tasks before calling the threadpool destructor.
- `ThreadPool::Stop`
Stop the pool.
void ThreadPool::Stop() {
{
std::unique_lock<std::mutex> lock(queue_mutex);
should_terminate = true;
}
mutex_condition.notify_all();
for (std::thread& active_thread : threads) {
active_thread.join();
}
threads.clear();
}
Once you integrate these ingredients, you have your own dynamic threading pool. These threads always run, waiting for job to do.
I apologize if there are some syntax errors, I typed this code and and I have a bad memory. Sorry that I cannot provide you the complete thread pool code; that would violate my job integrity.
Notes:
- The anonymous code blocks are used so that when they are exited, the `std::unique_lock` variables created within them go out of scope, unlocking the mutex.
- `ThreadPool::Stop` will not terminate any currently running jobs, it just waits for them to finish via `active_thread.join()`.
Problem
Relevant questions: About C++11: - C++11: std::thread pooled? - Will async(launch::async) in C++11 make thread pools obsolete for avoiding expensive thread creation? About Boost: - C++ boost thread reusing threads - boost::thread and creating a pool of them! How do I get a pool of threads to send tasks to, without creating and deleting them over and over again? This means persistent threads to resynchronize without joining. I have code that looks like this: ``` namespace { std::vector<std::thread> workers; int total = 4; int arr[4] = {0}; void each_thread_does(int i) { arr[i] += 2; } } int main(int argc, char *argv[]) { for (int i = 0; i < 8; ++i) { // for 8 iterations, for (int j = 0; j < 4; ++j) { workers.push_back(std::thread(each_thread_does, j)); } for (std::thread &t: workers) { if (t.joinable()) { t.join(); } } arr[4] = std::min_element(arr, arr+4); } return 0; } ``` Instead of creating and joining threads each iteration, I'd prefer to send tasks to my worker threads each iteration and only create them once.
Related problems
- C++11: std::thread pooled?
- Does async(launch::async) in C++11 make thread pools obsolete for avoiding expensive thread creation?
- C++ Thread Pool
- C++11 Dynamic Threadpool
- Extend the life of threads with synchronization (C++11)
- C++ boost thread reusing threads
- How can I wait on multiple things
- boost::thread and creating a pool of them!