C++17 High-Performance Thread Pool
C++17 High-Performance Thread Pool
Computing
Barak Shoshany
arXiv:2105.00613v4 [[Link]] 27 Dec 2023
# bshoshany@[Link] 0000-0003-2222-127X
Department of Physics, Brock University
+ 1812 Sir Isaac Brock Way, St. Catharines, Ontario, L2S 3A1, Canada
[Link] § @bshoshany
Abstract
We present a modern C++17-compatible thread pool implementation, built from scratch with high-
performance scientific computing in mind. The thread pool is implemented as a single lightweight and
self-contained class, and does not have any dependencies other than the C++17 standard library, thus
allowing a great degree of portability. In particular, our implementation does not utilize OpenMP or any
other high-level multithreading APIs, and thus gives the programmer precise low-level control over the
details of the parallelization, which permits more robust optimizations. The thread pool was extensively
tested on both AMD and Intel CPUs with up to 40 cores and 80 threads. This paper provides motivation,
detailed usage instructions, and performance tests.
Contents
1 Introduction 3
1.1 Motivation . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3
1.2 Overview of features . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3
1.3 Compiling and compatibility . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 5
2 Getting started 5
2.1 Installing the library . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 5
2.2 Constructors . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 6
2.3 Getting and resetting the number of threads in the pool . . . . . . . . . . . . . . . . . . . . . 6
2.4 Finding the version of the library . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 6
4 Parallelizing loops 15
4.1 Automatic parallelization of loops . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 15
4.2 Parallelizing loops without futures . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 17
4.3 Parallelizing individual indices vs. blocks . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 18
4.4 Loops with return values . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 19
4.5 Parallelizing sequences . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 20
1
4.6 More about BS::multi_future<T> . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 21
5 Utility classes 21
5.1 Synchronizing printing to a stream with BS::synced_stream . . . . . . . . . . . . . . . . . . 21
5.2 Measuring execution time with BS::timer . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 23
5.3 Sending simple signals between threads with BS::signaller . . . . . . . . . . . . . . . . . . 24
6 Managing tasks 24
6.1 Monitoring the tasks . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 24
6.2 Purging tasks . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 25
6.3 Exception handling . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 26
6.4 Getting information about the threads . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 28
6.5 Thread pool initialization functions . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 28
6.6 Passing task arguments by constant reference . . . . . . . . . . . . . . . . . . . . . . . . . . . 29
7 Optional features 30
7.1 Pausing the pool . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 30
7.2 Avoiding wait deadlocks . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 33
7.3 Accessing native thread handles . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 34
7.4 Setting task priority . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 35
2
1 Introduction
1.1 Motivation
Multithreading [1] is essential for modern high-performance computing. Since C++11, the C++ [2][3][4]
standard library has included built-in low-level multithreading support using constructs such as
std::thread. However, std::thread creates a new thread each time it is called, which can have a
significant performance overhead. Furthermore, it is possible to create more threads than the hardware
can handle simultaneously, potentially resulting in a substantial slowdown.
The library presented here contains a C++ thread pool class, BS::thread_pool, which avoids these issues
by creating a fixed pool of threads once and for all, and then continuously reusing the same threads to
perform different tasks throughout the lifetime of the program. By default, the number of threads in the
pool is equal to the maximum number of threads that the hardware can run in parallel.
The user submits tasks to be executed into a queue. Whenever a thread becomes available, it retrieves
the next task from the queue and executes it. The pool automatically produces an std::future for each
task, which allows the user to wait for the task to finish executing and/or obtain its eventual return value, if
applicable. Threads and tasks are autonomously managed by the pool in the background, without requiring
any input from the user aside from submitting the desired tasks.
The design of this library is guided by four important principles. First, compactness: the entire library con-
sists of just one self-contained header file, with no other components or dependencies, aside from a small
self-contained header file with optional utilities. Second, portability: the library only utilizes the C++17
standard library [5], without relying on any compiler extensions or 3rd-party libraries, and is therefore
compatible with any modern standards-conforming C++17 compiler on any platform. Third, ease of use:
the library is extensively documented, and programmers of any level should be able to use it right out of
the box.
The fourth and final guiding principle is performance: each and every line of code in this library was care-
fully designed with maximum performance in mind, and performance was tested and verified on a vari-
ety of compilers and platforms. Indeed, the library was originally designed for use in the author's own
computationally-intensive scientific computing projects, running both on high-end desktop/laptop com-
puters and high-performance computing nodes.
Other, more advanced multithreading libraries may offer more features and/or higher performance. How-
ever, they typically consist of a vast codebase with multiple components and dependencies, and involve
complex APIs that require a substantial time investment to learn. This library is not intended to replace
these more advanced libraries; instead, it was designed for users who don't require very advanced features,
and prefer a simple and lightweight library that is easy to learn and use and can be readily incorporated
into existing or new projects.
3
– Only 304 lines of code, excluding comments, blank lines, and lines containing only a single
brace, with all optional features disabled.
– Only 396 lines of code across both header files with all optional features enabled and including
all optional utilities.
• Easy to use:
– Very simple operation, using only a handful of member functions, with additional member func-
tions for more advanced use.
– Every task submitted to the queue using the submit_task() member function automatically
generates an std::future, which can be used to wait for the task to finish executing and/or
obtain its eventual return value.
– Loops can be automatically parallelized into any number of tasks using the submit_loop() mem-
ber function, which returns a BS::multi_future that can be used to track the execution of all
parallel tasks at once.
– If futures are not needed, tasks may be submitted using detach_task(), and loops can be par-
allelized using detach_loop() - sacrificing convenience for even greater performance. In that
case, wait(), wait_for(), and wait_until() can be used to wait for all the tasks in the queue
to complete.
– The code is thoroughly documented using Doxygen comments - not only the interface, but also
the implementation, in case the user would like to make modifications.
– The included test program BS_thread_pool_test.cpp can be used to perform exhaustive auto-
mated tests and benchmarks, and also serves as a comprehensive example of how to properly
use the library. The included PowerShell script BS_thread_pool_test.ps1 provides a portable
way to run the tests with multiple compilers.
• Utility classes:
– The optional header file BS_thread_pool_utils.hpp contains several useful utility classes.
– Send simple signals between threads using the BS::signaller utility class.
– Synchronize output to a stream from multiple threads in parallel using the BS::synced_stream
utility class.
– Easily measure execution time for benchmarking purposes using the BS::timer utility class.
• Additional features:
– Assign a priority to each task using the optional task priority feature. Tasks with higher priorities
will be executed first.
– Submit a sequence of tasks enumerated by indices to the queue using detach_sequence() and
submit_sequence().
– Change the number of threads in the pool safely and on-the-fly as needed using the reset()
member function.
– Monitor the number of queued and/or running tasks using the get_tasks_queued(),
get_tasks_running(), and get_tasks_total() member functions.
– Get the current thread count of the pool using get_thread_count().
– Freely pause and resume the pool using the pause(), unpause(), and is_paused() member
functions; when paused, threads do not retrieve new tasks out of the queue.
– Purge all tasks currently waiting in the queue with the purge() member function.
– Catch exceptions thrown by tasks submitted using submit_task() or submit_loop() from the
main thread through their futures.
– Run an initialization function in each thread before it starts to execute any submitted tasks.
– Get the pool index of the current thread using BS::this_thread::get_index() and a pointer
to the pool that owns the thread using BS::this_thread::get_pool().
– Get the unique thread IDs for all threads in the pool using get_thread_ids() or the
implementation-defined thread handles using the optional get_native_handles() member
function.
– Submit class member functions to the pool, either applied to a specific object or from within the
object itself.
– Pass arguments to tasks by value, reference, or constant reference.
– Under continuous and active development. Bug reports and feature requests are welcome, and
4
should be made via GitHub issues.
2 Getting started
2.1 Installing the library
To install BS::thread_pool, simply download the latest release from the GitHub repository, place the
header file BS_thread_pool.hpp from the include folder in the desired folder, and include it in your
program:
#include "BS_thread_pool.hpp"
The thread pool will now be accessible via the BS::thread_pool class. For an even quicker installation,
you can download the header file itself directly at this URL.
5
This library also comes with an independent utilities header file BS_thread_pool_utils.hpp, which is not
required to use the thread pool, but provides some utility classes that may be helpful for multithreading.
This header file also resides in the include folder. It can be downloaded directly at this URL.
This library is also available on various package managers and build system, including vcpkg, Conan,
Meson, and CMake with CPM. Please see below for more details.
2.2 Constructors
The default constructor creates a thread pool with as many threads as the hardware can handle concur-
rently, as reported by the implementation via std::thread::hardware_concurrency(). This is usually
determined by the number of cores in the CPU. If a core is hyperthreaded, it will count as two threads. For
example:
// Constructs a thread pool with as many threads as available in the hardware.
BS::thread_pool pool;
Optionally, a number of threads different from the hardware concurrency can be specified as an argument
to the constructor. However, note that adding more threads than the hardware can handle will not improve
performance, and in fact will most likely hinder it. This option exists in order to allow using less threads
than the hardware concurrency, in cases where you wish to leave some threads available for other processes.
For example:
// Constructs a thread pool with only 12 threads.
BS::thread_pool pool(12);
Usually, when the thread pool is used, a program's main thread should only submit tasks to the thread pool
and wait for them to finish, and should not perform any computationally intensive tasks on its own. In
that case, it is recommended to use the default value for the number of threads. This ensures that all of the
threads available in the hardware will be put to work while the main thread waits.
Sample output:
Thread pool library version is v4.0.0 (2023-12-27).
6
3 Submitting tasks to the queue
3.1 Submitting tasks with no arguments and receiving a future
In this section we will learn how to submit a task with no arguments, but potentially with a return value,
to the queue. Once a task has been submitted, it will be executed as soon as a thread becomes available.
Tasks are executed in the order that they were submitted (first-in, first-out), unless task priority is enabled
(see below).
For example, if the pool has 8 threads and an empty queue, and we submitted 16 tasks, then we should
expect the first 8 tasks to be executed in parallel, with the remaining tasks being picked up by the threads
one by one as each thread finishes executing its first task, until no tasks are left in the queue.
The member function submit_task() is used to submit tasks to the queue. It takes exactly one input, the
task to submit. This task must be a function with no arguments, but it can have a return value. The return
value is an std::future associated to the task.
If the submitted function has a return value of type T, then the future will be of type std::future<T>, and
will be set to the return value when the function finishes its execution. If the submitted function does not
have a return value, then the future will be an std::future<void>, which will not return any value but
may still be used to wait for the function to finish.
Using auto for the return value of submit_task() means the compiler will automatically detect which
instance of the template std::future to use. However, specifying the particular type std::future<T>, as
in the examples below, is recommended for increased readability.
To wait until the task finishes, use the member function wait() of the future. To obtain the return value,
use the member function get(), which will also automatically wait for the task to finish if it hasn't yet.
Here is a simple example:
#include "BS_thread_pool.hpp" // BS::thread_pool
#include <future> // std::future
#include <iostream> // std::cout
int the_answer()
{
return 42;
}
int main()
{
BS::thread_pool pool;
std::future<int> my_future = pool.submit_task(the_answer);
std::cout << my_future.get() << '\n';
}
In this example we submitted the function the_answer(), which returns an int. The member function
submit_task() of the pool therefore returned an std::future<int>. We then used used the get() member
function of the future to get the return value, and printed it out.
In addition to submitted a pre-defined function, we can also use a lambda expression to quickly define the
task on-the-fly. Rewriting the previous example in terms of a lambda expression, we get:
#include "BS_thread_pool.hpp" // BS::thread_pool
#include <future> // std::future
#include <iostream> // std::cout
int main()
7
{
BS::thread_pool pool;
std::future<int> my_future = pool.submit_task([]{ return 42; });
std::cout << my_future.get() << '\n';
}
Here, the lambda expression []{ return 42; } has two parts:
1. An empty capture clause, denoted by []. This signifies to the compiler that a lambda expression is
being defined.
2. A code block { return 42; } that simply returns the value 42.
It is generally simpler and faster to submit lambda expressions rather than pre-defined functions, especially
due to the ability to capture local variables, which we will discuss in the next section.
Of course, tasks do not have to return values. In the following example, we submit a function with no
return value and then using the future to wait for it to finish executing:
#include "BS_thread_pool.hpp" // BS::thread_pool
#include <chrono> // std::chrono
#include <future> // std::future
#include <iostream> // std::cout
#include <thread> // std::this_thread
int main()
{
BS::thread_pool pool;
const std::future<void> my_future = pool.submit_task(
[]
{
std::this_thread::sleep_for(std::chrono::milliseconds(500));
});
std::cout << "Waiting for the task to complete... ";
my_future.wait();
std::cout << "Done." << '\n';
}
Here we split the lambda into multiple lines to make it more readable. The command
std::this_thread::sleep_for(std::chrono::milliseconds(500))
instructs the task to simply sleep for 500 milliseconds, simulating a computationally-intensive task.
8
}
int main()
{
BS::thread_pool pool;
std::future<double> my_future = pool.submit_task(
[]
{
return multiply(6, 7);
});
std::cout << my_future.get() << '\n';
}
As you can see, to pass the arguments to multiply we simply called multiply(6, 7) explicitly inside
a lambda. If the arguments are not literals, we need to use the lambda capture clause to capture the
arguments from the local scope:
#include "BS_thread_pool.hpp" // BS::thread_pool
#include <future> // std::future
#include <iostream> // std::cout
int main()
{
BS::thread_pool pool;
constexpr double first = 6;
constexpr double second = 7;
std::future<double> my_future = pool.submit_task(
[first, second]
{
return multiply(first, second);
});
std::cout << my_future.get() << '\n';
}
We could even get rid of the multiply function entirely and put everything inside a lambda, if desired:
#include "BS_thread_pool.hpp" // BS::thread_pool
#include <future> // std::future
#include <iostream> // std::cout
int main()
{
BS::thread_pool pool;
constexpr double first = 6;
constexpr double second = 7;
std::future<double> my_future = pool.submit_task(
[first, second]
{
return first * second;
});
9
std::cout << my_future.get() << '\n';
}
int main()
{
BS::thread_pool pool;
int result = 0;
pool.detach_task(
[&result]
{
std::this_thread::sleep_for(std::chrono::milliseconds(100));
result = 42;
});
std::cout << result << '\n';
}
This program first defines a local variable named result and initializes it to 0. It then detaches a task in the
form of a lambda expression. Note that the lambda captures result by reference, as indicated by the & in
front of it. This means that the task can modify result, and any such modification will be reflected in the
main thread. The task changes result to 42, but it first sleeps for 100 milliseconds. When the main thread
prints out the value of result, the task has not yet had time to modify its value, since it is still sleeping.
Therefore, the program will print out the initial value 0.
To wait for the task to complete, we must use the wait() member function after detaching it:
#include "BS_thread_pool.hpp" // BS::thread_pool
#include <chrono> // std::chrono
#include <iostream> // std::cout
#include <thread> // std::this_thread
10
int main()
{
BS::thread_pool pool;
int result = 0;
pool.detach_task(
[&result]
{
std::this_thread::sleep_for(std::chrono::milliseconds(100));
result = 42;
});
[Link]();
std::cout << result << '\n';
}
Now the program will print out the value 42, as expected. Note, however, that wait() will wait for all the
tasks in the queue, including any other tasks that were potentially submitted before or after the one we
care about. If we want to wait for just one task, submit_task() would be a better choice.
int main()
{
BS::thread_pool pool;
const std::future<void> my_future = pool.submit_task(
[]
{
std::this_thread::sleep_for(std::chrono::milliseconds(1000));
std::cout << "Task done!\n";
});
while (true)
{
if (my_future.wait_for(std::chrono::milliseconds(200)) !=
11
std::future_status::ready)
std::cout << "Sorry, the task is not done yet.\n";
else
break;
}
}
For detached tasks, since we do not have a future for them, we cannot use this method. However,
BS::thread_pool has two member functions, also named wait_for() and wait_until(), which similarly
wait for a specified duration or until a specified time point, but do so for all tasks (whether submitted or
detached). Instead of an std::future_status, the thread pool's wait functions returns true if all tasks
finished running, or false if the duration expired or the time point was reached but some tasks are still
running.
Here is the same example as above, using detach_task() and pool.wait_for():
#include "BS_thread_pool.hpp" // BS::thread_pool
#include <chrono> // std::chrono
#include <iostream> // std::cout
#include <thread> // std::this_thread
int main()
{
BS::thread_pool pool;
pool.detach_task(
[]
{
std::this_thread::sleep_for(std::chrono::milliseconds(1000));
std::cout << "Task done!\n";
});
while (true)
{
if (!pool.wait_for(std::chrono::milliseconds(200)))
std::cout << "Sorry, the task is not done yet.\n";
else
break;
}
}
class flag_class
{
public:
12
[[nodiscard]] bool get_flag() const
{
return flag;
}
private:
bool flag = false;
};
int main()
{
flag_class flag_object;
flag_object.set_flag(true);
std::cout << std::boolalpha << flag_object.get_flag() << '\n';
}
This program creates a new object flag_object of the class flag_class, sets the flag to true using the
setter member function set_flag(), and then prints out the flag's value using the getter member function
get_flag().
What if we want to submit the member function set_flag() as a task to the thread pool? We simply wrap
the entire statement flag_object.set_flag(true); from line in a lambda, and pass flag_object to the
lambda by reference, as in this example:
#include "BS_thread_pool.hpp" // BS::thread_pool
#include <iostream> // std::cout, std::boolalpha
class flag_class
{
public:
[[nodiscard]] bool get_flag() const
{
return flag;
}
private:
bool flag = false;
};
int main()
{
BS::thread_pool pool;
flag_class flag_object;
pool.submit_task(
[&flag_object]
13
{
flag_object.set_flag(true);
})
.wait();
std::cout << std::boolalpha << flag_object.get_flag() << '\n';
}
Of course, this will also work with detach_task(), if we call wait() on the pool itself instead of on the
returned future.
Note that in this example, instead of getting a future from submit_task() and then waiting for that future,
we simply called wait() on that future straight away. This is a common way of waiting for a task to
complete if we have nothing else to do in the meantime. Note also that we passed flag_object by reference
to the lambda, since we want to set the flag on that same object, not a copy of it (passing by value wouldn't
have worked anyway, since variables captured by value are implicitly const).
Another thing you might want to do is call a member function from within the object itself, that is, from
another member function. This follows a similar syntax, except that you must also capture this (i.e. a
pointer to the current object) in the lambda. Here is an example:
#include "BS_thread_pool.hpp" // BS::thread_pool
#include <iostream> // std::cout, std::boolalpha
BS::thread_pool pool;
class flag_class
{
public:
[[nodiscard]] bool get_flag() const
{
return flag;
}
void set_flag_to_true()
{
pool.submit_task(
[this]
{
set_flag(true);
})
.wait();
}
private:
bool flag = false;
};
int main()
{
flag_class flag_object;
14
flag_object.set_flag_to_true();
std::cout << std::boolalpha << flag_object.get_flag() << '\n';
}
Note that in this example we defined the thread pool as a global object, so that it is accessible outside the
main() function.
4 Parallelizing loops
4.1 Automatic parallelization of loops
One of the most common and effective methods of parallelization is splitting a loop into smaller loops and
running them in parallel. It is most effective in "embarrassingly parallel" computations, such as vector or
matrix operations, where each iteration of the loop is completely independent of every other iteration.
For example, if we are summing up two vectors of 1000 elements each, and we have 10 threads, we could
split the summation into 10 blocks of 100 elements each, and run all the blocks in parallel, potentially
increasing performance by up to a factor of 10.
BS::thread_pool can automatically parallelize loops. To see how this works, consider the following
generic loop:
for (T i = start; i < end; ++i)
loop(i);
where:
• T is any signed or unsigned integer type.
• The loop is over the range [start, end), i.e. inclusive of start but exclusive of end.
• loop() is an operation performed for each loop index i, such as modifying an array with end - start
elements.
This loop may be automatically parallelized and submitted to the thread pool's queue using the member
function submit_loop(), which has the follows syntax:
pool.submit_loop(start, end, loop, num_blocks);
where:
• start is the first index in the range.
• end is the index after the last index in the range, such that the full range is [start, end). In other
words, the loop will be equivalent to the one above if start and end are the same.
– start and end must both be of the same integer type T. See below for examples of what to do
when they are not of the same type.
– Note that if end <= start, nothing will happen.
• loop() is the function that should run in every iteration of the loop, and takes one argument, the
loop index.
• num_blocks is the number of blocks of the form [a, b) to split the loop into. For example, if the
range is [0, 9) and there are 3 blocks, then the blocks will be the ranges [0, 3), [3, 6), and [6,
9).
– The internal algorithm ensures that each of the blocks has one of two sizes, differing by 1, with
the larger blocks always first, so that the tasks are as evenly distributed as possible. For example,
if the range [0, 100) is split into 15 blocks, the result will be 10 blocks of size 7, which will be
executed first, and 5 blocks of size 6.
– This argument can be omitted, in which case the number of blocks will be the number of threads
in the pool.
15
Each block will be submitted to the thread pool's queue as a separate task. Therefore, a loop that is split
into 3 blocks will be split into 3 individual tasks, which may run in parallel. If there is only one block, then
the entire loop will run as one task, and no parallelization will take place.
To parallelize the generic loop above, we use the following commands:
BS::multi_future<void> loop_future = pool.submit_loop(start, end, loop, num_blocks);
loop_future.wait();
submit_loop() returns an object of the helper class template BS::multi_future. This is essentially a spe-
cialization of std::vector<std::future<T>> with additional member functions. Each of the num_blocks
blocks will have an std::future assigned to it, and all these futures will be stored inside the returned
BS::multi_future. When loop_future.wait() is called, the main thread will wait until all tasks gener-
ated by submit_loop() finish executing, and only those tasks - not any other tasks that also happen to be
in the queue. This is essentially the role of the BS::multi_future class: to wait for a specific group of
tasks, in this case the tasks running the loop blocks.
What value should you use for num_blocks? Omitting this argument, so that the number of blocks will
be equal to the number of threads in the pool, is typically a good choice. For best performance, it is
recommended to do your own benchmarks to find the optimal number of blocks for each loop (you can
use the BS::timer utility class). Using fewer tasks than there are threads may be preferred if you are also
running other tasks in parallel. Using more tasks than there are threads may improve performance in some
cases, but parallelization with too many tasks will suffer from diminishing returns.
As a simple example, the following code calculates and prints a table of squares of all integers from 0 to 99:
#include <iomanip> // std::setw
#include <iostream> // std::cout
int main()
{
constexpr unsigned int max = 100;
unsigned int squares[max];
for (unsigned int i = 0; i < max; ++i)
squares[i] = i * i;
for (unsigned int i = 0; i < max; ++i)
std::cout << std::setw(2) << i << "ˆ2 = " << std::setw(4) << squares[i]
<< ((i % 5 != 4) ? " | " : "\n");
}
int main()
{
BS::thread_pool pool(10);
constexpr unsigned int max = 100;
unsigned int squares[max];
const BS::multi_future<void> loop_future = pool.submit_loop<unsigned int>(0, max,
[&squares](const unsigned int i)
{
squares[i] = i * i;
});
loop_future.wait();
16
for (unsigned int i = 0; i < max; ++i)
std::cout << std::setw(2) << i << "ˆ2 = " << std::setw(4) << squares[i]
<< ((i % 5 != 4) ? " | " : "\n");
}
Since there are 10 threads, and we omitted the num_blocks argument, the loop will be divided into 10
blocks, each calculating 10 squares.
Note that submit_loop() was executed with the explicit template parameter <unsigned int>. The reason
is that the two loop indices must be of the same type. However, here max is a unsigned int, while 0 is
a (signed) int, so the types do not match, and the code will not compile unless we force the 0 to be of
the right type. This can be done most elegantly by specifying the type of the indices explicitly using the
template parameter.
The reason this is not done automatically (e.g. using std::common_type is that it may result in accidentally
casting negative indices to an unsigned type, or integer indices to a too narrow integer type, which may
lead to an incorrect loop range.
We could also cast the 0 explicitly to unsigned int, but that doesn't look as nice:
pool.submit_loop(static_cast<unsigned int>(0), max, /* ... */);
As a side note, notice that here we parallelized the calculation of the squares, but we did not parallelize
printing the results. This is for two reasons:
1. We want to print out the squares in ascending order, and we have no guarantee that the blocks will
be executed in the correct order. This is very important; you must never expect that the parallelized
loop will execute at the same order as the non-parallelized loop.
2. If we did print out the squares from within the parallel tasks, we would get a huge mess, since all 10
blocks would print to the standard output at once. Later we will see how to synchronize printing to a
stream from multiple tasks at the same time.
int main()
{
BS::thread_pool pool(10);
constexpr unsigned int max = 100;
unsigned int squares[max];
pool.detach_loop<unsigned int>(0, max,
17
[&squares](const unsigned int i)
{
squares[i] = i * i;
});
[Link]();
for (unsigned int i = 0; i < max; ++i)
std::cout << std::setw(2) << i << "ˆ2 = " << std::setw(4) << squares[i]
<< ((i % 5 != 4) ? " | " : "\n");
}
Warning: Since detach_loop() does not return a BS::multi_future, there is no built-in way for the user
to know when the loop finishes executing. You must use either wait() as we did here, or some other
method such as condition variables, to ensure that the loop finishes executing before trying to use anything
that depends on its output. Otherwise, bad things will happen!
The start and end indices of each block are determined automatically by the pool. For example, in the
previous section, the loop from 0 to 100 was split into 10 blocks of 10 indices each: start = 0 to end =
10, start = 10 to end = 20, and so on; the blocks are not inclusive of the last index, since the for loop has
the condition i < end and not i <= end.
However, this also means that the loop() function is executed multiple times per block. This generates
additional overhead due to the multiple function calls. For short loops, this should not affect performance.
However, for very long loops, with millions of indices, the performance cost may be significate.
For this reason, the thread pool library provides two additional member functions for parallelizing loops:
detach_blocks() and submit_blocks(). While detach_loop() and submit_loop() execute a function
loop(i) once per index but multiple times per block, detach_blocks() and submit_blocks() execute a
function block(start, end) once per block.
The main advantage of this method is increased performance, but the main disadvantage is slightly more
complicated code. In particular, the user must define the loop from start to end manually within each
block. Here is the previous example using detach_blocks():
#include "BS_thread_pool.hpp" // BS::thread_pool
#include <iomanip> // std::setw
#include <iostream> // std::cout
int main()
{
BS::thread_pool pool(10);
constexpr unsigned int max = 100;
unsigned int squares[max];
pool.detach_blocks<unsigned int>(0, max,
[&squares](const unsigned int start, const unsigned int end)
{
for (unsigned int i = start; i < end; ++i)
squares[i] = i * i;
18
});
[Link]();
for (unsigned int i = 0; i < max; ++i)
std::cout << std::setw(2) << i << "ˆ2 = " << std::setw(4) << squares[i]
<< ((i % 5 != 4) ? " | " : "\n");
}
Note how the block function takes two arguments, and includes the internal loop.
Generally, compiler optimizations should be able to make detach_loop() and submit_loop() perform
roughly the same as detach_blocks() and submit_blocks(). However, you should perform your own
benchmarks to see which option works best for your particular use case.
BS::thread_pool pool;
int main()
{
19
std::cout << sum<std::uint64_t>(1, 1'000'000);
}
int main()
20
{
BS::thread_pool pool;
constexpr ui64 max = 20;
BS::multi_future<ui64> sequence_future =
pool.submit_sequence<ui64>(0, max + 1, factorial);
std::vector<ui64> factorials = sequence_future.get();
for (ui64 i = 0; i < max + 1; ++i)
std::cout << i << "! = " << factorials[i] << '\n';
}
5 Utility classes
The optional header file BS_thread_pool_utils.hpp contains several useful utility classes. These are not
necessary for using the thread pool itself; BS_thread_pool.hpp is the only header file required. However,
the utility classes can make writing multithreading code more convenient. The version of the utilities
header file can be found by checking the macro BS_THREAD_POOL_UTILS_VERSION.
BS::thread_pool pool;
int main()
21
{
pool.detach_sequence(0, 5,
[](int i)
{
std::cout << "Task no. " << i << " executing.\n";
});
}
The reason is that, although each individual insertion to std::cout is thread-safe, there is no mechanism
in place to ensure subsequent insertions from the same thread are printed contiguously.
The utility class BS::synced_stream is designed to eliminate such synchronization issues. The construc-
tor takes one optional argument, specifying the output stream to print to. If no argument is supplied,
std::cout will be used:
// Construct a synced stream that will print to std::cout.
BS::synced_stream sync_out;
// Construct a synced stream that will print to the output stream my_stream.
BS::synced_stream sync_out(my_stream);
The member function print() takes an arbitrary number of arguments, which are inserted into the stream
one by one, in the order they were given. println() does the same, but also prints a newline character \n
at the end, for convenience. A mutex is used to synchronize this process, so that any other calls to print()
or println() using the same BS::synced_stream object must wait until the previous call has finished.
As an example, this code:
#include "BS_thread_pool.hpp" // BS::thread_pool
#include "BS_thread_pool_utils.hpp" // BS::synced_stream
BS::synced_stream sync_out;
BS::thread_pool pool;
int main()
{
pool.detach_sequence(0, 5,
[](int i)
{
sync_out.println("Task no. ", i, " executing.");
});
}
22
Warning: Always create the BS::synced_stream object before the BS::thread_pool object, as we did in
this example. When the BS::thread_pool object goes out of scope, it waits for the remaining tasks to
be executed. If the BS::synced_stream object goes out of scope before the BS::thread_pool object, then
any tasks using the BS::synced_stream will crash. Since objects are destructed in the opposite order of
construction, creating the BS::synced_stream object before the BS::thread_pool object ensures that the
BS::synced_stream is always available to the tasks, even while the pool is destructing.
Most stream manipulators defined in the headers <ios> and <iomanip>, such as std::setw (set the char-
acter width of the next output), std::setprecision (set the precision of floating point numbers), and
std::fixed (display floating point numbers with a fixed number of digits), can be passed to print() and
println() just as you would pass them to a stream.
The only exceptions are the flushing manipulators std::endl and std::flush, which will not work be-
cause the compiler will not be able to figure out which template specializations to use. Instead, use
BS::synced_stream::endl and BS::synced_stream::flush. Here is an example:
#include "BS_thread_pool.hpp" // BS::thread_pool
#include "BS_thread_pool_utils.hpp" // BS::synced_stream
#include <cmath> // std::sqrt
#include <iomanip> // std::setprecision, std::setw
#include <ios> // std::fixed
BS::synced_stream sync_out;
BS::thread_pool pool;
int main()
{
sync_out.print(std::setprecision(10), std::fixed);
pool.detach_sequence(0, 16,
[](int i)
{
sync_out.print("The square root of ", std::setw(2), i, " is ",
std::sqrt(i), ".", BS::synced_stream::endl);
});
}
Note, however, that BS::synced_stream::endl should only be used if flushing is desired; otherwise, a
newline character should be used instead.
23
For example:
BS::timer tmr;
[Link]();
do_something();
[Link]();
std::cout << "The elapsed time was " << [Link]() << " ms.\n";
A practical application of the BS::timer class can be found in the benchmark portion of the test program
BS_thread_pool_test.cpp.
6 Managing tasks
6.1 Monitoring the tasks
Sometimes you may wish to monitor what is happening with the tasks you submitted to the pool. This may
be done using three member functions:
• get_tasks_queued() gets the number of tasks currently waiting in the queue to be executed by the
threads.
• get_tasks_running() gets the number of tasks currently being executed by the threads.
• get_tasks_total() gets the total number of unfinished tasks: either still in the queue, or running in
a thread.
• Note that get_tasks_total() == get_tasks_queued() + get_tasks_running().
These functions are demonstrated in the following program:
#include "BS_thread_pool.hpp" // BS::thread_pool
#include "BS_thread_pool_utils.hpp" // BS::synced_stream
#include <chrono> // std::chrono
#include <thread> // std::this_thread
BS::synced_stream sync_out;
BS::thread_pool pool(4);
void monitor_tasks()
{
sync_out.println(pool.get_tasks_total(), " tasks total, ",
pool.get_tasks_running(), " tasks running, ",
24
pool.get_tasks_queued(), " tasks queued.");
}
int main()
{
[Link]();
pool.detach_sequence(0, 12, sleep_half_second);
monitor_tasks();
std::this_thread::sleep_for(std::chrono::milliseconds(750));
monitor_tasks();
std::this_thread::sleep_for(std::chrono::milliseconds(500));
monitor_tasks();
std::this_thread::sleep_for(std::chrono::milliseconds(500));
monitor_tasks();
}
Assuming you have at least 4 hardware threads (so that 4 tasks can run concurrently), the output should
be similar to:
12 tasks total, 0 tasks running, 12 tasks queued.
Task 0 done.
Task 1 done.
Task 2 done.
Task 3 done.
8 tasks total, 4 tasks running, 4 tasks queued.
Task 4 done.
Task 5 done.
Task 6 done.
Task 7 done.
4 tasks total, 4 tasks running, 0 tasks queued.
Task 8 done.
Task 9 done.
Task 10 done.
Task 11 done.
0 tasks total, 0 tasks running, 0 tasks queued.
The reason we called [Link]() in the beginning is that when the thread pool is created, an initialization
task runs in each thread, so if we don't wait, the first line will say there are 16 tasks in total, including the
4 initialization tasks. See below for more details.
25
#include "BS_thread_pool.hpp" // BS::thread_pool
#include "BS_thread_pool_utils.hpp" // BS::synced_stream
#include <chrono> // std::chrono
#include <thread> // std::this_thread
BS::synced_stream sync_out;
BS::thread_pool pool(4);
int main()
{
for (size_t i = 0; i < 8; ++i)
{
pool.detach_task(
[i]
{
std::this_thread::sleep_for(std::chrono::milliseconds(100));
sync_out.println("Task ", i, " done.");
});
}
std::this_thread::sleep_for(std::chrono::milliseconds(50));
[Link]();
[Link]();
}
The program submit 8 tasks to the queue. Each task waits 100 milliseconds and then prints a message. The
thread pool has 4 threads, so it will execute the first 4 tasks in parallel, and then the remaining 4. We wait
50 milliseconds, to ensure that the first 4 tasks have all started running. Then we call purge() to purge
the remaining 4 tasks. As a result, these tasks never get executed. However, since the first 4 tasks are still
running when purge() is called, they will finish uninterrupted; purge() only discards tasks that have not
yet started running. The output of the program therefore only contains the messages from the first 4 tasks:
Task 0 done.
Task 1 done.
Task 2 done.
Task 3 done.
BS::synced_stream sync_out;
BS::thread_pool pool;
int main()
26
{
constexpr double num = 0;
std::future<double> my_future = pool.submit_task(inverse, num);
try
{
const double result = my_future.get();
sync_out.println("The inverse of ", num, " is ", result, ".");
}
catch (const std::exception& e)
{
sync_out.println("Caught exception: ", [Link]());
}
}
However, if you change num to any non-zero number, no exceptions will be thrown and the inverse will be
printed.
It is important to note that wait() does not throw any exceptions; only get() does. Therefore, even if your
task does not return anything, i.e. your future is an std::future<void>, you must still use get() on the
future obtained from it if you want to catch exceptions thrown by it. Here is an example:
#include "BS_thread_pool.hpp"
BS::synced_stream sync_out;
BS::thread_pool pool;
int main()
{
constexpr double num = 0;
std::future<void> my_future = pool.submit_task(print_inverse, num);
try
{
my_future.get();
}
catch (const std::exception& e)
{
sync_out.println("Caught exception: ", [Link]());
}
}
When using BS::multi_future to handle multiple futures at once, exception handling works the same
way: if any of the futures may throw exceptions, you may catch these exceptions when calling get(), even
in the case of BS::multi_future<void>.
27
6.4 Getting information about the threads
BS::thread_pool comes with a variety of methods to obtain information about the threads in the pool:
1. The namespace BS::this_thread provides functionality similar to std::this_thread. If the current
thread belongs to a BS::thread_pool object, then BS::this_thread::get_index() can be used to
get the index of the current thread, and BS::this_thread::get_pool() can be used to get the pointer
to the thread pool that owns the current thread. Please see the reference below for more details.
2. The member function get_thread_ids() returns a vector containing the unique identifiers for each
of the pool's threads, as obtained by std::thread::get_id(). These values are not so useful on their
own, but can be used for whatever the user wants to use them for.
3. The optional member function get_native_handles(), if enabled, returns a vector containing the
underlying implementation-defined thread handles for each of the pool's threads, as obtained by
std::thread::native_handle(). For more information, see the relevant section below.
BS::synced_stream sync_out;
thread_local std::mt19937_64 twister;
int main()
{
BS::thread_pool pool(
[]
{
[Link](std::random_device()());
});
pool.submit_sequence(0, 4,
[](int)
{
sync_out.println("I generated a random number: ", twister());
})
.wait();
}
28
In this example, we create a thread_local Mersenne twister engine, meaning that each thread has
its own independent engine. However, we did not seed the engine, so each thread will generate the
exact same sequence of pseudo-random numbers. To remedy this, we pass an initialization func-
tion to the BS::thread_pool constructor which seeds the twister in each thread with the (hopefully)
non-deterministic random number generator std::random_device.
BS::synced_stream sync_out;
void increment(int& x)
{
++x;
}
int main()
{
BS::thread_pool pool;
int n = 0;
pool.submit_task(
[&n]
{
increment(n);
})
.wait();
pool.submit_task(
[&n = std::as_const(n)]
{
print(n);
})
.wait();
}
The increment() function takes a reference to an integer, and increments that integer. Passing the argu-
ment by reference guarantees that n itself, in the scope of main(), will be incremented - rather than a copy
of it in the scope of increment().
Similarly, the print() function takes a constant reference to an integer, and prints that integer. Passing
the argument by constant reference guarantees that the variable will not be accidentally modified by the
function, even though we are accessing n itself, rather than a copy. If we replace print with increment, the
29
program won't compile, as increment cannot take constant references.
Generally, it is not really necessary to pass arguments by constant reference, but it is more "correct" to do
so, if we would like to guarantee that the variable being referenced is indeed never modified. This section
is therefore included here for completeness.
7 Optional features
7.1 Pausing the pool
Sometimes you may wish to temporarily pause the execution of tasks, or perhaps you want to submit tasks
to the queue in advance and only start executing them at a later time. You can do this using the member
functions pause(), unpause(), and is_paused().
However, these functions are disabled by default, and must be explicitly enabled by defining the macro
BS_THREAD_POOL_ENABLE_PAUSE before including BS_thread_pool.hpp. The reason is that pausing the
pool adds additional checks to the waiting and worker functions, which have a very small but non-zero
overhead.
When you call pause(), the workers will temporarily stop retrieving new tasks out of the queue. However,
any tasks already executed will keep running until they are done, since the thread pool has no control over
the internal code of your tasks. If you need to pause a task in the middle of its execution, you must do that
manually by programming your own pause mechanism into the task itself. To resume retrieving tasks, call
unpause(). To check whether the pool is currently paused, call is_paused().
Here is an example:
#define BS_THREAD_POOL_ENABLE_PAUSE
#include "BS_thread_pool.hpp" // BS::thread_pool
#include "BS_thread_pool_utils.hpp" // BS::synced_stream
#include <chrono> // std::chrono
#include <thread> // std::this_thread
BS::synced_stream sync_out;
BS::thread_pool pool(4);
void check_if_paused()
{
if (pool.is_paused())
sync_out.println("Pool paused.");
else
sync_out.println("Pool unpaused.");
}
int main()
{
pool.detach_sequence(0, 8, sleep_half_second);
sync_out.println("Submitted 8 tasks.");
std::this_thread::sleep_for(std::chrono::milliseconds(250));
[Link]();
30
check_if_paused();
std::this_thread::sleep_for(std::chrono::milliseconds(1000));
sync_out.println("Still paused...");
std::this_thread::sleep_for(std::chrono::milliseconds(1000));
pool.detach_sequence(8, 12, sleep_half_second);
sync_out.println("Submitted 4 more tasks.");
sync_out.println("Still paused...");
std::this_thread::sleep_for(std::chrono::milliseconds(1000));
[Link]();
check_if_paused();
}
Assuming you have at least 4 hardware threads, the output should be similar to:
Submitted 8 tasks.
Pool paused.
Task 0 done.
Task 1 done.
Task 2 done.
Task 3 done.
Still paused...
Submitted 4 more tasks.
Still paused...
Pool unpaused.
Task 4 done.
Task 5 done.
Task 6 done.
Task 7 done.
Task 8 done.
Task 9 done.
Task 10 done.
Task 11 done.
Here is what happened. We initially submitted a total of 8 tasks to the queue. Since we waited for 250ms
before pausing, the first 4 tasks have already started running, so they kept running until they finished.
While the pool was paused, we submitted 4 more tasks to the queue, but they just waited at the end of the
queue. When we unpaused, the remaining 4 initial tasks were executed, followed by the 4 new tasks.
While the workers are paused, wait() will wait for the running tasks instead of all tasks (otherwise it
would wait forever). This is demonstrated by the following program:
#define BS_THREAD_POOL_ENABLE_PAUSE
#include "BS_thread_pool.hpp" // BS::thread_pool
#include "BS_thread_pool_utils.hpp" // BS::synced_stream
#include <chrono> // std::chrono
#include <thread> // std::this_thread
BS::synced_stream sync_out;
BS::thread_pool pool(4);
31
void check_if_paused()
{
if (pool.is_paused())
sync_out.println("Pool paused.");
else
sync_out.println("Pool unpaused.");
}
int main()
{
pool.detach_sequence(0, 8, sleep_half_second);
sync_out.println("Submitted 8 tasks. Waiting for them to complete.");
[Link]();
pool.detach_sequence(8, 20, sleep_half_second);
sync_out.println("Submitted 12 more tasks.");
std::this_thread::sleep_for(std::chrono::milliseconds(250));
[Link]();
check_if_paused();
sync_out.println("Waiting for the ", pool.get_tasks_running(),
" running tasks to complete.");
[Link]();
sync_out.println("All running tasks completed. ", pool.get_tasks_queued(),
" tasks still queued.");
std::this_thread::sleep_for(std::chrono::milliseconds(1000));
sync_out.println("Still paused...");
std::this_thread::sleep_for(std::chrono::milliseconds(1000));
sync_out.println("Still paused...");
std::this_thread::sleep_for(std::chrono::milliseconds(1000));
[Link]();
check_if_paused();
std::this_thread::sleep_for(std::chrono::milliseconds(250));
sync_out.println("Waiting for the remaining ", pool.get_tasks_total(),
" tasks (", pool.get_tasks_running(), " running and ",
pool.get_tasks_queued(), " queued) to complete.");
[Link]();
sync_out.println("All tasks completed.");
}
32
Task 10 done.
Task 11 done.
All running tasks completed. 8 tasks still queued.
Still paused...
Still paused...
Pool unpaused.
Waiting for the remaining 8 tasks (4 running and 4 queued) to complete.
Task 12 done.
Task 13 done.
Task 14 done.
Task 15 done.
Task 16 done.
Task 17 done.
Task 18 done.
Task 19 done.
All tasks completed.
The first wait(), which was called while the pool was not paused, waited for all 8 tasks, both running and
queued. The second wait(), which was called after pausing the pool, only waited for the 4 running tasks,
while the other 8 tasks remained queued, and were not executed since the pool was paused. Finally, the
third wait(), which was called after unpausing the pool, waited for the remaining 8 tasks, both running
and queued.
Warning: If the thread pool is destroyed while paused, any tasks still in the queue will never be executed!
int main()
{
BS::thread_pool pool;
pool.detach_task(
[&pool]
{
[Link]();
std::cout << "Done waiting.\n";
});
}
This program creates a thread pool, and then detaches a task that waits for tasks in the same thread pool
to complete. If you run this program, it will never print the message "Done waiting", because the task will
wait for itself to complete. This causes a deadlock, and the program will wait forever.
Usually, in simple programs, this will never happen. However, in more complicated programs, per-
haps ones running multiple thread pools in parallel, wait deadlocks could potentially occur. In such
cases, the macro BS_THREAD_POOL_ENABLE_WAIT_DEADLOCK_CHECK can be defined before including
BS_thread_pool.hpp. wait() will then check whether the user tried to call it from within a thread of the
same pool, and if so, it will throw the exception BS::thread_pool::wait_deadlock instead of waiting.
This check is disabled by default because wait deadlocks are not something that happens often, and the
check adds a small but non-zero overhead every time wait() is called.
Here is an example:
33
#define BS_THREAD_POOL_ENABLE_WAIT_DEADLOCK_CHECK
#include "BS_thread_pool.hpp" // BS::thread_pool
#include <iostream> // std::cout
int main()
{
BS::thread_pool pool;
pool.detach_task(
[&pool]
{
try
{
[Link]();
std::cout << "Done waiting.\n";
}
catch (const BS::thread_pool::wait_deadlock&)
{
std::cout << "Error: Deadlock!\n";
}
});
}
This time, wait() will detect the deadlock, and will throw an exception, causing the output to be "Error:
Deadlock!".
BS::synced_stream sync_out;
BS::thread_pool pool(4);
int main()
{
std::vector<std::thread::native_handle_type> handles = pool.get_native_handles();
for (BS::concurrency_t i = 0; i < [Link](); ++i)
sync_out.println("Thread ", i, " native handle: ", handles[i]);
}
The output will depend on your compiler and operating system. Here is an example:
34
Thread 0 native handle: 000000F4
Thread 1 native handle: 000000F8
Thread 2 native handle: 000000EC
Thread 3 native handle: 000000FC
BS::synced_stream sync_out;
BS::thread_pool pool(1);
int main()
{
pool.detach_task([] { sync_out.println("This task will execute third."); },
BS::pr::normal);
pool.detach_task([] { sync_out.println("This task will execute fifth."); },
BS::pr::lowest);
pool.detach_task([] { sync_out.println("This task will execute second."); },
BS::pr::high);
pool.detach_task([] { sync_out.println("This task will execute first."); },
BS::pr::highest);
pool.detach_task([] { sync_out.println("This task will execute fourth."); },
BS::pr::low);
}
This program will print out the tasks in the correct priority order. Note that for simplicity, we used a pool
with just one thread, so the tasks will run one at a time. In a pool with 5 or more threads, all 5 tasks will
actually run more or less at the same time, because, for example, the task with the second-highest priority
will be picked up by another thread while the task with the highest priority is still running.
Of course, this is just a pedagogical example. In a realistic use case we may want, for example, to submit
tasks that must be completed immediately with high priority so they skip over other tasks already in the
queue, or background non-urgent tasks with low priority so they evaluate only after higher-priority tasks
are done.
Here are some subtleties to note when using task priority:
35
• Task priority is facilitated using std::priority_queue, which has O(log n) complexity for storing
new tasks, but only O(1) complexity for retrieving the next (i.e. highest-priority) task. This is in
contrast with std::queue, used if priority is disabled, which both stores and retrieves with O(1)
complexity.
• Due to this, enabling the priority queue can incur a very slight decrease in performance, depending on
the specific use case, which is why this feature is disabled by default. As usual, there is a trade-off here,
where you get functionality in exchange for performance. However, the difference in performance is
never substantial, and compiler optimizations can often reduce it to a negligible amount.
• When using the priority queue, tasks will not necessarily be executed in the same order they
were submitted, even if they all have the same priority. This is due to the implementation of
std::priority_queue as a binary heap, which means tasks are stored as a binary tree instead of
sequentially. To execute tasks in submission order, give them monotonically decreasing priorities.
• Technically, BS::priority_t is defined to be (std::int_least16_t), since this type is guaranteed
to be present on all systems, rather than std::int16_t, which is optional in the C++ standard.
This means that on some exotic systems BS::priority_t may actually have more than 16 bits.
However, the pre-defined priorities are 100% portable, and will always have the same values (e.g.:
BS::pr::highest = 32767) regardless of the actual bit width.
36
8.2 Performance tests
If all checks passed, BS_thread_pool_test.cpp performs simple benchmarks by filling a very large vector
with values using detach_blocks(). The program decides what the size of the vector should be by testing
how many elements are needed to reach a certain target duration when parallelizing using a number of
blocks equal to the number of threads. This ensures that the test takes approximately the same amount of
time on all systems, and is thus more consistent and portable.
Once the appropriate size of the vector has been determined, the program allocates the vector and fills it
with values, calculated according to a fixed prescription. This operation is performed both single-threaded
and multithreaded, with the multithreaded computation spread across multiple tasks submitted to the
pool.
Several different multithreaded tests are performed, with the number of tasks either equal to, smaller than,
or larger than the pool's thread count. Each test is repeated multiple times, with the run times averaged
over all runs of the same test. The program keeps increasing the number of blocks by a factor of 2 until
diminishing returns are encountered. The run times of the tests are compared, and the maximum speedup
obtained is calculated.
As an example, here are the results of the benchmarks from a Digital Research Alliance of Canada node
equipped with two 20-core / 40-thread Intel Xeon Gold 6148 CPUs (for a total of 40 cores and 80
threads), running CentOS Linux 7.9.2009. The tests were compiled using GCC v13.2.0 with the -O3 and
-march=native flags. The output was as follows:
======================
Performing benchmarks:
======================
Using 80 threads.
Determining the number of elements to generate in order to achieve an approximate mean
execution time of 50 ms with 80 tasks...
Each test will be repeated up to 30 times to collect reliable statistics.
Generating 27962000 elements:
[......]
Single-threaded, mean execution time was 2815.2 ms with standard deviation 3.5 ms.
[......]
With 2 tasks, mean execution time was 1431.3 ms with standard deviation 10.1 ms.
[.......]
With 4 tasks, mean execution time was 722.1 ms with standard deviation 11.4 ms.
[..............]
With 8 tasks, mean execution time was 364.9 ms with standard deviation 10.9 ms.
[............................]
With 16 tasks, mean execution time was 181.9 ms with standard deviation 8.0 ms.
[..............................]
With 32 tasks, mean execution time was 110.6 ms with standard deviation 1.8 ms.
[..............................]
With 64 tasks, mean execution time was 64.0 ms with standard deviation 6.3 ms.
[..............................]
With 128 tasks, mean execution time was 59.8 ms with standard deviation 0.8 ms.
[..............................]
With 256 tasks, mean execution time was 59.0 ms with standard deviation 0.0 ms.
[..............................]
With 512 tasks, mean execution time was 52.8 ms with standard deviation 0.4 ms.
[..............................]
With 1024 tasks, mean execution time was 50.7 ms with standard deviation 0.9 ms.
[..............................]
With 2048 tasks, mean execution time was 50.0 ms with standard deviation 0.5 ms.
37
[..............................]
With 4096 tasks, mean execution time was 49.4 ms with standard deviation 0.5 ms.
[..............................]
With 8192 tasks, mean execution time was 50.2 ms with standard deviation 0.4 ms.
Maximum speedup obtained by multithreading vs. single-threading: 56.9x, using 4096 tasks.
+++++++++++++++++++++++++++++++++++++++
Thread pool performance test completed!
+++++++++++++++++++++++++++++++++++++++
These two CPUs have 40 physical cores in total, with each core providing two separate logical cores via
hyperthreading, for a total of 80 threads. Without hyperthreading, we would expect a maximum theoretical
speedup of 40x. With hyperthreading, one might naively expect to achieve up to an 80x speedup, but this
is in fact impossible, as each pair of hyperthreaded logical cores share the same physical core's resources.
However, generally we would expect at most an estimated 30% additional speedup [6] from hyperthreading,
which amounts to around 52x in this case. The speedup of 56.9x in our performance test exceeds this
estimate.
On Windows:
.\vcpkg install bshoshany-thread-pool:x86-windows bshoshany-thread-pool:x64-windows
To update the package to the latest version, simply change the version number. Please refer to
this package's page on ConanCenter for more information.
38
meson wrap update bshoshany-thread-pool
This will automatically download the indicated version of the package from this GitHub repository and
include it in your project.
It is also possible to use CPM without installing it first, by adding the following lines to [Link]
before CPMAddPackage:
set(CPM_DOWNLOAD_LOCATION "${CMAKE_BINARY_DIR}/cmake/[Link]")
if(NOT(EXISTS ${CPM_DOWNLOAD_LOCATION}))
message(STATUS "Downloading [Link]")
file(DOWNLOAD
[Link]
${CPM_DOWNLOAD_LOCATION})
endif()
include(${CPM_DOWNLOAD_LOCATION})
Here is an example of a complete [Link] for a project named my_project consisting of a single
source file [Link] which uses BS_thread_pool.hpp:
cmake_minimum_required(VERSION 3.19)
project(my_project LANGUAGES CXX)
set(CMAKE_CXX_STANDARD 17)
set(CMAKE_CXX_STANDARD_REQUIRED ON)
set(CMAKE_CXX_EXTENSIONS OFF)
set(CPM_DOWNLOAD_LOCATION "${CMAKE_BINARY_DIR}/cmake/[Link]")
if(NOT(EXISTS ${CPM_DOWNLOAD_LOCATION}))
message(STATUS "Downloading [Link]")
file(DOWNLOAD
[Link]
${CPM_DOWNLOAD_LOCATION})
endif()
include(${CPM_DOWNLOAD_LOCATION})
CPMAddPackage(
NAME BS_thread_pool
GITHUB_REPOSITORY bshoshany/thread-pool
VERSION 4.0.0)
add_library(BS_thread_pool INTERFACE)
target_include_directories(BS_thread_pool INTERFACE ${BS_thread_pool_SOURCE_DIR}/include)
add_executable(my_project [Link])
target_link_libraries(my_project BS_thread_pool)
With both [Link] and [Link] in the same folder, type the following commands to build the
project:
39
cmake -S . -B build
cmake --build build
40
• Task submission without futures (T and F are template parameters):
– void detach_task(F&& task): Submit a function with no arguments and no return value into
the task queue. To push a function with arguments, enclose it in a lambda expression. Does
not return a future, so the user must use wait() or some other method to ensure that the task
finishes executing, otherwise bad things will happen.
– void detach_blocks(T first_index, T index_after_last, F&& block, size_t num_blocks
= 0): Parallelize a loop by automatically splitting it into blocks and submitting each block
separately to the queue. The block function takes two arguments, the start and end of the
block, so that it is only called only once per block, but it is up to the user make sure the block
function correctly deals with all the indices in each block. Does not return a BS::multi_future,
so the user must use wait() or some other method to ensure that the loop finishes executing,
otherwise bad things will happen.
– void detach_loop(T first_index, T index_after_last, F&& loop, size_t num_blocks
= 0): Parallelize a loop by automatically splitting it into blocks and submitting each block
separately to the queue. The loop function takes one argument, the loop index, so that it is
called many times per block. Does not return a BS::multi_future, so the user must use wait()
or some other method to ensure that the loop finishes executing, otherwise bad things will
happen.
– void detach_sequence(T first_index, T index_after_last, F&& sequence): Submit a
sequence of tasks enumerated by indices to the queue. Does not return a BS::multi_future, so
the user must use wait() or some other method to ensure that the sequence finishes executing,
otherwise bad things will happen.
• Task submission with futures (T, F, and R are template parameters):
– std::future<R> submit_task(F&& task): Submit a function with no arguments into the task
queue. To submit a function with arguments, enclose it in a lambda expression. If the function
has a return value, get a future for the eventual returned value. If the function has no return
value, get an std::future<void> which can be used to wait until the task finishes.
– BS::multi_future<R> submit_blocks(T first_index, T index_after_last, F&& block,
size_t num_blocks = 0): Parallelize a loop by automatically splitting it into blocks and
submitting each block separately to the queue. The block function takes two arguments, the
start and end of the block, so that it is only called only once per block, but it is up to the
user make sure the block function correctly deals with all the indices in each block. Returns a
BS::multi_future that contains the futures for all of the blocks.
– BS::multi_future<void> submit_loop(T first_index, T index_after_last, F&& loop,
size_t num_blocks = 0): Parallelize a loop by automatically splitting it into blocks and
submitting each block separately to the queue. The loop function takes one argument, the
loop index, so that it is called many times per block. It must have no return value. Returns a
BS::multi_future that contains the futures for all of the blocks.
– BS::multi_future<R> submit_sequence(T first_index, T index_after_last, F&& se-
quence): Submit a sequence of tasks enumerated by indices to the queue. Returns a
BS::multi_future that contains the futures for all of the tasks.
• Task management:
– void purge(): Purge all the tasks waiting in the queue. Tasks that are currently running will
not be affected, but any tasks still waiting in the queue will be discarded, and will never be
executed by the threads. Please note that there is no way to restore the purged tasks.
• Waiting for tasks (R and P, C, and D are template parameters):
– void wait(): Wait for tasks to be completed. Normally, this function waits for all tasks, both
those that are currently running in the threads and those that are still waiting in the queue.
However, if the pool is paused, this function only waits for the currently running tasks (otherwise
it would wait forever). Note: To wait for just one specific task, use submit_task() instead, and
call the wait() member function of the generated future.
– bool wait_for(std::chrono::duration<R, P>& duration): Wait for tasks to be completed,
but stop waiting after the specified duration has passed. Returns true if all tasks finished run-
ning, false if the duration expired but some tasks are still running.
41
– bool wait_until(std::chrono::time_point<C, D>& timeout_time): Wait for tasks to be
completed, but stop waiting after the specified time point has been reached. Returns true if all
tasks finished running, false if the time point was reached but some tasks are still running.
42
the member functions that can be used on an std::vector can also be used on a BS::multi_future. For
example, you may use a range-based for loop with a BS::multi_future, since it has iterators.
In addition to inherited member functions, BS::multi_future has the following specialized member func-
tions (R and P, C, and D are template parameters):
• [void or std::vector<T>] get(): Get the results from all the futures stored in this
BS::multi_future, rethrowing any stored exceptions. If the futures return void, this function re-
turns void as well. If the futures return a type T, this function returns a vector containing the results.
• size_t ready_count(): Check how many of the futures stored in this BS::multi_future are ready.
• bool valid(): Check if all the futures stored in this BS::multi_future are valid.
• void wait(): Wait for all the futures stored in this BS::multi_future.
• bool wait_for(std::chrono::duration<R, P>& duration): Wait for all the futures stored in this
BS::multi_future, but stop waiting after the specified duration has passed. Returns true if all
futures have been waited for before the duration expired, false otherwise.
• bool wait_until(std::chrono::time_point<C, D>& timeout_time): Wait for all the futures
stored in this multi_future object, but stop waiting after the specified time point has been reached.
Returns true if all futures have been waited for before the time point was reached, false otherwise.
43
• void start(): Start (or restart) measuring time. Note that the timer starts ticking as soon as the
object is created, so this is only necessary if we want to restart the clock later.
• void stop(): Stop measuring time and store the elapsed time since the object was constructed or
since start() was last called.
• std::chrono::milliseconds::rep current_ms(): Get the number of milliseconds that have
elapsed since the object was constructed or since start() was last called, but keep the timer ticking.
• std::chrono::milliseconds::rep ms(): Get the number of milliseconds stored when stop() was
last called.
11.2 Acknowledgements
Many GitHub users have helped improve this project, directly or indirectly, via issues, pull requests, com-
ments, and/or personal correspondence. Please see [Link] for links to specific issues and pull
requests that have been the most helpful. Thank you all for your contribution! :)
44
journal = {arXiv e-prints},
keywords = {Computer Science - Distributed, Parallel, and Cluster Computing,
D.1.3, D.1.5},
month = {May},
primaryclass = {[Link]},
title = {{A C++17 Thread Pool for High-Performance Scientific Computing}},
year = {2021}
}
Please note that the companion paper on arXiv is updated infrequently. The paper is intended to fa-
cilitate discovery of the library by scientists who may find it useful for scientific computing purposes
and to allow citing the library in scientific research, but most users should read the [Link] file on
the GitHub repository instead, as it is guaranteed to always be up to date.
References
[1] A. Williams, C++ Concurrency in Action. Manning Publications, 2019.
[2] B. Stroustrup, The C++ Programming Language. Pearson Education, 2013.
[3] B. Stroustrup, Programming: Principles and Practice Using C++. Pearson Education, 2014.
[4] B. Stroustrup, A Tour of C++. C++ In-Depth Series, Pearson Education, 2018.
[5] I. 14882:2017, Programming languages - C++. ISO, 2017.
[6] S. D. Casey, “How to determine the effectiveness of hyper-threading technology with an application,”
Intel Technology Journal, vol. 6, no. 1 p.11, 2011.
45