17 : pool_(pool), partitions_(std::move(
partitions)) {
18 for (
const auto& partition : partitions_) {
19 if (!partition.node_id.has_value())
continue;
20 if (std::find(node_ids_.begin(), node_ids_.end(), *partition.node_id) == node_ids_.end())
21 node_ids_.push_back(*partition.node_id);
23 std::sort(node_ids_.begin(), node_ids_.end());
27 std::size_t sequence = 0) const noexcept {
28 if (pool_.
size() == 0 || node_ids_.empty())
return sequence % std::max<std::size_t>(1, pool_.
size());
29 const auto it = std::find(node_ids_.begin(), node_ids_.end(), node_id);
30 const std::size_t node_rank = it == node_ids_.end()
31 ? node_id % node_ids_.size()
32 :
static_cast<std::size_t
>(it - node_ids_.begin());
33 const std::size_t node_count = node_ids_.size();
34 const std::size_t local_workers = (pool_.
size() + node_count - 1) / node_count;
35 return (node_rank + (sequence % local_workers) * node_count) % pool_.
size();
52 std::size_t sequence = 0;
53 for (
const auto& partition : partitions_) {
54 const auto queue = partition.node_id.has_value()
56 : sequence++ % std::max<std::size_t>(1, pool_.
size());
57 pool_.
submit([partition, &fn] { fn(partition); }, queue);