80 Iterator first, Iterator last, Callback &&callback,
81 const std::size_t parallelism = std::thread::hardware_concurrency(),
82 const std::size_t stack_size_bytes = 0) ->
void {
83 const auto effective_parallelism{(std::max)(parallelism, 1uz)};
92 if (effective_parallelism == 1 || std::next(first) == last) {
93 std::size_t cursor{1};
94 for (
auto iterator = first; iterator != last; ++iterator, ++cursor) {
95 callback(*iterator, effective_parallelism, cursor);
101 std::queue<Iterator> tasks;
102 for (
auto iterator = first; iterator != last; ++iterator) {
103 tasks.push(iterator);
106 std::mutex queue_mutex;
107 std::mutex exception_mutex;
109 auto effective_callback = std::forward<Callback>(callback);
111 std::exception_ptr exception =
nullptr;
112 auto handle_exception = [&exception_mutex,
113 &exception](
const std::exception_ptr &pointer) {
114 std::scoped_lock lock{exception_mutex};
120 std::vector<std::thread> workers;
121 workers.reserve(effective_parallelism);
123 const auto total{tasks.size()};
128 auto worker_callable = [&tasks, &queue_mutex, &effective_callback,
129 &handle_exception, effective_parallelism, total] {
133 std::size_t cursor{0};
135 std::scoped_lock lock{queue_mutex};
139 iterator = tasks.front();
140 cursor = total - tasks.size() + 1;
143 effective_callback(*iterator, effective_parallelism, cursor);
146 handle_exception(std::current_exception());
150 const char *creation_error{
nullptr};
155 for (std::size_t index = 0; index < effective_parallelism; ++index) {
156 auto *heap_function =
new std::function<void()>(worker_callable);
157 if (stack_size_bytes >
static_cast<std::size_t
>(UINT_MAX)) {
158 delete heap_function;
160 "The requested stack size is too large for this platform";
164 auto raw_handle = _beginthreadex(
165 nullptr,
static_cast<unsigned>(stack_size_bytes),
166 ¶llel_for_each_windows_thread_start, heap_function, 0,
nullptr);
167 if (raw_handle == 0) {
168 delete heap_function;
169 creation_error =
"Could not create thread";
173 HANDLE thread_handle =
reinterpret_cast<HANDLE
>(raw_handle);
174 workers.emplace_back([thread_handle] {
175 WaitForSingleObject(thread_handle, INFINITE);
176 CloseHandle(thread_handle);
180 for (std::size_t index = 0; index < effective_parallelism; ++index) {
184 pthread_attr_init(&attr);
189 if (stack_size_bytes > 0 &&
190 pthread_attr_setstacksize(&attr, stack_size_bytes) != 0) {
191 pthread_attr_destroy(&attr);
192 creation_error =
"The requested stack size is not supported by this "
197 auto *heap_function =
new std::function<void()>(worker_callable);
198 pthread_t pthread_handle;
199 auto raw_handle = pthread_create(
200 &pthread_handle, &attr,
201 [](
void *arg) ->
void * {
202 auto *function_ptr =
static_cast<std::function<
void()
> *>(arg);
208 if (raw_handle != 0) {
209 pthread_attr_destroy(&attr);
210 delete heap_function;
211 creation_error =
"Could not create thread";
214 workers.emplace_back(
215 [pthread_handle] { pthread_join(pthread_handle,
nullptr); });
216 pthread_attr_destroy(&attr);
225 if (creation_error !=
nullptr) {
226 std::scoped_lock lock{queue_mutex};
227 std::queue<Iterator> empty;
231 for (
auto &worker_thread : workers) {
232 worker_thread.join();
235 if (creation_error !=
nullptr) {
236 throw std::runtime_error(creation_error);
240 std::rethrow_exception(exception);