VeloGraphX
High-performance dynamic graph analytics in C++20
Loading...
Searching...
No Matches
numa_scheduler.hpp
Go to the documentation of this file.
1#pragma once
2
5
6#include <algorithm>
7#include <cstddef>
8#include <optional>
9#include <utility>
10#include <vector>
11
12namespace velographx {
13
15 public:
16 NumaLocalScheduler(WorkStealingPool& pool, std::vector<NumaVertexPartition> partitions)
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);
22 }
23 std::sort(node_ids_.begin(), node_ids_.end());
24 }
25
26 std::size_t preferred_queue_for_node(std::size_t node_id,
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();
36 }
37
38 std::size_t preferred_queue_for_vertex(std::size_t vertex,
39 std::size_t sequence = 0) const noexcept {
40 const auto node = numa_node_for_vertex(vertex, partitions_);
41 if (!node.has_value()) return sequence % std::max<std::size_t>(1, pool_.size());
42 return preferred_queue_for_node(*node, sequence);
43 }
44
45 void submit_for_vertex(std::size_t vertex, WorkStealingPool::Task task,
46 std::size_t sequence = 0) {
47 pool_.submit(std::move(task), preferred_queue_for_vertex(vertex, sequence));
48 }
49
50 template <class Fn>
52 std::size_t sequence = 0;
53 for (const auto& partition : partitions_) {
54 const auto queue = partition.node_id.has_value()
55 ? preferred_queue_for_node(*partition.node_id, sequence++)
56 : sequence++ % std::max<std::size_t>(1, pool_.size());
57 pool_.submit([partition, &fn] { fn(partition); }, queue);
58 }
59 pool_.wait_idle();
60 }
61
62 const std::vector<NumaVertexPartition>& partitions() const noexcept { return partitions_; }
63
64 private:
65 WorkStealingPool& pool_;
66 std::vector<NumaVertexPartition> partitions_;
67 std::vector<std::size_t> node_ids_;
68};
69
70} // namespace velographx
void submit_for_vertex(std::size_t vertex, WorkStealingPool::Task task, std::size_t sequence=0)
std::size_t preferred_queue_for_vertex(std::size_t vertex, std::size_t sequence=0) const noexcept
std::size_t preferred_queue_for_node(std::size_t node_id, std::size_t sequence=0) const noexcept
const std::vector< NumaVertexPartition > & partitions() const noexcept
NumaLocalScheduler(WorkStealingPool &pool, std::vector< NumaVertexPartition > partitions)
std::size_t size() const noexcept
void submit(Task task, std::size_t locality_hint=0)
std::optional< std::size_t > numa_node_for_vertex(std::size_t vertex, const std::vector< NumaVertexPartition > &partitions) noexcept