Skip to content

Commit 8635ba2

Browse files
Changed multiple files
1 parent 5266ad1 commit 8635ba2

4 files changed

Lines changed: 127 additions & 87 deletions

File tree

src/query/algorithms/triangles/SheepTriangles.cpp

Lines changed: 26 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -22,8 +22,8 @@ Logger sheep_triangle_logger;
2222
long SheepTriangles::run(JasmineGraphHashMapLocalStore &graphDB,
2323
JasmineGraphHashMapCentralStore &centralStore,
2424
JasmineGraphHashMapDuplicateCentralStore &duplicateCentralStore,
25-
std::string graphId,
26-
std::string partitionId) {
25+
const std::string &graphId,
26+
const std::string &partitionId) {
2727
OTEL_TRACE_FUNCTION();
2828

2929
std::string workerInfo = "worker_" + graphId + "_partition_" + partitionId;
@@ -71,23 +71,38 @@ long SheepTriangles::run(JasmineGraphHashMapLocalStore &graphDB,
7171
void SheepTriangles::mergeStores(
7272
std::map<long, std::unordered_set<long>> &localMap,
7373
std::map<long, std::unordered_set<long>> &centralMap,
74-
std::map<long, std::unordered_set<long>> &duplicateMap) {
74+
const std::map<long, std::unordered_set<long>> &duplicateMap) {
7575

7676
// Merge duplicate central into central (avoiding duplicates)
77-
for (const auto &entry : duplicateMap) {
78-
long vertex = entry.first;
79-
const auto &edges = entry.second;
77+
for (const auto &[vertex, edges] : duplicateMap) {
8078
centralMap[vertex].insert(edges.begin(), edges.end());
8179
}
8280

8381
// Merge central into local
84-
for (const auto &entry : centralMap) {
85-
long vertex = entry.first;
86-
const auto &edges = entry.second;
82+
for (const auto &[vertex, edges] : centralMap) {
8783
localMap[vertex].insert(edges.begin(), edges.end());
8884
}
8985
}
9086

87+
static long countTrianglesForPair(long u, long v, const std::unordered_set<long> &u_neighbors,
88+
const std::unordered_set<long> &v_neighbors, bool returnTriangles,
89+
std::ostringstream &triangleStream) {
90+
long count = 0;
91+
for (long w : u_neighbors) {
92+
if (w <= v) continue; // Enforce v < w
93+
94+
// Check if v is connected to w (completing the triangle)
95+
if (v_neighbors.find(w) != v_neighbors.end()) {
96+
count++;
97+
98+
if (returnTriangles) {
99+
triangleStream << u << "," << v << "," << w << ":";
100+
}
101+
}
102+
}
103+
return count;
104+
}
105+
91106
SheepTriangleResult SheepTriangles::countTriangles(
92107
std::map<long, std::unordered_set<long>> &edgeMap,
93108
bool returnTriangles) {
@@ -111,20 +126,8 @@ SheepTriangleResult SheepTriangles::countTriangles(
111126
if (v_it == edgeMap.end()) continue;
112127
const auto &v_neighbors = v_it->second;
113128

114-
// For each neighbor w of u where w > v (enforce u < v < w)
115-
for (long w : u_neighbors) {
116-
if (w <= v) continue; // Enforce v < w
117-
118-
// Check if v is connected to w (completing the triangle)
119-
if (v_neighbors.find(w) != v_neighbors.end()) {
120-
// Found triangle (u, v, w) with u < v < w
121-
result.count++;
122-
123-
if (returnTriangles) {
124-
triangleStream << u << "," << v << "," << w << ":";
125-
}
126-
}
127-
}
129+
result.count += countTrianglesForPair(u, v, u_neighbors, v_neighbors,
130+
returnTriangles, triangleStream);
128131
}
129132
}
130133

src/query/algorithms/triangles/SheepTriangles.h

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -48,8 +48,8 @@ class SheepTriangles {
4848
static long run(JasmineGraphHashMapLocalStore &graphDB,
4949
JasmineGraphHashMapCentralStore &centralStore,
5050
JasmineGraphHashMapDuplicateCentralStore &duplicateCentralStore,
51-
std::string graphId,
52-
std::string partitionId);
51+
const std::string &graphId,
52+
const std::string &partitionId);
5353

5454
private:
5555
/**
@@ -66,7 +66,7 @@ class SheepTriangles {
6666
static void mergeStores(
6767
std::map<long, std::unordered_set<long>> &localMap,
6868
std::map<long, std::unordered_set<long>> &centralMap,
69-
std::map<long, std::unordered_set<long>> &duplicateMap);
69+
const std::map<long, std::unordered_set<long>> &duplicateMap);
7070
};
7171

7272
#endif // JASMINEGRAPH_SHEEPTRIANGLES_H

src/server/JasmineGraphInstanceService.cpp

Lines changed: 83 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -73,31 +73,44 @@ static void delete_graph_command(int connFd, bool *loop_exit_p);
7373
static void delete_graph_fragment_command(int connFd, bool *loop_exit_p);
7474
static void duplicate_centralstore_command(int connFd, int serverPort, bool *loop_exit_p);
7575
static void worker_in_degree_distribution_command(
76-
int connFd, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
77-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores, bool *loop_exit_p);
76+
int connFd,
77+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
78+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
79+
bool *loop_exit_p);
7880
static void in_degree_distribution_command(
79-
int connFd, int serverPort, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
80-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores, bool *loop_exit_p);
81+
int connFd, int serverPort,
82+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
83+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
84+
bool *loop_exit_p);
8185
static void worker_out_degree_distribution_command(
82-
int connFd, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
83-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores, bool *loop_exit_p);
86+
int connFd,
87+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
88+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
89+
bool *loop_exit_p);
8490
static void out_degree_distribution_command(
85-
int connFd, int serverPort, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
86-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores, bool *loop_exit_p);
87-
static void page_rank_command(int connFd, int serverPort,
88-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
89-
bool *loop_exit_p);
91+
int connFd, int serverPort,
92+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
93+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
94+
bool *loop_exit_p);
95+
static void page_rank_command(
96+
int connFd, int serverPort,
97+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
98+
bool *loop_exit_p);
9099
static void worker_page_rank_distribution_command(
91-
int connFd, int serverPort, std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
100+
int connFd, int serverPort,
101+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
102+
bool *loop_exit_p);
103+
static void egonet_command(
104+
int connFd, int serverPort,
105+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
106+
bool *loop_exit_p);
107+
static void worker_egonet_command(
108+
int connFd, int serverPort,
109+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
92110
bool *loop_exit_p);
93-
static void egonet_command(int connFd, int serverPort,
94-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
95-
bool *loop_exit_p);
96-
static void worker_egonet_command(int connFd, int serverPort,
97-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
98-
bool *loop_exit_p);
99111
static void triangles_command(
100-
int connFd, int serverPort, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
112+
int connFd, int serverPort,
113+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
101114
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
102115
std::map<std::string, JasmineGraphHashMapDuplicateCentralStore, std::less<>> &graphDBMapDuplicateCentralStores,
103116
bool *loop_exit_p);
@@ -135,10 +148,11 @@ static void graph_stream_start_command(int connFd, InstanceStreamHandler &instan
135148
static void send_priority_command(int connFd, bool *loop_exit_p);
136149
static std::string initiate_command_common(int connFd, bool *loop_exit_p);
137150
static void batch_upload_common(int connFd, bool *loop_exit_p, bool batch_upload);
138-
static void degree_distribution_common(int connFd, int serverPort,
139-
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
140-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
141-
bool *loop_exit_p, bool in);
151+
static void degree_distribution_common(
152+
int connFd, int serverPort,
153+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
154+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
155+
bool *loop_exit_p, bool in);
142156
static void push_partition_command(int connFd, bool *loop_exit_p);
143157
static void push_file_command(int connFd, bool *loop_exit_p);
144158
static void send_edges_command(int connFd, bool *loop_exit_p);
@@ -183,11 +197,12 @@ void *instanceservicesession(void *dummyPt) {
183197
instanceservicesessionargs sessionargs = *sessionargs_p;
184198
delete sessionargs_p;
185199
int connFd = sessionargs.connFd;
186-
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> *graphDBMapLocalStores = sessionargs.graphDBMapLocalStores;
187-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> *graphDBMapCentralStores =
188-
sessionargs.graphDBMapCentralStores;
189-
std::map<std::string, JasmineGraphHashMapDuplicateCentralStore, std::less<>> *graphDBMapDuplicateCentralStores =
190-
sessionargs.graphDBMapDuplicateCentralStores;
200+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>>
201+
*graphDBMapLocalStores = sessionargs.graphDBMapLocalStores;
202+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>>
203+
*graphDBMapCentralStores = sessionargs.graphDBMapCentralStores;
204+
std::map<std::string, JasmineGraphHashMapDuplicateCentralStore, std::less<>>
205+
*graphDBMapDuplicateCentralStores = sessionargs.graphDBMapDuplicateCentralStores;
191206
std::map<std::string, JasmineGraphIncrementalLocalStore *> &incrementalLocalStoreMap =
192207
*(sessionargs.incrementalLocalStore);
193208
InstanceStreamHandler streamHandler(incrementalLocalStoreMap);
@@ -1281,10 +1296,11 @@ bool JasmineGraphInstanceService::duplicateCentralStore(int thisWorkerPort, int
12811296
return true;
12821297
}
12831298

1284-
map<long, long> calculateOutDegreeDist(string graphID, string partitionID, int serverPort,
1285-
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
1286-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
1287-
std::vector<string> &workerSockets) {
1299+
map<long, long> calculateOutDegreeDist(
1300+
string graphID, string partitionID, int serverPort,
1301+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
1302+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
1303+
std::vector<string> &workerSockets) {
12881304
map<long, long> degreeDistribution =
12891305
calculateLocalOutDegreeDist(graphID, partitionID, graphDBMapLocalStores, graphDBMapCentralStores);
12901306

@@ -1305,7 +1321,8 @@ map<long, long> calculateOutDegreeDist(string graphID, string partitionID, int s
13051321
}
13061322

13071323
map<long, long> calculateLocalOutDegreeDist(
1308-
string graphID, string partitionID, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
1324+
string graphID, string partitionID,
1325+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
13091326
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores) {
13101327
auto t_start = std::chrono::high_resolution_clock::now();
13111328

@@ -1352,8 +1369,9 @@ map<long, long> calculateLocalOutDegreeDist(
13521369
}
13531370

13541371
map<long, long> calculateLocalInDegreeDist(
1355-
string graphID, string partitionID, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
1356-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores) {
1372+
string graphID, string partitionID,
1373+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
1374+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>>) {
13571375
JasmineGraphHashMapLocalStore graphDB;
13581376

13591377
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>>::iterator it;
@@ -1370,10 +1388,11 @@ map<long, long> calculateLocalInDegreeDist(
13701388
return degreeDistribution;
13711389
}
13721390

1373-
map<long, long> calculateInDegreeDist(string graphID, string partitionID, int serverPort,
1374-
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
1375-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
1376-
std::vector<string> &workerSockets, string workerList) {
1391+
map<long, long> calculateInDegreeDist(
1392+
string graphID, string partitionID, int serverPort,
1393+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
1394+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
1395+
std::vector<string> &workerSockets, string workerList) {
13771396
auto t_start = std::chrono::high_resolution_clock::now();
13781397

13791398
map<long, long> degreeDistribution =
@@ -2443,10 +2462,11 @@ static void worker_in_degree_distribution_command(
24432462
*loop_exit_p = true;
24442463
}
24452464

2446-
static void degree_distribution_common(int connFd, int serverPort,
2447-
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
2448-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
2449-
bool *loop_exit_p, bool in) {
2465+
static void degree_distribution_common(
2466+
int connFd, int serverPort,
2467+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
2468+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
2469+
bool *loop_exit_p, bool in) {
24502470
if (!Utils::send_str_wrapper(connFd, JasmineGraphInstanceProtocol::OK)) {
24512471
*loop_exit_p = true;
24522472
return;
@@ -2496,14 +2516,18 @@ static void degree_distribution_common(int connFd, int serverPort,
24962516
}
24972517

24982518
static void in_degree_distribution_command(
2499-
int connFd, int serverPort, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
2500-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores, bool *loop_exit_p) {
2519+
int connFd, int serverPort,
2520+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
2521+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
2522+
bool *loop_exit_p) {
25012523
degree_distribution_common(connFd, serverPort, graphDBMapLocalStores, graphDBMapCentralStores, loop_exit_p, true);
25022524
}
25032525

25042526
static void worker_out_degree_distribution_command(
2505-
int connFd, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
2506-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores, bool *loop_exit_p) {
2527+
int connFd,
2528+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
2529+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
2530+
bool *loop_exit_p) {
25072531
if (!Utils::send_str_wrapper(connFd, JasmineGraphInstanceProtocol::OK)) {
25082532
*loop_exit_p = true;
25092533
return;
@@ -2538,8 +2562,10 @@ static void worker_out_degree_distribution_command(
25382562
}
25392563

25402564
static void out_degree_distribution_command(
2541-
int connFd, int serverPort, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
2542-
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores, bool *loop_exit_p) {
2565+
int connFd, int serverPort,
2566+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
2567+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
2568+
bool *loop_exit_p) {
25432569
degree_distribution_common(connFd, serverPort, graphDBMapLocalStores, graphDBMapCentralStores, loop_exit_p, false);
25442570
}
25452571

@@ -2697,7 +2723,8 @@ static void page_rank_command(int connFd, int serverPort,
26972723
}
26982724

26992725
static void worker_page_rank_distribution_command(
2700-
int connFd, int serverPort, std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
2726+
int connFd, int serverPort,
2727+
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
27012728
bool *loop_exit_p) {
27022729
if (!Utils::send_str_wrapper(connFd, JasmineGraphInstanceProtocol::OK)) {
27032730
*loop_exit_p = true;
@@ -2956,7 +2983,8 @@ static void worker_egonet_command(int connFd, int serverPort, std::map<std::stri
29562983

29572984
template <typename TriangleCountFn>
29582985
static void triangle_counting_command_common(
2959-
int connFd, int serverPort, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
2986+
int connFd, int serverPort,
2987+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
29602988
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
29612989
std::map<std::string, JasmineGraphHashMapDuplicateCentralStore, std::less<>> &graphDBMapDuplicateCentralStores,
29622990
bool *loop_exit_p, TriangleCountFn countFn) {
@@ -3033,7 +3061,8 @@ static void triangle_counting_command_common(
30333061
}
30343062

30353063
static void triangles_command(
3036-
int connFd, int serverPort, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
3064+
int connFd, int serverPort,
3065+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
30373066
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
30383067
std::map<std::string, JasmineGraphHashMapDuplicateCentralStore, std::less<>> &graphDBMapDuplicateCentralStores,
30393068
bool *loop_exit_p) {
@@ -3042,7 +3071,8 @@ static void triangles_command(
30423071
}
30433072

30443073
static void sheep_triangles_command(
3045-
int connFd, int serverPort, std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
3074+
int connFd, int serverPort,
3075+
std::map<std::string, JasmineGraphHashMapLocalStore, std::less<>> &graphDBMapLocalStores,
30463076
std::map<std::string, JasmineGraphHashMapCentralStore, std::less<>> &graphDBMapCentralStores,
30473077
std::map<std::string, JasmineGraphHashMapDuplicateCentralStore, std::less<>> &graphDBMapDuplicateCentralStores,
30483078
bool *loop_exit_p) {

0 commit comments

Comments
 (0)