-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathDijkstraParallelSolver.cpp
More file actions
154 lines (123 loc) · 5.15 KB
/
Copy pathDijkstraParallelSolver.cpp
File metadata and controls
154 lines (123 loc) · 5.15 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
#include "DijkstraParallelSolver.h"
/**
* @brief Konstruktor, który inicjalizuje cały stan współdzielony.
*/
DijkstraParallelSolver::DijkstraParallelSolver(const AdjacencyList& g, int n_v, int n_t, int start)
: adj(g),
num_vertices(n_v),
num_threads(n_t),
start_node(start),
distances(n_v, INFINITY_COST),
local_pqs(n_t),
barrier(n_t),
processed(n_v),
algorithm_done(false),
current_global_node({INFINITY_COST, -1})
{
for (int i = 0; i < n_v; ++i) processed[i].store(false);
distances[start_node] = 0.0;
int start_owner = get_owner(start_node);
local_pqs[start_owner].insert({0.0, start_node});
}
/** Zwraca ID wątku, który jest "właścicielem" wierzchołka v */
int DijkstraParallelSolver::get_owner(int v) const {
return v % num_threads;
}
/**
* @brief Uruchamia algorytm i zwraca wektor kosztów.
*/
std::vector<double> DijkstraParallelSolver::solve() {
std::vector<std::thread> threads;
for (int i = 0; i < num_threads; ++i) {
threads.emplace_back(&DijkstraParallelSolver::worker_thread, this, i);
}
for (auto& t : threads) {
t.join();
}
return distances;
}
/**
* @brief Główna funkcja robocza (Wersja 5 - Poprawiona logika relaksacji)
*/
void DijkstraParallelSolver::worker_thread(int thread_id) {
while (true) {
// RÓWNOLEGŁA RELAKSACJA
int u = current_global_node.second;
double u_dist = current_global_node.first;
if (u != -1) {
for (const auto& edge : adj[u]) {
int v = edge.first;
int v_owner = get_owner(v);
if (v_owner == thread_id) {
double weight = edge.second;
double new_dist = u_dist + weight;
// Ten odczyt jest bezpieczny (atomowy)
if (processed[v].load(std::memory_order_acquire)) continue;
// Ten odczyt/zapis distances[v] jest technicznie wyścigiem,
// ale jest to łagodny wyścig.
// W najgorszym wypadku dodamy do kolejki wpis, który
// i tak zostanie przefiltrowany.
if (new_dist < distances[v]) {
distances[v] = new_dist;
local_pqs[v_owner].insert({new_dist, v});
}
}
}
}
//SYNCHRONIZACJA I REDUKCJA
barrier.wait(); // Czekaj na zakończenie relaksacji
// Wątek 0 wykonuje fazę redukcji
if (thread_id == 0) {
std::pair<double, int> next_global_min = {INFINITY_COST, -1};
bool found_unvisited = false;
while (!found_unvisited) {
next_global_min = {INFINITY_COST, -1};
int min_owner = -1;
// 1. Znajdź globalne minimum (bez blokad)
for (int i = 0; i < num_threads; ++i) {
if (!local_pqs[i].empty()) {
auto top = *local_pqs[i].begin();
//DODATKOWE ZABEZPIECZENIE
// Sprawdź, czy wpis na szczycie nie jest przestarzały
// (tzn. czy jego koszt jest gorszy niż już znany dystans)
// To pomaga w czyszczeniu kolejki.
if (top.first > distances[top.second]) {
local_pqs[i].erase(local_pqs[i].begin());
i--; // Sprawdź tę samą kolejkę ponownie
continue;
}
//KONIEC ZABEZPIECZENIA
if (top.first < next_global_min.first) {
next_global_min = top;
min_owner = i;
}
}
}
// 2. Sprawdź warunek zakończenia
if (min_owner == -1) {
algorithm_done.store(true);
found_unvisited = true;
} else {
int u_next = next_global_min.second;
// 3. Sprawdź, czy nie jest *przetworzony* (odwiedzony)
if (processed[u_next].load(std::memory_order_acquire)) {
local_pqs[min_owner].erase(local_pqs[min_owner].begin());
// i pętla while() szuka dalej
} else {
// 4. Znaleziono poprawny, nowy węzeł
found_unvisited = true;
processed[u_next].store(true, std::memory_order_release);
local_pqs[min_owner].erase(local_pqs[min_owner].begin());
}
}
} // koniec while (!found_unvisited)
// 5. Rozgłoś nowy węzeł
current_global_node = next_global_min;
} // koniec if (thread_id == 0)
//OCZEKIWANIE NA ZAKOŃCZENIE REDUKCJI I SPRAWDZENIE KOŃCA
barrier.wait(); // Czekaj na zakończenie redukcji przez wątek 0
if (algorithm_done.load()) {
break;
}
} // koniec while(true)
}