Skip to content

Commit 3cb9e59

Browse files
authored
refactor: Avoid unnecessary compact for csr before dump (#738)
Fix #736 Fix #701
1 parent 29d7c29 commit 3cb9e59

9 files changed

Lines changed: 217 additions & 2 deletions

File tree

include/neug/storages/csr/csr_base.h

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,8 @@
2929

3030
namespace neug {
3131

32+
class EdgeTable;
33+
3234
// CsrType is defined in csr_view.h to avoid a csr_base <-> csr_view
3335
// include cycle (CsrView may need to know about CsrType in future).
3436

@@ -93,6 +95,13 @@ class CsrBase : public Module {
9395
/// (e.g. PropertyGraphCowState) is responsible for tracking which
9496
/// adjlists have been detached.
9597
virtual void DetachVertex(vid_t vid, Allocator& alloc) = 0;
98+
99+
private:
100+
friend class EdgeTable;
101+
102+
// Called after a timestamped edge update so compact() cannot take the
103+
// "nothing to do" fast path.
104+
virtual void mark_compaction_required() = 0;
96105
};
97106

98107
template <typename EDATA_T>

include/neug/storages/csr/immutable_csr.h

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -116,6 +116,9 @@ class ImmutableCsr : public TypedCsrBase<EDATA_T> {
116116
cow_clone->nbr_list_buffer_ = nbr_list_buffer_;
117117
cow_clone->unsorted_since_ = unsorted_since_;
118118
cow_clone->edge_num_ = edge_num_.load();
119+
cow_clone->needs_compact_.store(
120+
needs_compact_.load(std::memory_order_relaxed),
121+
std::memory_order_relaxed);
119122
return cow_clone;
120123
}
121124

@@ -131,11 +134,16 @@ class ImmutableCsr : public TypedCsrBase<EDATA_T> {
131134
}
132135

133136
private:
137+
void mark_compaction_required() override {
138+
needs_compact_.store(true, std::memory_order_relaxed);
139+
}
140+
134141
std::shared_ptr<IDataContainer> adj_list_buffer_;
135142
std::shared_ptr<IDataContainer> degree_list_buffer_;
136143
std::shared_ptr<IDataContainer> nbr_list_buffer_;
137144
timestamp_t unsorted_since_;
138145
std::atomic<uint64_t> edge_num_{0};
146+
std::atomic<bool> needs_compact_{false};
139147
CsrPrefetchPolicy prefetch_policy_;
140148

141149
void refresh_prefetch_policy();
@@ -239,6 +247,8 @@ class SingleImmutableCsr : public TypedCsrBase<EDATA_T> {
239247
}
240248

241249
private:
250+
void mark_compaction_required() override {}
251+
242252
void refresh_prefetch_policy();
243253
CsrPrefetchPolicy prefetch_policy_;
244254
std::shared_ptr<IDataContainer> nbr_list_buffer_;

include/neug/storages/csr/mutable_csr.h

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -142,6 +142,9 @@ class MutableCsr : public TypedCsrBase<EDATA_T> {
142142
nbr.data = data;
143143
nbr.timestamp.store(ts);
144144
edge_num_.fetch_add(1);
145+
if (ts != 0) {
146+
needs_compact_.store(true, std::memory_order_relaxed);
147+
}
145148
// invalidate sort flag
146149
if (ts < unsorted_since_) {
147150
unsorted_since_ = 0;
@@ -217,6 +220,9 @@ class MutableCsr : public TypedCsrBase<EDATA_T> {
217220
cow_clone->nbr_list_ = nbr_list_;
218221
cow_clone->unsorted_since_ = unsorted_since_;
219222
cow_clone->edge_num_ = edge_num_.load();
223+
cow_clone->needs_compact_.store(
224+
needs_compact_.load(std::memory_order_relaxed),
225+
std::memory_order_relaxed);
220226
return cow_clone;
221227
}
222228

@@ -235,13 +241,18 @@ class MutableCsr : public TypedCsrBase<EDATA_T> {
235241
}
236242

237243
private:
244+
void mark_compaction_required() override {
245+
needs_compact_.store(true, std::memory_order_relaxed);
246+
}
247+
238248
std::unique_ptr<SpinLock[]> locks_;
239249
std::shared_ptr<IDataContainer> adj_list_buffer_;
240250
std::shared_ptr<IDataContainer> degree_list_;
241251
std::shared_ptr<IDataContainer> cap_list_;
242252
std::shared_ptr<IDataContainer> nbr_list_;
243253
timestamp_t unsorted_since_;
244254
std::atomic<uint64_t> edge_num_{0};
255+
std::atomic<bool> needs_compact_{false};
245256
CsrPrefetchPolicy prefetch_policy_;
246257

247258
void refresh_prefetch_policy();
@@ -331,6 +342,9 @@ class SingleMutableCsr : public TypedCsrBase<EDATA_T> {
331342
CHECK_EQ(nbrs[src].timestamp, std::numeric_limits<timestamp_t>::max());
332343
nbrs[src].timestamp.store(ts);
333344
edge_num_.fetch_add(1, std::memory_order_relaxed);
345+
if (ts != 0) {
346+
needs_compact_.store(true, std::memory_order_relaxed);
347+
}
334348
return {0, static_cast<const void*>(&nbrs[src].data)};
335349
}
336350

@@ -367,6 +381,9 @@ class SingleMutableCsr : public TypedCsrBase<EDATA_T> {
367381
auto cow_clone = std::make_unique<SingleMutableCsr<EDATA_T>>();
368382
cow_clone->nbr_list_ = nbr_list_;
369383
cow_clone->edge_num_ = edge_num_.load();
384+
cow_clone->needs_compact_.store(
385+
needs_compact_.load(std::memory_order_relaxed),
386+
std::memory_order_relaxed);
370387
return cow_clone;
371388
}
372389

@@ -382,8 +399,13 @@ class SingleMutableCsr : public TypedCsrBase<EDATA_T> {
382399
}
383400

384401
private:
402+
void mark_compaction_required() override {
403+
needs_compact_.store(true, std::memory_order_relaxed);
404+
}
405+
385406
std::shared_ptr<IDataContainer> nbr_list_;
386407
std::atomic<uint64_t> edge_num_{0};
408+
std::atomic<bool> needs_compact_{false};
387409
CsrPrefetchPolicy prefetch_policy_;
388410

389411
void refresh_prefetch_policy();
@@ -478,6 +500,9 @@ class EmptyCsr : public TypedCsrBase<EDATA_T> {
478500
static std::string type_name() {
479501
return "empty_csr<" + type_name_string<EDATA_T>() + ">";
480502
}
503+
504+
private:
505+
void mark_compaction_required() override {}
481506
};
482507

483508
} // namespace neug

src/storages/csr/immutable_csr.cc

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ void ImmutableCsr<EDATA_T>::Open(Checkpoint& ckp, const ModuleDescriptor& desc,
4242
MemoryLevel memory_level) {
4343
unsorted_since_ = std::stoull(desc.get("unsorted_since").value_or("0"));
4444
edge_num_.store(std::stoull(desc.get("edge_num").value_or("0")));
45+
needs_compact_.store(false, std::memory_order_relaxed);
4546
degree_list_buffer_ = std::shared_ptr<IDataContainer>(ckp.OpenFile(
4647
desc.get_path(ModuleDescriptor::kDegreeListPath).value_or(""),
4748
memory_level));
@@ -90,6 +91,9 @@ void ImmutableCsr<EDATA_T>::Dump(Checkpoint& ckp, CheckpointManifest& meta,
9091

9192
template <typename EDATA_T>
9293
void ImmutableCsr<EDATA_T>::compact() {
94+
if (!needs_compact_.load(std::memory_order_relaxed)) {
95+
return;
96+
}
9397
// For current adj_list where the dst vertex is invalid, swap it to the end.
9498
vid_t vnum = size();
9599
if (vnum <= 0) {
@@ -126,6 +130,7 @@ void ImmutableCsr<EDATA_T>::compact() {
126130
adj_arr[i] = ptr;
127131
ptr += deg_arr[i];
128132
}
133+
needs_compact_.store(false, std::memory_order_relaxed);
129134
}
130135

131136
template <typename EDATA_T>
@@ -183,6 +188,9 @@ void ImmutableCsr<EDATA_T>::batch_sort_by_edge_data(timestamp_t ts) {
183188
template <typename EDATA_T>
184189
void ImmutableCsr<EDATA_T>::batch_delete_vertices(
185190
const std::set<vid_t>& src_set, const std::set<vid_t>& dst_set) {
191+
if (!src_set.empty() || !dst_set.empty()) {
192+
needs_compact_.store(true, std::memory_order_relaxed);
193+
}
186194
vid_t vnum = size();
187195
auto** adj_arr = reinterpret_cast<nbr_t**>(adj_list_buffer_->GetData());
188196
auto* deg_arr = reinterpret_cast<int*>(degree_list_buffer_->GetData());
@@ -233,6 +241,9 @@ void ImmutableCsr<EDATA_T>::batch_delete_vertices(
233241
template <typename EDATA_T>
234242
void ImmutableCsr<EDATA_T>::batch_delete_edges(
235243
const std::vector<vid_t>& src_list, const std::vector<vid_t>& dst_list) {
244+
if (!src_list.empty()) {
245+
needs_compact_.store(true, std::memory_order_relaxed);
246+
}
236247
std::map<vid_t, std::set<vid_t>> src_dst_map;
237248
vid_t vnum = size();
238249
for (size_t i = 0; i < src_list.size(); ++i) {
@@ -271,6 +282,9 @@ void ImmutableCsr<EDATA_T>::batch_delete_edges(
271282
template <typename EDATA_T>
272283
void ImmutableCsr<EDATA_T>::batch_delete_edges(
273284
const std::vector<std::pair<vid_t, int32_t>>& edges) {
285+
if (!edges.empty()) {
286+
needs_compact_.store(true, std::memory_order_relaxed);
287+
}
274288
std::map<vid_t, std::set<int32_t>> src_offset_map;
275289
vid_t vnum = size();
276290
auto** adj_arr = reinterpret_cast<nbr_t**>(adj_list_buffer_->GetData());
@@ -315,6 +329,7 @@ void ImmutableCsr<EDATA_T>::delete_edge(vid_t src, int32_t offset,
315329
}
316330
nbrs[offset].neighbor = std::numeric_limits<vid_t>::max();
317331
edge_num_.fetch_sub(1, std::memory_order_relaxed);
332+
needs_compact_.store(true, std::memory_order_relaxed);
318333
unsorted_since_ = 0;
319334
}
320335

src/storages/csr/mutable_csr.cc

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,7 @@ void MutableCsr<EDATA_T>::Open(Checkpoint& ckp,
8585
" but degree list implies " + std::to_string(edge_count) +
8686
", desc: " + descriptor.ToJsonString());
8787
}
88+
needs_compact_.store(false, std::memory_order_relaxed);
8889
refresh_prefetch_policy();
8990
}
9091

@@ -171,6 +172,9 @@ void MutableCsr<EDATA_T>::Dump(Checkpoint& ckp, CheckpointManifest& meta,
171172

172173
template <typename EDATA_T>
173174
void MutableCsr<EDATA_T>::compact() {
175+
if (!needs_compact_.load(std::memory_order_relaxed)) {
176+
return;
177+
}
174178
// Remove deleted edges and reset timestamps on surviving edges.
175179
size_t vnum = vertex_capacity();
176180
auto** buf_arr = reinterpret_cast<nbr_t**>(adj_list_buffer_->GetData());
@@ -209,6 +213,7 @@ void MutableCsr<EDATA_T>::compact() {
209213
std::to_string(edge_num_.load()) + ", actual " +
210214
std::to_string(total_edge_num));
211215
}
216+
needs_compact_.store(false, std::memory_order_relaxed);
212217
}
213218

214219
template <typename EDATA_T>
@@ -278,6 +283,9 @@ void MutableCsr<EDATA_T>::batch_sort_by_edge_data(timestamp_t ts) {
278283
template <typename EDATA_T>
279284
void MutableCsr<EDATA_T>::batch_delete_vertices(
280285
const std::set<vid_t>& src_set, const std::set<vid_t>& dst_set) {
286+
if (!src_set.empty() || !dst_set.empty()) {
287+
needs_compact_.store(true, std::memory_order_relaxed);
288+
}
281289
vid_t vnum = static_cast<vid_t>(vertex_capacity());
282290
auto** buf_arr = reinterpret_cast<nbr_t**>(adj_list_buffer_->GetData());
283291
auto* sz_arr = reinterpret_cast<std::atomic<int>*>(degree_list_->GetData());
@@ -333,6 +341,9 @@ void MutableCsr<EDATA_T>::batch_delete_vertices(
333341
template <typename EDATA_T>
334342
void MutableCsr<EDATA_T>::batch_delete_edges(
335343
const std::vector<vid_t>& src_list, const std::vector<vid_t>& dst_list) {
344+
if (!src_list.empty()) {
345+
needs_compact_.store(true, std::memory_order_relaxed);
346+
}
336347
std::map<vid_t, std::set<vid_t>> src_dst_map;
337348
vid_t vnum = static_cast<vid_t>(vertex_capacity());
338349
for (size_t i = 0; i < src_list.size(); ++i) {
@@ -367,6 +378,9 @@ void MutableCsr<EDATA_T>::batch_delete_edges(
367378
template <typename EDATA_T>
368379
void MutableCsr<EDATA_T>::batch_delete_edges(
369380
const std::vector<std::pair<vid_t, int32_t>>& edges) {
381+
if (!edges.empty()) {
382+
needs_compact_.store(true, std::memory_order_relaxed);
383+
}
370384
std::map<vid_t, std::set<int32_t>> src_offset_map;
371385
vid_t vnum = static_cast<vid_t>(vertex_capacity());
372386
auto* sz_arr = reinterpret_cast<std::atomic<int>*>(degree_list_->GetData());
@@ -413,6 +427,7 @@ void MutableCsr<EDATA_T>::delete_edge(vid_t src, int32_t offset,
413427
if (old_ts <= ts) {
414428
nbrs[offset].timestamp.store(std::numeric_limits<timestamp_t>::max());
415429
edge_num_.fetch_sub(1, std::memory_order_relaxed);
430+
needs_compact_.store(true, std::memory_order_relaxed);
416431
unsorted_since_ = 0;
417432
} else if (old_ts == std::numeric_limits<timestamp_t>::max()) {
418433
LOG(ERROR) << "Attempting to delete already deleted edge.";
@@ -442,6 +457,9 @@ void MutableCsr<EDATA_T>::revert_delete_edge(vid_t src, vid_t nbr,
442457
assert(nbrs[offset].neighbor == nbr);
443458
nbrs[offset].timestamp.store(ts);
444459
edge_num_.fetch_add(1, std::memory_order_relaxed);
460+
if (ts != 0) {
461+
needs_compact_.store(true, std::memory_order_relaxed);
462+
}
445463
} else {
446464
THROW_INVALID_ARGUMENT_EXCEPTION(
447465
"Attempting to revert delete on edge that is not deleted.");
@@ -519,6 +537,9 @@ void MutableCsr<EDATA_T>::batch_put_edges(const std::vector<vid_t>& src_list,
519537
added_edge_num++;
520538
}
521539
edge_num_.fetch_add(added_edge_num, std::memory_order_relaxed);
540+
if (ts != 0 && added_edge_num > 0) {
541+
needs_compact_.store(true, std::memory_order_relaxed);
542+
}
522543
// invalidate sort flag
523544
if (ts < unsorted_since_) {
524545
unsorted_since_ = 0;
@@ -535,6 +556,7 @@ void SingleMutableCsr<EDATA_T>::Open(Checkpoint& ckp,
535556
nbr_list_ = std::shared_ptr<IDataContainer>(ckp.OpenFile(
536557
descriptor.get_path(ModuleDescriptor::kNbrListPath).value_or(""), level));
537558
edge_num_.store(std::stoull(descriptor.get("edge_num").value_or("0")));
559+
needs_compact_.store(false, std::memory_order_relaxed);
538560
refresh_prefetch_policy();
539561
}
540562

@@ -558,6 +580,9 @@ void SingleMutableCsr<EDATA_T>::Dump(Checkpoint& ckp, CheckpointManifest& meta,
558580

559581
template <typename EDATA_T>
560582
void SingleMutableCsr<EDATA_T>::compact() {
583+
if (!needs_compact_.load(std::memory_order_relaxed)) {
584+
return;
585+
}
561586
if (!nbr_list_) {
562587
return;
563588
}
@@ -568,6 +593,7 @@ void SingleMutableCsr<EDATA_T>::compact() {
568593
data[i].timestamp.store(0, std::memory_order_relaxed);
569594
}
570595
}
596+
needs_compact_.store(false, std::memory_order_relaxed);
571597
}
572598

573599
template <typename EDATA_T>
@@ -601,6 +627,9 @@ void SingleMutableCsr<EDATA_T>::batch_delete_vertices(
601627
if (!nbr_list_) {
602628
return;
603629
}
630+
if (!src_set.empty() || !dst_set.empty()) {
631+
needs_compact_.store(true, std::memory_order_relaxed);
632+
}
604633
nbr_t* data = reinterpret_cast<nbr_t*>(nbr_list_->GetData());
605634
vid_t vnum = static_cast<vid_t>(vertex_capacity());
606635
for (auto src : src_set) {
@@ -629,6 +658,9 @@ void SingleMutableCsr<EDATA_T>::batch_delete_edges(
629658
if (!nbr_list_) {
630659
return;
631660
}
661+
if (!src_list.empty()) {
662+
needs_compact_.store(true, std::memory_order_relaxed);
663+
}
632664
nbr_t* data = reinterpret_cast<nbr_t*>(nbr_list_->GetData());
633665
vid_t vnum = static_cast<vid_t>(vertex_capacity());
634666
for (size_t i = 0; i != src_list.size(); ++i) {
@@ -653,6 +685,9 @@ void SingleMutableCsr<EDATA_T>::batch_delete_edges(
653685
if (!nbr_list_) {
654686
return;
655687
}
688+
if (!edge_list.empty()) {
689+
needs_compact_.store(true, std::memory_order_relaxed);
690+
}
656691
nbr_t* data = reinterpret_cast<nbr_t*>(nbr_list_->GetData());
657692
vid_t vnum = static_cast<vid_t>(vertex_capacity());
658693
for (const auto& edge : edge_list) {
@@ -685,6 +720,7 @@ void SingleMutableCsr<EDATA_T>::delete_edge(vid_t src, int32_t offset,
685720
if (nbr.timestamp.load() <= ts) {
686721
nbr.timestamp.store(std::numeric_limits<timestamp_t>::max());
687722
edge_num_.fetch_sub(1, std::memory_order_relaxed);
723+
needs_compact_.store(true, std::memory_order_relaxed);
688724
} else if (nbr.timestamp.load() == std::numeric_limits<timestamp_t>::max()) {
689725
LOG(ERROR) << "Fail to delete edge, already deleted.";
690726
} else {
@@ -711,6 +747,9 @@ void SingleMutableCsr<EDATA_T>::revert_delete_edge(vid_t src, vid_t nbr_vid,
711747
if (nbr.timestamp.load() == std::numeric_limits<timestamp_t>::max()) {
712748
nbr.timestamp.store(ts);
713749
edge_num_.fetch_add(1, std::memory_order_relaxed);
750+
if (ts != 0) {
751+
needs_compact_.store(true, std::memory_order_relaxed);
752+
}
714753
} else {
715754
THROW_INVALID_ARGUMENT_EXCEPTION(
716755
"Attempting to revert delete on edge that is not deleted.");
@@ -736,6 +775,9 @@ void SingleMutableCsr<EDATA_T>::batch_put_edges(
736775
nbr.data = data_list[i];
737776
nbr.timestamp.store(ts);
738777
edge_num_.fetch_add(1, std::memory_order_relaxed);
778+
if (ts != 0) {
779+
needs_compact_.store(true, std::memory_order_relaxed);
780+
}
739781
}
740782
}
741783

src/storages/graph/edge_table.cc

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -603,6 +603,9 @@ void EdgeTable::UpdateEdgeProperty(vid_t src_lid, vid_t dst_lid,
603603
if (oe_iter == oe_edges.end()) {
604604
THROW_INVALID_ARGUMENT_EXCEPTION("invalid oe offset ");
605605
}
606+
if (ts != 0) {
607+
out_csr_->mark_compaction_required();
608+
}
606609
accessor.set_data(oe_iter, prop, ts);
607610
if (meta_->is_bundled()) {
608611
auto ie_edges = in_csr_->get_generic_view(ts).get_edges(dst_lid);
@@ -611,6 +614,9 @@ void EdgeTable::UpdateEdgeProperty(vid_t src_lid, vid_t dst_lid,
611614
if (ie_iter == ie_edges.end()) {
612615
THROW_INVALID_ARGUMENT_EXCEPTION("invalid ie offset ");
613616
}
617+
if (ts != 0) {
618+
in_csr_->mark_compaction_required();
619+
}
614620
accessor.set_data(ie_iter, prop, ts);
615621
}
616622
}

0 commit comments

Comments
 (0)