120#include <type_traits>
132 namespace parallel_detail
160 if (n == 0)
return 1;
162 const size_t chunks = num_threads * 4;
168 template <
typename Container>
171 using It =
decltype(std::begin(std::declval<Container &>()));
172 return std::is_base_of_v<std::random_access_iterator_tag,
173 typename std::iterator_traits<It>::iterator_category>;
178 template <
typename Container>
184 return std::make_unique<std::vector<typename Container::value_type>>(std::begin(c), std::end(c));
188 template <
typename T>
191 if constexpr (std::is_pointer_v<std::decay_t<T>>)
219 if (
options.cancel_token.stop_requested())
225 return pool.num_threads() <= 1;
268 template <
typename ResultT =
void,
typename Container,
typename Op>
270 size_t chunk_size = 0)
274 options.chunk_size = chunk_size;
278 template <
typename ResultT =
void,
typename Container,
typename Op>
283 using InputT = std::decay_t<
decltype(*std::begin(c))>;
285 std::is_void_v<ResultT>,
286 std::invoke_result_t<Op, const InputT &>,
289 const size_t n = std::distance(std::begin(c), std::end(c));
291 return std::vector<ActualResultT>{};
295 std::vector<ActualResultT> result;
297 for (
auto it = std::begin(c); it != std::end(c); ++it)
300 result.push_back(op(*it));
313 std::vector<std::optional<ActualResultT>> slots(n);
314 std::vector<std::future<void>>
futures;
323 token =
options.cancel_token]()
325 auto in_it = std::begin(data);
326 std::advance(in_it, offset);
327 for (size_t i = offset; i < chunk_end; ++i, ++in_it)
329 parallel_detail::throw_if_parallel_canceled(token);
330 slots[i] = op(*in_it);
338 std::exception_ptr
ep;
340 try { f.get(); }
catch (...) {
if (
not ep)
ep = std::current_exception(); }
342 std::rethrow_exception(
ep);
345 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
347 std::vector<ActualResultT> result;
349 for (
auto & slot : slots)
350 result.push_back(
std::move(*slot));
386 template <
typename Container,
typename Pred>
388 size_t chunk_size = 0)
392 options.chunk_size = chunk_size;
396 template <
typename Container,
typename Pred>
400 ThreadPool & pool = parallel_detail::selected_parallel_pool(
options);
401 using T = std::decay_t<
decltype(*std::begin(c))>;
403 const size_t n = std::distance(std::begin(c), std::end(c));
405 return std::vector<T>{};
407 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
409 std::vector<T> result;
410 for (
auto it = std::begin(c); it != std::end(c); ++it)
412 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
414 result.push_back(*it);
416 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
421 parallel_detail::effective_parallel_chunk_size(n, pool,
options);
423 auto data_holder = parallel_detail::ensure_random_access(c);
424 const auto & data = parallel_detail::deref(data_holder);
426 std::vector<std::future<std::vector<T>>> futures;
427 const size_t num_chunks = parallel_detail::chunk_count(n, chunk_size);
428 futures.reserve(num_chunks);
430 for (
size_t chunk_idx = 0; chunk_idx < num_chunks; ++chunk_idx)
432 const auto bounds = parallel_detail::bounds_for_chunk(chunk_idx, n, chunk_size);
434 futures.push_back(pool.enqueue([&data,
pred,
436 chunk_end = bounds.end,
437 token =
options.cancel_token]()
439 std::vector<T> chunk_result;
440 auto it = std::begin(data);
441 std::advance(it, offset);
442 for (size_t i = offset; i < chunk_end; ++i, ++it)
444 parallel_detail::throw_if_parallel_canceled(token);
446 chunk_result.push_back(*it);
452 std::vector<T> result;
453 for (
auto & f: futures)
455 auto chunk_result = f.get();
456 result.insert(result.end(),
457 std::make_move_iterator(chunk_result.begin()),
458 std::make_move_iterator(chunk_result.end()));
461 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
507 template <
typename T,
typename Container,
typename BinaryOp>
509 size_t chunk_size = 0)
513 options.chunk_size = chunk_size;
517 template <
typename T,
typename Container,
typename BinaryOp>
521 ThreadPool & pool = parallel_detail::selected_parallel_pool(
options);
522 const size_t n = std::distance(std::begin(c), std::end(c));
526 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
529 for (
auto it = std::begin(c); it != std::end(c); ++it)
531 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
532 result = op(result, *it);
534 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
539 parallel_detail::effective_parallel_chunk_size(n, pool,
options);
541 auto data_holder = parallel_detail::ensure_random_access(c);
542 const auto & data = parallel_detail::deref(data_holder);
544 std::vector<std::future<T>> futures;
549 size_t chunk_end = std::min(
offset + chunk_size, n);
551 futures.push_back(pool.enqueue([&data, op,
offset, chunk_end,
552 token =
options.cancel_token]()
554 auto it = std::begin(data);
555 std::advance(it, offset);
556 parallel_detail::throw_if_parallel_canceled(token);
558 for (size_t i = offset + 1; i < chunk_end; ++i, ++it)
560 parallel_detail::throw_if_parallel_canceled(token);
561 local = op(local, *it);
571 for (
auto & f: futures)
572 result = op(result, f.get());
574 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
613 template <
typename Container,
typename Op>
618 options.chunk_size = chunk_size;
622 template <
typename Container,
typename Op>
625 ThreadPool & pool = parallel_detail::selected_parallel_pool(
options);
626 const size_t n = std::distance(std::begin(c), std::end(c));
630 const size_t chunk_size =
631 parallel_detail::effective_parallel_chunk_size(n, pool,
options);
633 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
635 for (
auto it = std::begin(c); it != std::end(c); ++it)
637 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
640 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
644 std::vector<std::future<void>> futures;
649 size_t chunk_end = std::min(
offset + chunk_size, n);
651 futures.push_back(pool.enqueue([&c, op,
offset, chunk_end,
652 token =
options.cancel_token]()
654 auto it = std::begin(c);
655 std::advance(it, offset);
656 for (size_t i = offset; i < chunk_end; ++i, ++it)
658 parallel_detail::throw_if_parallel_canceled(token);
666 for (
auto & f: futures)
669 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
684 template <
typename Container,
typename Op>
689 options.chunk_size = chunk_size;
693 template <
typename Container,
typename Op>
696 ThreadPool & pool = parallel_detail::selected_parallel_pool(
options);
697 const size_t n = std::distance(std::begin(c), std::end(c));
701 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
703 for (
auto it = std::begin(c); it != std::end(c); ++it)
705 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
708 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
713 parallel_detail::effective_parallel_chunk_size(n, pool,
options);
715 auto data_holder = parallel_detail::ensure_random_access(c);
716 const auto & data = parallel_detail::deref(data_holder);
718 std::vector<std::future<void>> futures;
723 size_t chunk_end = std::min(
offset + chunk_size, n);
725 futures.push_back(pool.enqueue([&data, op,
offset, chunk_end,
726 token =
options.cancel_token]()
728 auto it = std::begin(data);
729 std::advance(it, offset);
730 for (size_t i = offset; i < chunk_end; ++i, ++it)
732 parallel_detail::throw_if_parallel_canceled(token);
740 for (
auto & f: futures)
743 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
772 template <
typename Container,
typename Pred>
774 size_t chunk_size = 0)
778 options.chunk_size = chunk_size;
782 template <
typename Container,
typename Pred>
786 ThreadPool & pool = parallel_detail::selected_parallel_pool(
options);
787 const size_t n = std::distance(std::begin(c), std::end(c));
791 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
793 for (
auto it = std::begin(c); it != std::end(c); ++it)
795 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
799 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
804 parallel_detail::effective_parallel_chunk_size(n, pool,
options);
806 auto data_holder = parallel_detail::ensure_random_access(c);
807 const auto & data = parallel_detail::deref(data_holder);
809 std::atomic<bool> found_false{
false};
810 std::vector<std::future<void>> futures;
815 size_t chunk_end = std::min(
offset + chunk_size, n);
817 futures.push_back(pool.enqueue([&data,
pred, &found_false,
offset, chunk_end,
818 token =
options.cancel_token]()
820 parallel_detail::throw_if_parallel_canceled(token);
821 if (found_false.load(std::memory_order_relaxed))
824 auto it = std::begin(data);
825 std::advance(it, offset);
826 for (size_t i = offset; i < chunk_end; ++i, ++it)
830 found_false.store(true, std::memory_order_relaxed);
833 parallel_detail::throw_if_parallel_canceled(token);
834 if (found_false.load(std::memory_order_relaxed))
842 for (
auto & f: futures)
845 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
846 return not found_false.load();
870 template <
typename Container,
typename Pred>
872 size_t chunk_size = 0)
876 options.chunk_size = chunk_size;
880 template <
typename Container,
typename Pred>
884 ThreadPool & pool = parallel_detail::selected_parallel_pool(
options);
885 const size_t n = std::distance(std::begin(c), std::end(c));
889 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
891 for (
auto it = std::begin(c); it != std::end(c); ++it)
893 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
897 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
902 parallel_detail::effective_parallel_chunk_size(n, pool,
options);
904 auto data_holder = parallel_detail::ensure_random_access(c);
905 const auto & data = parallel_detail::deref(data_holder);
907 std::atomic<bool> found{
false};
908 std::vector<std::future<void>> futures;
913 size_t chunk_end = std::min(
offset + chunk_size, n);
915 futures.push_back(pool.enqueue([&data,
pred, &found,
offset, chunk_end,
916 token =
options.cancel_token]()
918 parallel_detail::throw_if_parallel_canceled(token);
919 if (found.load(std::memory_order_relaxed))
922 auto it = std::begin(data);
923 std::advance(it, offset);
924 for (size_t i = offset; i < chunk_end; ++i, ++it)
928 found.store(true, std::memory_order_relaxed);
931 parallel_detail::throw_if_parallel_canceled(token);
932 if (found.load(std::memory_order_relaxed))
940 for (
auto & f: futures)
943 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
960 template <
typename Container,
typename Pred>
962 size_t chunk_size = 0)
967 template <
typename Container,
typename Pred>
999 template <
typename Container,
typename Pred>
1001 size_t chunk_size = 0)
1005 options.chunk_size = chunk_size;
1009 template <
typename Container,
typename Pred>
1013 ThreadPool & pool = parallel_detail::selected_parallel_pool(
options);
1014 const size_t n = std::distance(std::begin(c), std::end(c));
1018 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
1021 for (
auto it = std::begin(c); it != std::end(c); ++it)
1023 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1027 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1032 parallel_detail::effective_parallel_chunk_size(n, pool,
options);
1034 auto data_holder = parallel_detail::ensure_random_access(c);
1035 const auto & data = parallel_detail::deref(data_holder);
1037 std::vector<std::future<size_t>> futures;
1042 size_t chunk_end = std::min(
offset + chunk_size, n);
1044 futures.push_back(pool.enqueue([&data,
pred,
offset, chunk_end,
1045 token =
options.cancel_token]()
1048 auto it = std::begin(data);
1049 std::advance(it, offset);
1050 for (size_t i = offset; i < chunk_end; ++i, ++it)
1052 parallel_detail::throw_if_parallel_canceled(token);
1063 for (
auto & f: futures)
1066 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1099 template <
typename Container,
typename Pred>
1101 Pred
pred,
size_t chunk_size = 0)
1105 options.chunk_size = chunk_size;
1109 template <
typename Container,
typename Pred>
1114 ThreadPool & pool = parallel_detail::selected_parallel_pool(
options);
1115 const size_t n = std::distance(std::begin(c), std::end(c));
1117 return std::nullopt;
1119 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
1122 for (
auto it = std::begin(c); it != std::end(c); ++it, ++idx)
1124 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1128 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1129 return std::nullopt;
1133 parallel_detail::effective_parallel_chunk_size(n, pool,
options);
1135 auto data_holder = parallel_detail::ensure_random_access(c);
1136 const auto & data = parallel_detail::deref(data_holder);
1139 std::atomic<size_t> min_index{n};
1140 std::vector<std::future<void>> futures;
1145 size_t chunk_end = std::min(
offset + chunk_size, n);
1147 futures.push_back(pool.enqueue([&data,
pred, &min_index,
offset, chunk_end,
1148 token =
options.cancel_token]()
1150 parallel_detail::throw_if_parallel_canceled(token);
1152 if (min_index.load(std::memory_order_relaxed) <= offset)
1155 auto it = std::begin(data);
1156 std::advance(it, offset);
1157 for (size_t i = offset; i < chunk_end; ++i, ++it)
1160 if (min_index.load(std::memory_order_relaxed) <= i)
1163 parallel_detail::throw_if_parallel_canceled(token);
1167 size_t expected = min_index.load(std::memory_order_relaxed);
1168 while (i < expected and
1169 not min_index.compare_exchange_weak(expected, i,
1170 std::memory_order_relaxed));
1179 for (
auto & f: futures)
1182 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1183 if (
size_t result = min_index.load(); result < n)
1185 return std::nullopt;
1211 template <
typename Container,
typename Pred>
1213 Pred
pred,
size_t chunk_size = 0)
1217 options.chunk_size = chunk_size;
1221 template <
typename Container,
typename Pred>
1226 using T = std::decay_t<
decltype(*std::begin(c))>;
1230 return std::optional<T>{std::nullopt};
1232 auto it = std::begin(c);
1233 std::advance(it, *idx);
1234 return std::optional<T>{*it};
1261 typename T = std::decay_t<decltype(*std::begin(std::declval<Container>()))>>
1263 size_t chunk_size = 0)
1265 return pfoldl(pool, c, init, std::plus<T>{},
chunk_size);
1269 typename T = std::decay_t<decltype(*std::begin(std::declval<Container>()))>>
1287 typename T = std::decay_t<decltype(*std::begin(std::declval<Container>()))>>
1289 size_t chunk_size = 0)
1291 return pfoldl(pool, c, init, std::multiplies<T>{},
chunk_size);
1295 typename T = std::decay_t<decltype(*std::begin(std::declval<Container>()))>>
1311 template <
typename Container>
1316 options.chunk_size = chunk_size;
1320 template <
typename Container>
1323 ThreadPool & pool = parallel_detail::selected_parallel_pool(
options);
1324 using T = std::decay_t<
decltype(*std::begin(c))>;
1326 const size_t n = std::distance(std::begin(c), std::end(c));
1328 return std::optional<T>{std::nullopt};
1330 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
1332 auto it = std::begin(c);
1334 for (; it != std::end(c); ++it)
1336 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1340 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1341 return std::optional<T>{result};
1345 parallel_detail::effective_parallel_chunk_size(n, pool,
options);
1347 auto data_holder = parallel_detail::ensure_random_access(c);
1348 const auto & data = parallel_detail::deref(data_holder);
1350 std::vector<std::future<T>> futures;
1355 size_t chunk_end = std::min(
offset + chunk_size, n);
1357 futures.push_back(pool.enqueue([&data,
offset, chunk_end,
1358 token =
options.cancel_token]()
1360 auto it = std::begin(data);
1361 std::advance(it, offset);
1362 parallel_detail::throw_if_parallel_canceled(token);
1363 T local_min = *it++;
1364 for (size_t i = offset + 1; i < chunk_end; ++i, ++it)
1366 parallel_detail::throw_if_parallel_canceled(token);
1367 if (*it < local_min)
1376 T result = futures[0].get();
1377 for (
size_t i = 1; i < futures.size(); ++i)
1379 T val = futures[i].get();
1384 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1385 return std::optional<T>{result};
1398 template <
typename Container>
1403 options.chunk_size = chunk_size;
1407 template <
typename Container>
1410 ThreadPool & pool = parallel_detail::selected_parallel_pool(
options);
1411 using T = std::decay_t<
decltype(*std::begin(c))>;
1413 const size_t n = std::distance(std::begin(c), std::end(c));
1415 return std::optional<T>{std::nullopt};
1417 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
1419 auto it = std::begin(c);
1421 for (; it != std::end(c); ++it)
1423 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1427 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1428 return std::optional<T>{result};
1432 parallel_detail::effective_parallel_chunk_size(n, pool,
options);
1434 auto data_holder = parallel_detail::ensure_random_access(c);
1435 const auto & data = parallel_detail::deref(data_holder);
1437 std::vector<std::future<T>> futures;
1442 size_t chunk_end = std::min(
offset + chunk_size, n);
1444 futures.push_back(pool.enqueue([&data,
offset, chunk_end,
1445 token =
options.cancel_token]()
1447 auto it = std::begin(data);
1448 std::advance(it, offset);
1449 parallel_detail::throw_if_parallel_canceled(token);
1450 T local_max = *it++;
1451 for (size_t i = offset + 1; i < chunk_end; ++i, ++it)
1453 parallel_detail::throw_if_parallel_canceled(token);
1454 if (*it > local_max)
1463 T result = futures[0].get();
1464 for (
size_t i = 1; i < futures.size(); ++i)
1466 T val = futures[i].get();
1471 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1472 return std::optional<T>{result};
1485 template <
typename Container>
1488 using T = std::decay_t<
decltype(*std::begin(c))>;
1490 const size_t n = std::distance(std::begin(c), std::end(c));
1492 return std::optional<std::pair<T, T>>{std::nullopt};
1494 if (chunk_size == 0)
1495 chunk_size = parallel_detail::chunk_size(n, pool.
num_threads());
1497 auto data_holder = parallel_detail::ensure_random_access(c);
1498 const auto & data = parallel_detail::deref(data_holder);
1500 std::vector<std::future<std::pair<T, T>>> futures;
1505 size_t chunk_end = std::min(
offset + chunk_size, n);
1509 auto it = std::begin(data);
1510 std::advance(it, offset);
1512 T local_max = *it++;
1513 for (size_t i = offset + 1; i < chunk_end; ++i, ++it)
1515 if (*it < local_min) local_min = *it;
1516 if (*it > local_max) local_max = *it;
1518 return std::make_pair(local_min, local_max);
1524 auto result = futures[0].get();
1525 for (
size_t i = 1; i < futures.size(); ++i)
1527 auto [mi, ma] = futures[i].get();
1528 if (mi < result.first) result.first = mi;
1529 if (ma > result.second) result.second = ma;
1532 return std::optional<std::pair<T, T>>{result};
1566 template <
typename Container,
typename Compare = std::less<>>
1568 const size_t min_parallel_size = 1024)
1572 options.min_size = min_parallel_size;
1576 template <
typename Container,
typename Compare = std::less<>>
1578 const ParallelOptions &
options = {})
1580 const size_t n = std::distance(std::begin(c), std::end(c));
1584 auto & pool = parallel_detail::selected_parallel_pool(
options);
1585 const size_t min_parallel_size =
options.min_size == 0 ? 1024 :
options.min_size;
1588 if (parallel_detail::use_sequential_parallel_path(n, pool,
options)
1589 or n <= min_parallel_size)
1591 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1597 size_t num_chunks = std::min(pool.
num_threads() * 2, n / min_parallel_size);
1599 num_chunks = std::min(num_chunks,
options.max_tasks);
1600 num_chunks = std::max<size_t>(1, num_chunks);
1601 const size_t chunk_size = (n + num_chunks - 1) / num_chunks;
1604 std::vector<std::future<void>> futures;
1607 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1608 size_t end = std::min(i + chunk_size, n);
1609 auto begin_it = std::begin(c);
1610 std::advance(begin_it, i);
1611 auto end_it = std::begin(c);
1612 std::advance(end_it, end);
1614 futures.push_back(pool.
enqueue([begin_it, end_it,
cmp, token =
options.cancel_token]()
1616 parallel_detail::throw_if_parallel_canceled(token);
1617 std::sort(begin_it, end_it, cmp);
1621 for (
auto & f: futures)
1625 using T = std::decay_t<
decltype(*std::begin(c))>;
1626 std::vector<std::optional<T>> buffer(n);
1628 ParallelOptions merge_options =
options;
1629 merge_options.pool = &pool;
1631 for (
size_t width = chunk_size; width < n; width *= 2)
1633 for (
size_t i = 0; i < n; i += 2 * width)
1635 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1636 size_t mid = std::min(i + width, n);
1637 size_t end = std::min(i + 2 * width, n);
1641 auto begin_it = std::begin(c);
1642 std::advance(begin_it, i);
1643 auto mid_it = std::begin(c);
1644 std::advance(mid_it, mid);
1645 auto end_it = std::begin(c);
1646 std::advance(end_it, end);
1649 buffer.begin() + i,
cmp, merge_options);
1654 auto begin_it = std::begin(c);
1655 std::advance(begin_it, i);
1656 auto end_it = std::begin(c);
1657 std::advance(end_it, mid);
1658 std::copy(begin_it, end_it, buffer.begin() + i);
1663 auto it = std::begin(c);
1664 for (
size_t i = 0; i < n; ++i, ++it)
1665 *it = std::move(*buffer[i]);
1700 template <
typename Container1,
typename Container2,
typename Op>
1702 Op op,
size_t chunk_size = 0)
1706 options.chunk_size = chunk_size;
1710 template <
typename Container1,
typename Container2,
typename Op>
1714 const size_t n1 = std::distance(std::begin(c1), std::end(c1));
1715 const size_t n2 = std::distance(std::begin(c2), std::end(c2));
1716 const size_t n = std::min(n1, n2);
1721 auto & pool = parallel_detail::selected_parallel_pool(
options);
1722 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
1724 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1725 auto it1 = std::begin(c1);
1726 auto it2 = std::begin(c2);
1727 for (
size_t i = 0; i < n; ++i, ++it1, ++it2)
1729 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1735 size_t chunk_size = parallel_detail::effective_parallel_chunk_size(n, pool,
options,
1738 auto h1 = parallel_detail::ensure_random_access(c1);
1739 auto h2 = parallel_detail::ensure_random_access(c2);
1740 const auto & d1 = parallel_detail::deref(h1);
1741 const auto & d2 = parallel_detail::deref(h2);
1743 std::vector<std::future<void>> futures;
1748 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1749 size_t chunk_end = std::min(
offset + chunk_size, n);
1751 futures.push_back(pool.
enqueue([&d1, &d2, op,
offset, chunk_end,
1752 token =
options.cancel_token]()
1754 auto it1 = std::begin(d1);
1755 auto it2 = std::begin(d2);
1756 std::advance(it1, offset);
1757 std::advance(it2, offset);
1758 for (size_t i = offset; i < chunk_end; ++i, ++it1, ++it2)
1760 parallel_detail::throw_if_parallel_canceled(token);
1768 for (
auto & f: futures)
1796 template <
typename Container1,
typename Container2,
typename Op>
1798 const Container2 & c2, Op op,
size_t chunk_size = 0)
1802 options.chunk_size = chunk_size;
1806 template <
typename Container1,
typename Container2,
typename Op>
1807 [[nodiscard]]
auto pzip_maps(
const Container1 & c1,
const Container2 & c2, Op op,
1810 using T1 = std::decay_t<
decltype(*std::begin(c1))>;
1811 using T2 = std::decay_t<
decltype(*std::begin(c2))>;
1812 using ResultT = std::invoke_result_t<Op, const T1 &, const T2 &>;
1814 const size_t n1 = std::distance(std::begin(c1), std::end(c1));
1815 const size_t n2 = std::distance(std::begin(c2), std::end(c2));
1816 const size_t n = std::min(n1, n2);
1819 return std::vector<ResultT>{};
1821 auto & pool = parallel_detail::selected_parallel_pool(
options);
1822 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
1824 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1825 std::vector<ResultT> result;
1827 auto it1 = std::begin(c1);
1828 auto it2 = std::begin(c2);
1829 for (
size_t i = 0; i < n; ++i, ++it1, ++it2)
1831 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1832 result.push_back(op(*it1, *it2));
1837 size_t chunk_size = parallel_detail::effective_parallel_chunk_size(n, pool,
options,
1840 auto h1 = parallel_detail::ensure_random_access(c1);
1841 auto h2 = parallel_detail::ensure_random_access(c2);
1842 const auto & d1 = parallel_detail::deref(h1);
1843 const auto & d2 = parallel_detail::deref(h2);
1845 std::vector<std::optional<ResultT>> slots(n);
1846 std::vector<std::future<void>> futures;
1851 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1852 size_t chunk_end = std::min(
offset + chunk_size, n);
1854 futures.push_back(pool.
enqueue([&slots, &d1, &d2, op,
offset, chunk_end,
1855 token =
options.cancel_token]()
1857 auto it1 = std::begin(d1);
1858 auto it2 = std::begin(d2);
1859 std::advance(it1, offset);
1860 std::advance(it2, offset);
1861 for (size_t i = offset; i < chunk_end; ++i, ++it1, ++it2)
1863 parallel_detail::throw_if_parallel_canceled(token);
1864 slots[i] = op(*it1, *it2);
1871 for (
auto & f: futures)
1874 std::vector<ResultT> result;
1876 for (
auto & slot : slots)
1877 result.push_back(
std::move(*slot));
1910 template <
typename Container1,
typename Container2,
typename T,
typename Op>
1912 const Container2 & c2,
T init, Op op,
1913 size_t chunk_size = 0)
1917 options.chunk_size = chunk_size;
1921 template <
typename Container1,
typename Container2,
typename T,
typename Op>
1922 [[nodiscard]]
T pzip_foldl(
const Container1 & c1,
const Container2 & c2,
T init, Op op,
1925 const size_t n1 = std::distance(std::begin(c1), std::end(c1));
1926 const size_t n2 = std::distance(std::begin(c2), std::end(c2));
1927 const size_t n = std::min(n1, n2);
1932 auto & pool = parallel_detail::selected_parallel_pool(
options);
1933 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
1935 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1936 auto it1 = std::begin(c1);
1937 auto it2 = std::begin(c2);
1939 for (
size_t i = 0; i < n; ++i, ++it1, ++it2)
1941 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1942 result = op(result, *it1, *it2);
1947 size_t chunk_size = parallel_detail::effective_parallel_chunk_size(n, pool,
options,
1950 auto h1 = parallel_detail::ensure_random_access(c1);
1951 auto h2 = parallel_detail::ensure_random_access(c2);
1952 const auto & d1 = parallel_detail::deref(h1);
1953 const auto & d2 = parallel_detail::deref(h2);
1955 std::vector<std::future<T>> futures;
1960 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
1961 size_t chunk_end = std::min(
offset + chunk_size, n);
1963 futures.push_back(pool.
enqueue([&d1, &d2, &init, op,
offset, chunk_end,
1964 token =
options.cancel_token]()
1966 auto it1 = std::begin(d1);
1967 auto it2 = std::begin(d2);
1968 std::advance(it1, offset);
1969 std::advance(it2, offset);
1971 parallel_detail::throw_if_parallel_canceled(token);
1972 T local = op(init, *it1++, *it2++);
1973 for (size_t i = offset + 1; i < chunk_end; ++i, ++it1, ++it2)
1975 parallel_detail::throw_if_parallel_canceled(token);
1976 local = op(local, *it1, *it2);
1986 T result = futures[0].get();
1987 for (
size_t i = 1; i < futures.size(); ++i)
1989 T val = futures[i].get();
1991 result = result + val - init;
2022 template <
typename Container,
typename Pred>
2024 size_t chunk_size = 0)
2028 options.chunk_size = chunk_size;
2032 template <
typename Container,
typename Pred>
2036 using T = std::decay_t<
decltype(*std::begin(c))>;
2037 ThreadPool & pool = parallel_detail::selected_parallel_pool(
options);
2039 const size_t n = std::distance(std::begin(c), std::end(c));
2041 return std::make_pair(std::vector<T>{}, std::vector<T>{});
2043 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
2045 std::vector<T> yes_result, no_result;
2046 for (
auto it = std::begin(c); it != std::end(c); ++it)
2048 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2050 yes_result.push_back(*it);
2052 no_result.push_back(*it);
2054 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2055 return std::make_pair(std::move(yes_result), std::move(no_result));
2059 parallel_detail::effective_parallel_chunk_size(n, pool,
options);
2061 auto data_holder = parallel_detail::ensure_random_access(c);
2062 const auto & data = parallel_detail::deref(data_holder);
2064 const size_t num_chunks = parallel_detail::chunk_count(n, chunk_size);
2065 std::vector<size_t> yes_counts(num_chunks, 0);
2066 std::vector<size_t> no_counts(num_chunks, 0);
2067 std::vector<std::future<void>> count_futures;
2068 count_futures.reserve(num_chunks);
2070 for (
size_t chunk_idx = 0; chunk_idx < num_chunks; ++chunk_idx)
2072 const auto bounds = parallel_detail::bounds_for_chunk(chunk_idx, n, chunk_size);
2074 count_futures.push_back(pool.enqueue([&data, &yes_counts, &no_counts,
pred,
2077 chunk_end = bounds.end,
2078 token =
options.cancel_token]()
2082 auto it = std::begin(data);
2083 std::advance(it, offset);
2084 for (size_t i = offset; i < chunk_end; ++i, ++it)
2086 parallel_detail::throw_if_parallel_canceled(token);
2092 yes_counts[chunk_idx] = yes;
2093 no_counts[chunk_idx] = no;
2097 for (
auto & f: count_futures)
2100 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2102 std::vector<size_t> yes_offsets(num_chunks, 0);
2103 std::vector<size_t> no_offsets(num_chunks, 0);
2107 size_t{0}, std::plus<size_t>{},
options);
2109 size_t{0}, std::plus<size_t>{},
options);
2112 const size_t yes_total = num_chunks == 0 ? 0 : yes_offsets.back() + yes_counts.back();
2113 const size_t no_total = num_chunks == 0 ? 0 : no_offsets.back() + no_counts.back();
2115 std::vector<std::optional<T>> yes_slots(yes_total);
2116 std::vector<std::optional<T>> no_slots(no_total);
2117 std::vector<std::future<void>> fill_futures;
2118 fill_futures.reserve(num_chunks);
2120 for (
size_t chunk_idx = 0; chunk_idx < num_chunks; ++chunk_idx)
2122 const auto bounds = parallel_detail::bounds_for_chunk(chunk_idx, n, chunk_size);
2123 const size_t yes_offset = yes_offsets[chunk_idx];
2124 const size_t no_offset = no_offsets[chunk_idx];
2126 fill_futures.push_back(pool.enqueue([&data, &yes_slots, &no_slots,
pred,
2128 chunk_end = bounds.end,
2129 yes_offset, no_offset,
2130 token =
options.cancel_token]()
2132 auto it = std::begin(data);
2133 std::advance(it, offset);
2134 size_t yes_pos = yes_offset;
2135 size_t no_pos = no_offset;
2136 for (size_t i = offset; i < chunk_end; ++i, ++it)
2138 parallel_detail::throw_if_parallel_canceled(token);
2140 yes_slots[yes_pos++].emplace(*it);
2142 no_slots[no_pos++].emplace(*it);
2147 for (
auto & f: fill_futures)
2150 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2152 std::vector<T> yes_result, no_result;
2153 yes_result.reserve(yes_total);
2154 no_result.reserve(no_total);
2156 for (
auto & slot: yes_slots)
2157 yes_result.push_back(
std::move(*slot));
2158 for (
auto & slot: no_slots)
2159 no_result.push_back(
std::move(*slot));
2161 return std::make_pair(std::move(yes_result), std::move(no_result));
2191 template <
typename Container,
typename BinaryOp>
2193 size_t chunk_size = 0)
2197 options.chunk_size = chunk_size;
2221 template <
typename Container,
typename BinaryOp>
2225 using T = std::decay_t<
decltype(*std::begin(c))>;
2227 auto holder = parallel_detail::ensure_random_access(c);
2228 const auto & data = parallel_detail::deref(holder);
2229 const size_t n =
static_cast<size_t>(std::distance(std::begin(data), std::end(data)));
2230 std::vector<std::optional<T>> slots(n);
2234 std::vector<T> result;
2236 for (
auto & slot: slots)
2237 result.push_back(
std::move(*slot));
2267 template <
typename Container,
typename T,
typename BinaryOp>
2269 BinaryOp op,
size_t chunk_size = 0)
2273 options.chunk_size = chunk_size;
2297 template <
typename Container,
typename T,
typename BinaryOp>
2301 auto holder = parallel_detail::ensure_random_access(c);
2302 const auto & data = parallel_detail::deref(holder);
2303 const size_t n =
static_cast<size_t>(std::distance(std::begin(data), std::end(data)));
2304 std::vector<std::optional<T>> slots(n);
2309 std::vector<T> result;
2311 for (
auto & slot: slots)
2312 result.push_back(
std::move(*slot));
2341 template <
typename Container1,
typename Container2,
typename Compare = std::less<>>
2343 Compare comp = Compare{},
size_t chunk_size = 0)
2347 options.chunk_size = chunk_size;
2348 return pmerge(c1, c2, std::move(comp),
options);
2371 template <
typename Container1,
typename Container2,
typename Compare = std::less<>>
2372 [[nodiscard]]
auto pmerge(
const Container1 & c1,
const Container2 & c2,
2373 Compare comp = Compare{},
2374 const ParallelOptions &
options = {})
2376 using T1 = std::decay_t<
decltype(*std::begin(c1))>;
2377 using T2 = std::decay_t<
decltype(*std::begin(c2))>;
2378 using ResultT = std::common_type_t<T1, T2>;
2380 auto h1 = parallel_detail::ensure_random_access(c1);
2381 auto h2 = parallel_detail::ensure_random_access(c2);
2382 const auto & d1 = parallel_detail::deref(h1);
2383 const auto & d2 = parallel_detail::deref(h2);
2385 const size_t total =
static_cast<size_t>(std::distance(std::begin(d1), std::end(d1))
2386 + std::distance(std::begin(d2), std::end(d2)));
2387 std::vector<std::optional<ResultT>> slots(total);
2390 std::begin(d2), std::end(d2),
2391 slots.begin(), comp,
options);
2393 std::vector<ResultT> result;
2394 result.reserve(total);
2395 for (
auto & slot: slots)
2396 result.push_back(
std::move(*slot));
2404 namespace parallel_zip_detail
2414 template <
typename Container>
2417 using value_type = std::decay_t<decltype(*std::begin(std::declval<Container &>()))>;
2419 parallel_detail::has_random_access<Container>(),
2421 std::unique_ptr<std::vector<value_type>>>;
2428 if constexpr (parallel_detail::has_random_access<Container>())
2432 cached_size =
static_cast<size_t>(std::distance(std::begin(c), std::end(c)));
2437 data = std::make_unique<std::vector<value_type>>(std::begin(c), std::end(c));
2438 cached_size = data->size();
2444 if constexpr (parallel_detail::has_random_access<Container>())
2451 [[nodiscard]]
size_t size() const noexcept {
return cached_size; }
2453 auto begin()
const {
return std::begin(get()); }
2454 auto end()
const {
return std::end(get()); }
2458 template <
typename... Holders,
size_t... Is>
2460 std::index_sequence<Is...>)
2462 return std::min({std::get<Is>(holders).size()...});
2465 template <
typename... Holders>
2472 template <
typename... Holders,
size_t... Is>
2474 std::index_sequence<Is...>)
2476 return std::make_tuple([&]()
2478 auto it = std::get<Is>(holders).begin();
2479 std::advance(it,
offset);
2485 template <
typename... Iters,
size_t... Is>
2488 (++std::get<Is>(iters), ...);
2492 template <
typename... Iters,
size_t... Is>
2495 return std::make_tuple(*std::get<Is>(iters)...);
2531 template <
typename Op,
typename... Containers>
2539 template <
typename Op,
typename... Containers>
2542 static_assert(
sizeof...(Containers) >= 2,
2543 "pzip_for_each requires at least 2 containers");
2550 const size_t n = parallel_zip_detail::min_holder_size(holders);
2554 auto & pool = parallel_detail::selected_parallel_pool(
options);
2555 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
2557 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2558 constexpr size_t N =
sizeof...(Containers);
2559 auto iters = parallel_zip_detail::make_iterators_at(
2560 0, holders, std::make_index_sequence<N>{});
2561 for (
size_t i = 0; i < n; ++i)
2563 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2564 std::apply(op, parallel_zip_detail::deref_all_iters(
2565 iters, std::make_index_sequence<N>{}));
2566 parallel_zip_detail::advance_all_iters(iters, std::make_index_sequence<N>{});
2571 const size_t chunk_size = parallel_detail::effective_parallel_chunk_size(
2574 std::vector<std::future<void>> futures;
2579 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2580 size_t chunk_end = std::min(
offset + chunk_size, n);
2582 futures.push_back(pool.
enqueue([&holders, op,
offset, chunk_end,
2583 token =
options.cancel_token]()
2585 constexpr size_t N = sizeof...(Containers);
2586 auto iters = parallel_zip_detail::make_iterators_at(
2587 offset, holders, std::make_index_sequence<N>{});
2589 for (
size_t i =
offset; i < chunk_end; ++i)
2591 parallel_detail::throw_if_parallel_canceled(token);
2592 std::apply(op, parallel_zip_detail::deref_all_iters(
2593 iters, std::make_index_sequence<N>{}));
2594 parallel_zip_detail::advance_all_iters(
2595 iters, std::make_index_sequence<N>{});
2602 for (
auto & f: futures)
2637 template <
typename Op,
typename... Containers>
2645 template <
typename Op,
typename... Containers>
2647 const Containers &... cs)
2649 static_assert(
sizeof...(Containers) >= 2,
2650 "pzip_maps requires at least 2 containers");
2653 using ResultT = std::invoke_result_t<Op,
2654 std::decay_t<
decltype(*std::begin(cs))>...>;
2660 const size_t n = parallel_zip_detail::min_holder_size(holders);
2662 return std::vector<ResultT>{};
2664 auto & pool = parallel_detail::selected_parallel_pool(
options);
2665 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
2667 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2668 std::vector<std::optional<ResultT>> slots(n);
2669 constexpr size_t N =
sizeof...(Containers);
2670 auto iters = parallel_zip_detail::make_iterators_at(
2671 0, holders, std::make_index_sequence<N>{});
2672 for (
size_t i = 0; i < n; ++i)
2674 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2675 slots[i] = std::apply(op, parallel_zip_detail::deref_all_iters(
2676 iters, std::make_index_sequence<N>{}));
2677 parallel_zip_detail::advance_all_iters(iters, std::make_index_sequence<N>{});
2680 std::vector<ResultT> result;
2682 for (
auto & slot: slots)
2683 result.push_back(std::move(*slot));
2687 size_t chunk_size = parallel_detail::effective_parallel_chunk_size(n, pool,
options,
2690 std::vector<std::optional<ResultT>> slots(n);
2691 std::vector<std::future<void>> futures;
2696 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2697 size_t chunk_end = std::min(
offset + chunk_size, n);
2699 futures.push_back(pool.
enqueue([&slots, &holders, op,
offset, chunk_end,
2700 token =
options.cancel_token]()
2702 constexpr size_t N = sizeof...(Containers);
2703 auto iters = parallel_zip_detail::make_iterators_at(
2704 offset, holders, std::make_index_sequence<N>{});
2706 for (
size_t i =
offset; i < chunk_end; ++i)
2708 parallel_detail::throw_if_parallel_canceled(token);
2709 slots[i] = std::apply(op, parallel_zip_detail::deref_all_iters(
2710 iters, std::make_index_sequence<N>{}));
2711 parallel_zip_detail::advance_all_iters(iters, std::make_index_sequence<N>{});
2717 for (
auto & f: futures)
2720 std::vector<ResultT> result;
2722 for (
auto & slot: slots)
2723 result.push_back(
std::move(*slot));
2766 template <
typename T,
typename Op,
typename Combiner,
typename... Containers>
2768 const Containers &... cs)
2775 template <
typename T,
typename Op,
typename Combiner,
typename... Containers>
2779 static_assert(
sizeof...(Containers) >= 2,
2780 "pzip_foldl requires at least 2 containers");
2786 const size_t n = parallel_zip_detail::min_holder_size(holders);
2790 auto & pool = parallel_detail::selected_parallel_pool(
options);
2791 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
2793 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2794 constexpr size_t N =
sizeof...(Containers);
2795 auto iters = parallel_zip_detail::make_iterators_at(
2796 0, holders, std::make_index_sequence<N>{});
2798 for (
size_t i = 0; i < n; ++i)
2800 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2801 auto tuple = parallel_zip_detail::deref_all_iters(iters, std::make_index_sequence<N>{});
2802 result = std::apply([&op, &result](
auto &&... args)
2804 return op(result, std::forward<
decltype(args)>(args)...);
2806 parallel_zip_detail::advance_all_iters(iters, std::make_index_sequence<N>{});
2811 size_t chunk_size = parallel_detail::effective_parallel_chunk_size(n, pool,
options,
2814 std::vector<std::future<T>> futures;
2819 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2820 size_t chunk_end = std::min(
offset + chunk_size, n);
2822 futures.push_back(pool.
enqueue([&holders, init, op,
offset, chunk_end,
2823 token =
options.cancel_token]()
2825 constexpr size_t N = sizeof...(Containers);
2826 auto iters = parallel_zip_detail::make_iterators_at(
2827 offset, holders, std::make_index_sequence<N>{});
2830 parallel_detail::throw_if_parallel_canceled(token);
2831 auto first_tuple = parallel_zip_detail::deref_all_iters(
2832 iters, std::make_index_sequence<N>{});
2833 T local = std::apply([&op, &init](
auto &&... args)
2835 return op(init, std::forward<
decltype(args)>(args)
2838 parallel_zip_detail::advance_all_iters(
2839 iters, std::make_index_sequence<N>{});
2842 for (
size_t i =
offset + 1; i < chunk_end; ++i)
2844 parallel_detail::throw_if_parallel_canceled(token);
2845 auto tuple = parallel_zip_detail::deref_all_iters(
2846 iters, std::make_index_sequence<N>{});
2847 local = std::apply([&op, &local](
auto &&... args)
2850 std::forward<
decltype(args)>(args)
2853 parallel_zip_detail::advance_all_iters(
2854 iters, std::make_index_sequence<N>{});
2864 T result = futures[0].get();
2865 for (
size_t i = 1; i < futures.size(); ++i)
2866 result = combiner(result, futures[i].get());
2899 template <
typename Pred,
typename... Containers>
2907 template <
typename Pred,
typename... Containers>
2909 const Containers &... cs)
2911 static_assert(
sizeof...(Containers) >= 2,
2912 "pzip_all requires at least 2 containers");
2915 auto holders = std::make_tuple(
2919 const size_t n = parallel_zip_detail::min_holder_size(holders);
2923 auto & pool = parallel_detail::selected_parallel_pool(
options);
2924 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
2926 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2927 constexpr size_t N =
sizeof...(Containers);
2928 auto iters = parallel_zip_detail::make_iterators_at(
2929 0, holders, std::make_index_sequence<N>{});
2930 for (
size_t i = 0; i < n; ++i)
2932 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2933 auto tuple = parallel_zip_detail::deref_all_iters(iters, std::make_index_sequence<N>{});
2934 if (! std::apply(
pred, tuple))
2936 parallel_zip_detail::advance_all_iters(iters, std::make_index_sequence<N>{});
2941 const size_t chunk_size = parallel_detail::effective_parallel_chunk_size(
2944 std::atomic<bool> found_false{
false};
2945 std::vector<std::future<void>> futures;
2950 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
2951 size_t chunk_end = std::min(
offset + chunk_size, n);
2953 futures.push_back(pool.
enqueue([&holders,
pred, &found_false,
offset, chunk_end,
2954 token =
options.cancel_token]()
2956 if (found_false.load(std::memory_order_relaxed))
2959 constexpr size_t N = sizeof...(Containers);
2960 auto iters = parallel_zip_detail::make_iterators_at(
2961 offset, holders, std::make_index_sequence<N>{});
2963 for (
size_t i =
offset; i < chunk_end; ++i)
2965 parallel_detail::throw_if_parallel_canceled(token);
2966 auto tuple = parallel_zip_detail::deref_all_iters(
2967 iters, std::make_index_sequence<N>{});
2968 if (! std::apply(
pred, tuple))
2970 found_false.store(
true, std::memory_order_relaxed);
2973 if (found_false.load(std::memory_order_relaxed))
2975 parallel_zip_detail::advance_all_iters(
2976 iters, std::make_index_sequence<N>{});
2983 for (
auto & f: futures)
2986 return not found_false.load();
3017 template <
typename Pred,
typename... Containers>
3025 template <
typename Pred,
typename... Containers>
3027 const Containers &... cs)
3029 static_assert(
sizeof...(Containers) >= 2,
3030 "pzip_exists requires at least 2 containers");
3036 const size_t n = parallel_zip_detail::min_holder_size(holders);
3040 auto & pool = parallel_detail::selected_parallel_pool(
options);
3041 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
3043 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3044 constexpr size_t N =
sizeof...(Containers);
3045 auto iters = parallel_zip_detail::make_iterators_at(
3046 0, holders, std::make_index_sequence<N>{});
3047 for (
size_t i = 0; i < n; ++i)
3049 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3050 auto tuple = parallel_zip_detail::deref_all_iters(iters, std::make_index_sequence<N>{});
3051 if (std::apply(
pred, tuple))
3053 parallel_zip_detail::advance_all_iters(iters, std::make_index_sequence<N>{});
3058 const size_t chunk_size = parallel_detail::effective_parallel_chunk_size(
3061 std::atomic<bool> found{
false};
3062 std::vector<std::future<void>> futures;
3067 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3068 size_t chunk_end = std::min(
offset + chunk_size, n);
3071 token =
options.cancel_token]()
3073 if (found.load(std::memory_order_relaxed))
3076 constexpr size_t N = sizeof...(Containers);
3077 auto iters = parallel_zip_detail::make_iterators_at(
3078 offset, holders, std::make_index_sequence<N>{});
3080 for (
size_t i =
offset; i < chunk_end; ++i)
3082 parallel_detail::throw_if_parallel_canceled(token);
3083 auto tuple = parallel_zip_detail::deref_all_iters(
3084 iters, std::make_index_sequence<N>{});
3085 if (std::apply(
pred, tuple))
3087 found.store(
true, std::memory_order_relaxed);
3090 if (found.load(std::memory_order_relaxed))
3092 parallel_zip_detail::advance_all_iters(
3093 iters, std::make_index_sequence<N>{});
3100 for (
auto & f: futures)
3103 return found.load();
3121 template <
typename Pred,
typename... Containers>
3123 const Containers &... cs)
3130 template <
typename Pred,
typename... Containers>
3132 const Containers &... cs)
3134 static_assert(
sizeof...(Containers) >= 2,
3135 "pzip_count_if requires at least 2 containers");
3141 const size_t n = parallel_zip_detail::min_holder_size(holders);
3145 auto & pool = parallel_detail::selected_parallel_pool(
options);
3146 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
3148 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3149 constexpr size_t N =
sizeof...(Containers);
3150 auto iters = parallel_zip_detail::make_iterators_at(
3151 0, holders, std::make_index_sequence<N>{});
3153 for (
size_t i = 0; i < n; ++i)
3155 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3156 auto tuple = parallel_zip_detail::deref_all_iters(iters, std::make_index_sequence<N>{});
3157 if (std::apply(
pred, tuple))
3159 parallel_zip_detail::advance_all_iters(iters, std::make_index_sequence<N>{});
3164 size_t chunk_size = parallel_detail::effective_parallel_chunk_size(n, pool,
options,
3167 std::vector<std::future<size_t>> futures;
3172 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3173 size_t chunk_end = std::min(
offset + chunk_size, n);
3176 token =
options.cancel_token]()
3178 constexpr size_t N = sizeof...(Containers);
3179 auto iters = parallel_zip_detail::make_iterators_at(
3180 offset, holders, std::make_index_sequence<N>{});
3183 for (
size_t i =
offset; i < chunk_end; ++i)
3185 parallel_detail::throw_if_parallel_canceled(token);
3186 auto tuple = parallel_zip_detail::deref_all_iters(
3187 iters, std::make_index_sequence<N>{});
3188 if (std::apply(
pred, tuple))
3190 parallel_zip_detail::advance_all_iters(
3191 iters, std::make_index_sequence<N>{});
3200 for (
auto & f: futures)
3235 template <
typename Container,
typename Op>
3237 size_t chunk_size = 0)
3241 options.chunk_size = chunk_size;
3245 template <
typename Container,
typename Op>
3249 const size_t n = std::distance(std::begin(c), std::end(c));
3253 auto & pool = parallel_detail::selected_parallel_pool(
options);
3254 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
3256 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3257 auto it = std::begin(c);
3258 for (
size_t i = 0; i < n; ++i, ++it)
3260 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3266 size_t chunk_size = parallel_detail::effective_parallel_chunk_size(n, pool,
options,
3269 std::vector<std::future<void>> futures;
3274 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3275 size_t chunk_end = std::min(
offset + chunk_size, n);
3278 token =
options.cancel_token]()
3280 auto it = std::begin(c);
3281 std::advance(it, offset);
3282 for (size_t i = offset; i < chunk_end; ++i, ++it)
3284 parallel_detail::throw_if_parallel_canceled(token);
3292 for (
auto & f: futures)
3308 template <
typename Container,
typename Op>
3310 size_t chunk_size = 0)
3314 options.chunk_size = chunk_size;
3318 template <
typename Container,
typename Op>
3322 const size_t n = std::distance(std::begin(c), std::end(c));
3326 auto & pool = parallel_detail::selected_parallel_pool(
options);
3327 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
3329 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3330 auto it = std::begin(c);
3331 for (
size_t i = 0; i < n; ++i, ++it)
3333 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3339 size_t chunk_size = parallel_detail::effective_parallel_chunk_size(n, pool,
options,
3342 auto data_holder = parallel_detail::ensure_random_access(c);
3343 const auto & data = parallel_detail::deref(data_holder);
3345 std::vector<std::future<void>> futures;
3350 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3351 size_t chunk_end = std::min(
offset + chunk_size, n);
3353 futures.push_back(pool.
enqueue([&data, op,
offset, chunk_end,
3354 token =
options.cancel_token]()
3356 auto it = std::begin(data);
3357 std::advance(it, offset);
3358 for (size_t i = offset; i < chunk_end; ++i, ++it)
3360 parallel_detail::throw_if_parallel_canceled(token);
3368 for (
auto & f: futures)
3400 template <
typename Container,
typename Op>
3402 size_t chunk_size = 0)
3406 options.chunk_size = chunk_size;
3410 template <
typename Container,
typename Op>
3414 using T = std::decay_t<
decltype(*std::begin(c))>;
3415 using ResultT = std::invoke_result_t<Op, size_t, const T &>;
3417 const size_t n = std::distance(std::begin(c), std::end(c));
3419 return std::vector<ResultT>{};
3421 auto & pool = parallel_detail::selected_parallel_pool(
options);
3422 if (parallel_detail::use_sequential_parallel_path(n, pool,
options))
3424 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3425 std::vector<ResultT> result;
3427 auto it = std::begin(c);
3428 for (
size_t i = 0; i < n; ++i, ++it)
3430 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3431 result.push_back(op(i, *it));
3436 size_t chunk_size = parallel_detail::effective_parallel_chunk_size(n, pool,
options,
3439 auto data_holder = parallel_detail::ensure_random_access(c);
3440 const auto & data = parallel_detail::deref(data_holder);
3442 std::vector<std::optional<ResultT>> slots(n);
3443 std::vector<std::future<void>> futures;
3448 parallel_detail::throw_if_parallel_canceled(
options.cancel_token);
3449 size_t chunk_end = std::min(
offset + chunk_size, n);
3451 futures.push_back(pool.
enqueue([&slots, &data, op,
offset, chunk_end,
3452 token =
options.cancel_token]()
3454 auto it = std::begin(data);
3455 std::advance(it, offset);
3456 for (size_t i = offset; i < chunk_end; ++i, ++it)
3458 parallel_detail::throw_if_parallel_canceled(token);
3459 slots[i] = op(i, *it);
3465 for (
auto & f: futures)
3468 std::vector<ResultT> result;
3470 for (
auto & slot: slots)
3471 result.push_back(
std::move(*slot));
3492#ifdef AH_PARALLEL_USE_DEFAULT_POOL
3494#define PMAP(c, op) pmaps(parallel_default_pool(), c, op)
3495#define PFILTER(c, pred) pfilter(parallel_default_pool(), c, pred)
3496#define PFOLD(c, init, op) pfoldl(parallel_default_pool(), c, init, op)
3497#define PFOR_EACH(c, op) pfor_each(parallel_default_pool(), c, op)
3498#define PALL(c, pred) pall(parallel_default_pool(), c, pred)
3499#define PEXISTS(c, pred) pexists(parallel_default_pool(), c, pred)
3500#define PSUM(c) psum(parallel_default_pool(), c)
C++20 Ranges support and adaptors for Aleph-w containers.
Read-only cooperative cancellation token.
void throw_if_cancellation_requested() const
Throw operation_canceled if cancellation was requested.
A reusable thread pool for efficient parallel task execution.
size_t num_threads() const noexcept
Get the number of worker threads.
auto enqueue(F &&f, Args &&... args) -> std::future< std::invoke_result_t< F, Args... > >
Submit a task for execution and get a future for the result.
int cmp(const __gmp_expr< T, U > &expr1, const __gmp_expr< V, W > &expr2)
size_t blossom_maximum_cardinality_matching(const GT &g, DynDlist< typename GT::Arc * > &matching, SA sa=SA())
Alias of compute_maximum_cardinality_general_matching().
Freq_Node * pred
Predecessor node in level-order traversal.
const long double offset[]
Offset values indexed by symbol string length (bounded by MAX_OFFSET_INDEX)
decltype(auto) deref(T &&ptr)
Get reference from pointer or unique_ptr.
bool use_sequential_parallel_path(const size_t n, const ThreadPool &pool, const ParallelOptions &options) noexcept
void throw_if_parallel_canceled(const CancellationToken &token)
size_t chunk_count(const size_t n, const size_t chunk_size) noexcept
ThreadPool & selected_parallel_pool(const ParallelOptions &options)
auto ensure_random_access(const Container &c)
For containers with random access, just return a pointer to it For non-random access,...
size_t effective_parallel_chunk_size(const size_t n, const ThreadPool &pool, const ParallelOptions &options, const size_t min_chunk=64)
size_t chunk_size(const size_t n, const size_t num_threads, const size_t min_chunk=64)
Calculate optimal chunk size based on data size and thread count.
constexpr bool has_random_access()
Check if container supports random access.
chunk_bounds bounds_for_chunk(const size_t idx, const size_t n, const size_t chunk_size) noexcept
size_t min_holder_size_impl(const std::tuple< Holders... > &holders, std::index_sequence< Is... >)
Get minimum size from tuple of holders - always O(1) per holder.
void advance_all_iters(std::tuple< Iters... > &iters, std::index_sequence< Is... >)
Advance all iterators in tuple.
auto make_iterators_at(size_t offset, const std::tuple< Holders... > &holders, std::index_sequence< Is... >)
Create tuple of iterators at given offset.
size_t min_holder_size(const std::tuple< Holders... > &holders)
auto deref_all_iters(const std::tuple< Iters... > &iters, std::index_sequence< Is... >)
Dereference all iterators and make tuple.
Main namespace for Aleph-w library functions.
bool pall(ThreadPool &pool, const Container &c, Pred pred, size_t chunk_size=0)
Parallel all predicate (short-circuit).
bool pnone(ThreadPool &pool, const Container &c, Pred pred, size_t chunk_size=0)
Parallel none predicate.
auto pmaps(ThreadPool &pool, const Container &c, Op op, size_t chunk_size=0)
Parallel map operation.
auto pmin(ThreadPool &pool, const Container &c, size_t chunk_size=0)
Parallel minimum element.
T pzip_foldl_n(ThreadPool &pool, T init, Op op, Combiner combiner, const Containers &... cs)
Parallel fold/reduce over N zipped containers (variadic).
ThreadPool & default_pool()
Return the default shared thread pool instance.
T pzip_foldl(ThreadPool &pool, const Container1 &c1, const Container2 &c2, T init, Op op, size_t chunk_size=0)
Parallel zip + fold.
void pzip_for_each(ThreadPool &pool, const Container1 &c1, const Container2 &c2, Op op, size_t chunk_size=0)
Parallel zip + for_each.
size_t size(Node *root) noexcept
ThreadPool & parallel_default_pool()
Global default pool for parallel operations.
auto penumerate_maps(ThreadPool &pool, const Container &c, Op op, size_t chunk_size=0)
Parallel enumerate with map.
void sort_range(Range &r, Cmp cmp)
Sort a whole range in place, portably across the ranges divide.
void penumerate_for_each(ThreadPool &pool, Container &c, Op op, size_t chunk_size=0)
Parallel for_each with index (enumerate).
bool pzip_all_n(ThreadPool &pool, Pred pred, const Containers &... cs)
Parallel all predicate over N zipped containers (variadic).
size_t pcount_if(ThreadPool &pool, const Container &c, Pred pred, size_t chunk_size=0)
Parallel count_if operation.
auto pexclusive_scan(ThreadPool &pool, const Container &c, T init, BinaryOp op, size_t chunk_size=0)
Parallel exclusive scan over a container.
std::optional< size_t > pfind(ThreadPool &pool, const Container &c, Pred pred, size_t chunk_size=0)
Parallel find operation (returns index).
and
Check uniqueness with explicit hash + equality functors.
bool pzip_exists_n(ThreadPool &pool, Pred pred, const Containers &... cs)
Parallel exists predicate over N zipped containers (variadic).
std::decay_t< typename HeadC::Item_Type > T
void psort(ThreadPool &pool, Container &c, Compare cmp=Compare{}, const size_t min_parallel_size=1024)
Parallel sort (in-place).
void pfor_each(ThreadPool &pool, Container &c, Op op, size_t chunk_size=0)
Parallel for_each operation.
auto pfilter(ThreadPool &pool, const Container &c, Pred pred, size_t chunk_size=0)
Parallel filter operation.
auto ppartition(ThreadPool &pool, const Container &c, Pred pred, size_t chunk_size=0)
Parallel partition (stable).
auto pmerge(ThreadPool &pool, const Container1 &c1, const Container2 &c2, Compare comp=Compare{}, size_t chunk_size=0)
Parallel merge of two sorted containers.
T pproduct(ThreadPool &pool, const Container &c, T init=T{1}, size_t chunk_size=0)
Parallel product of elements.
auto pmax(ThreadPool &pool, const Container &c, size_t chunk_size=0)
Parallel maximum element.
void pzip_for_each_n(ThreadPool &pool, Op op, const Containers &... cs)
Parallel for_each over N zipped containers (variadic).
auto pzip_maps_n(ThreadPool &pool, Op op, const Containers &... cs)
Parallel map over N zipped containers (variadic).
auto pscan(ThreadPool &pool, const Container &c, BinaryOp op, size_t chunk_size=0)
Parallel inclusive scan over a container.
T pfoldl(ThreadPool &pool, const Container &c, T init, BinaryOp op, size_t chunk_size=0)
Parallel left fold (reduce).
auto pfind_value(ThreadPool &pool, const Container &c, Pred pred, size_t chunk_size=0)
Parallel find with value return.
auto pzip_maps(ThreadPool &pool, const Container1 &c1, const Container2 &c2, Op op, size_t chunk_size=0)
Parallel zip + map.
T psum(ThreadPool &pool, const Container &c, T init=T{}, size_t chunk_size=0)
Parallel sum of elements.
auto pminmax(ThreadPool &pool, const Container &c, size_t chunk_size=0)
Parallel min and max elements.
size_t pzip_count_if_n(ThreadPool &pool, Pred pred, const Containers &... cs)
Parallel count over N zipped containers (variadic).
bool pexists(ThreadPool &pool, const Container &c, Pred pred, size_t chunk_size=0)
Parallel exists predicate (short-circuit).
Itor::difference_type count(const Itor &beg, const Itor &end, const T &value)
Count elements equal to a value.
static struct argp_option options[]
Common configuration object for parallel algorithms.
ThreadPool * pool
Executor to use (nullptr = default_pool()).
Holder for converted containers (either pointer or unique_ptr to vector).
size_t cached_size
Cached size for O(1) access.
std::decay_t< decltype(*std::begin(std::declval< Container & >()))> value_type
std::conditional_t< parallel_detail::has_random_access< Container >(), const Container *, std::unique_ptr< std::vector< value_type > > > holder_type
ContainerHolder(const Container &c)
decltype(auto) get() const
size_t size() const noexcept
Size is always O(1) - either from random access or from cached vector size.
Filter_Iterator< DynList< int >, DynList< int >::Iterator, Par > It
A modern, efficient thread pool for parallel task execution.