From 2db06a96c1ab1287fcd08c45bc34c3c478870f68 Mon Sep 17 00:00:00 2001 From: "xiaolei.zl" Date: Mon, 14 Sep 2026 12:29:24 +0800 Subject: [PATCH 1/4] perf: skip no-op edge CSR compaction after COPY --- include/neug/storages/graph/edge_table.h | 10 ++ src/main/checkpoint_coordinator.cc | 3 +- src/storages/graph/edge_table.cc | 30 +++++ src/transaction/cow_graph_workspace.cc | 6 +- tests/storage/test_ap_index.cc | 120 ++++++++++++++++++++ tests/storage/test_edge_table.cc | 63 ++++++++++ tools/python_bind/tests/test_persistence.py | 28 +++++ 7 files changed, 258 insertions(+), 2 deletions(-) diff --git a/include/neug/storages/graph/edge_table.h b/include/neug/storages/graph/edge_table.h index 308a11dee..cf16beead 100644 --- a/include/neug/storages/graph/edge_table.h +++ b/include/neug/storages/graph/edge_table.h @@ -180,6 +180,11 @@ class NEUG_API EdgeTable { int32_t ie_offset, int32_t col_id, const Value& new_prop, timestamp_t ts); + bool NeedsCompaction( + const std::optional& sort_key_for_nbr) const noexcept { + return sort_key_for_nbr.has_value() || needs_csr_compaction_.load(); + } + void Compact(const std::optional& sort_key_for_nbr); size_t PropTableSize() const; @@ -208,6 +213,11 @@ class NEUG_API EdgeTable { std::unique_ptr table_; std::atomic table_idx_{0}; std::atomic capacity_{0}; + // Batch COPY appends checkpoint-normalized CSR entries at timestamp zero. + // Ordinary writes and deletes require a later CSR scan to normalize + // timestamps or remove tombstones. Persist this bit across incremental + // checkpoint reopen so a later bulk load cannot skip required compaction. + std::atomic needs_csr_compaction_{false}; friend class PropertyGraph; friend class EdgeTableView; diff --git a/src/main/checkpoint_coordinator.cc b/src/main/checkpoint_coordinator.cc index 9ec946b58..dca8a2c8c 100644 --- a/src/main/checkpoint_coordinator.cc +++ b/src/main/checkpoint_coordinator.cc @@ -151,7 +151,8 @@ Status CheckpointCoordinator::CommitCowWrite( // Finalize only persistent COPY targets before checkpoint consumption, // while ordinary rollback remains safe. Vertex COPY has a timestamp-zero - // tail; edge COPY needs compaction only when it has a neighbor sort key. + // tail. An edge target needs compaction when it has a neighbor sort key or + // inherited ordinary writes left timestamps/tombstones to normalize. // Keeping the target sets transaction-local avoids compacting unrelated // dirty tables inherited by the private COW graph. workspace.FinalizeBulkTablesForCheckpoint(); diff --git a/src/storages/graph/edge_table.cc b/src/storages/graph/edge_table.cc index 20f03dd78..084868299 100644 --- a/src/storages/graph/edge_table.cc +++ b/src/storages/graph/edge_table.cc @@ -473,6 +473,7 @@ EdgeTable::EdgeTable(EdgeTable&& edge_table) table_ = std::move(edge_table.table_); table_idx_ = edge_table.table_idx_.load(); capacity_ = edge_table.capacity_.load(); + needs_csr_compaction_ = edge_table.needs_csr_compaction_.load(); } EdgeTable& EdgeTable::operator=(EdgeTable&& other) noexcept { @@ -485,6 +486,7 @@ EdgeTable& EdgeTable::operator=(EdgeTable&& other) noexcept { table_ = std::move(other.table_); table_idx_ = other.table_idx_.load(); capacity_ = other.capacity_.load(); + needs_csr_compaction_ = other.needs_csr_compaction_.load(); } return *this; } @@ -502,6 +504,9 @@ void EdgeTable::Swap(EdgeTable& edge_table) { auto cap = capacity_.load(); capacity_.store(edge_table.capacity_.load()); edge_table.capacity_.store(cap); + auto needs_compaction = needs_csr_compaction_.load(); + needs_csr_compaction_.store(edge_table.needs_csr_compaction_.load()); + edge_table.needs_csr_compaction_.store(needs_compaction); } EdgeTable EdgeTable::Clone() const { @@ -519,6 +524,7 @@ EdgeTable EdgeTable::Clone() const { cow_clone.table_idx_ = table_idx_.load(); cow_clone.capacity_ = capacity_.load(); + cow_clone.needs_csr_compaction_ = needs_csr_compaction_.load(); return cow_clone; } @@ -558,12 +564,18 @@ void EdgeTable::SortByEdgeData(timestamp_t ts) { void EdgeTable::BatchDeleteVertices(const std::set& src_set, const std::set& dst_set) { + if (!src_set.empty() || !dst_set.empty()) { + needs_csr_compaction_.store(true); + } out_csr_->batch_delete_vertices(src_set, dst_set); in_csr_->batch_delete_vertices(dst_set, src_set); } void EdgeTable::BatchDeleteEdges(const std::vector& src_list, const std::vector& dst_list) { + if (!src_list.empty() || !dst_list.empty()) { + needs_csr_compaction_.store(true); + } out_csr_->batch_delete_edges(src_list, dst_list); in_csr_->batch_delete_edges(dst_list, src_list); } @@ -571,12 +583,16 @@ void EdgeTable::BatchDeleteEdges(const std::vector& src_list, void EdgeTable::BatchDeleteEdges( const std::vector>& oe_edges, const std::vector>& ie_edges) { + if (!oe_edges.empty() || !ie_edges.empty()) { + needs_csr_compaction_.store(true); + } out_csr_->batch_delete_edges(oe_edges); in_csr_->batch_delete_edges(ie_edges); } void EdgeTable::DeleteEdge(vid_t src_lid, vid_t dst_lid, int32_t oe_offset, int32_t ie_offset, timestamp_t ts) { + needs_csr_compaction_.store(true); out_csr_->delete_edge(src_lid, oe_offset, ts); in_csr_->delete_edge(dst_lid, ie_offset, ts); } @@ -653,6 +669,9 @@ void EdgeTable::UpdateEdgeProperty(vid_t src_lid, vid_t dst_lid, int32_t oe_offset, int32_t ie_offset, int32_t col_id, const Value& prop, timestamp_t ts) { + if (ts != 0) { + needs_csr_compaction_.store(true); + } auto accessor = get_edge_data_accessor(col_id); auto oe_edges = out_csr_->get_generic_view(ts).get_edges(src_lid); auto oe_iter = oe_edges.begin(); @@ -813,6 +832,9 @@ void EdgeTable::DeleteProperties(Checkpoint& ckp, std::pair EdgeTable::AddEdge( vid_t src_lid, vid_t dst_lid, const std::vector& edge_data, timestamp_t ts, Allocator& alloc, bool insert_safe) { + if (ts != 0) { + needs_csr_compaction_.store(true); + } return internal::insert_edge_into_csr_internal( *out_csr_, *in_csr_, *table_.get(), table_idx_, *meta_, src_lid, dst_lid, edge_data, ts, alloc, insert_safe); @@ -999,6 +1021,7 @@ void EdgeTable::Compact(const std::optional& sort_key_for_nbr) { out_csr_->batch_sort_by_edge_data(1); in_csr_->batch_sort_by_edge_data(1); } + needs_csr_compaction_.store(false); } size_t EdgeTable::PropTableSize() const { @@ -1224,6 +1247,10 @@ EdgeTable EdgeTable::OpenFrom(std::shared_ptr ckp, et.SetCapacity( meta.GetScalarAs(ScalarKey(src, edge, dst, "capacity")) .value_or(0)); + et.needs_csr_compaction_.store( + meta.GetScalarAs( + ScalarKey(src, edge, dst, "needs_csr_compaction")) + .value_or(meta.base_timestamp() == 0 ? 0 : 1) != 0); return et; } @@ -1249,6 +1276,8 @@ void EdgeTable::DisassembleTo(ModuleBroker& store, CheckpointManifest& meta, } meta.SetScalar(ScalarKey(src, edge, dst, "capacity"), std::to_string(GetCapacity())); + meta.SetScalar(ScalarKey(src, edge, dst, "needs_csr_compaction"), + needs_csr_compaction_.load() ? "1" : "0"); } void EdgeTable::ReuseCheckpointModules(Checkpoint& ckp, @@ -1271,6 +1300,7 @@ void EdgeTable::ReuseCheckpointModules(Checkpoint& ckp, meta.CopyScalarFrom(prev, ScalarKey(src, edge, dst, "table_idx")); } meta.CopyScalarFrom(prev, ScalarKey(src, edge, dst, "capacity")); + meta.CopyScalarFrom(prev, ScalarKey(src, edge, dst, "needs_csr_compaction")); } } // namespace neug diff --git a/src/transaction/cow_graph_workspace.cc b/src/transaction/cow_graph_workspace.cc index e1b712b37..22431f1d7 100644 --- a/src/transaction/cow_graph_workspace.cc +++ b/src/transaction/cow_graph_workspace.cc @@ -58,7 +58,11 @@ void CowGraphWorkspace::FinalizeBulkTablesForCheckpoint() { graph.schema().parse_edge_label(edge_triplet_id); const auto& sort_key = graph.schema().get_sort_key_for_nbr(src_label, dst_label, edge_label); - graph.get_edge_table_by_index(edge_triplet_id).Compact(sort_key); + auto& edge_table = graph.get_edge_table_by_index(edge_triplet_id); + if (!edge_table.NeedsCompaction(sort_key)) { + continue; + } + edge_table.Compact(sort_key); } } diff --git a/tests/storage/test_ap_index.cc b/tests/storage/test_ap_index.cc index 183f56028..db97dfedf 100644 --- a/tests/storage/test_ap_index.cc +++ b/tests/storage/test_ap_index.cc @@ -29,6 +29,7 @@ #include "neug/common/types/value.h" #include "neug/storages/checkpoint_manager.h" #include "neug/storages/container/i_container.h" +#include "neug/storages/csr/mutable_csr.h" #include "neug/storages/graph/graph_interface.h" #include "neug/storages/graph/graph_view.h" #include "neug/storages/graph/property_graph.h" @@ -76,6 +77,20 @@ class FailingIndex : public ExampleIndex { static inline FailurePoint failure_point_{FailurePoint::kNone}; }; +class CountingMutableCsr : public MutableCsr { + public: + explicit CountingMutableCsr(std::shared_ptr compact_count) + : compact_count_(std::move(compact_count)) {} + + void compact() override { + ++(*compact_count_); + MutableCsr::compact(); + } + + private: + std::shared_ptr compact_count_; +}; + TEST(ModuleDescriptorTest, RequiredDefaultsTrueAndRoundTripsFalse) { ModuleDescriptor required; EXPECT_TRUE(required.required); @@ -960,15 +975,120 @@ TEST_F(APIndexTest, BatchLoadFinalizesVertexTimestampAndEdgeOrder) { EXPECT_EQ( graph_->get_vertex_table(item).get_vertex_timestamp().InitVertexNum(), 0); EXPECT_EQ(plain_edge_count(0), 1); + EXPECT_TRUE( + graph_->get_edge_table(item, item, plain).NeedsCompaction(std::nullopt)); + EXPECT_TRUE(graph_->get_edge_table(item, item, weighted) + .NeedsCompaction(std::string("weight"))); // CommitCowWrite finalizes the recorded COPY targets right before the // checkpoint consumes the private graph; drive the same code path here. workspace_->FinalizeBulkTablesForCheckpoint(); + EXPECT_FALSE( + graph_->get_edge_table(item, item, plain).NeedsCompaction(std::nullopt)); expect_finalized(); CheckpointDirtyAndReopen(); + EXPECT_FALSE( + graph_->get_edge_table(item, item, plain).NeedsCompaction(std::nullopt)); expect_finalized(); } +TEST_F(APIndexTest, EdgeCompactionStateSurvivesCheckpointReopen) { + CreateItemTable(); + const auto item = graph_->schema().get_vertex_label_id("Item"); + auto vertices = ap_->BatchAddVertices( + item, MakeItemSupplier({{1, 10}, {2, 20}, {3, 30}})); + ASSERT_TRUE(vertices) << vertices.error().ToString(); + + CreateEdgeTypeParamBuilder edge_builder; + auto create_edge = ap_->CreateEdgeType(edge_builder.SrcLabel("Item") + .DstLabel("Item") + .EdgeLabel("plain") + .Build()); + ASSERT_TRUE(create_edge.ok()) << create_edge.ToString(); + const auto plain = graph_->schema().get_edge_label_id("plain"); + + vid_t source = 0; + vid_t destination = 0; + ASSERT_TRUE(ap_->GetVertexIndex(item, Value::INT32(1), source)); + ASSERT_TRUE(ap_->GetVertexIndex(item, Value::INT32(2), destination)); + + CowGraphStorage dml(*workspace_, 0, 7, allocator_); + const void* property = nullptr; + ASSERT_TRUE( + dml.AddEdge(item, source, item, destination, plain, {}, property).ok()); + EXPECT_TRUE( + graph_->get_edge_table(item, item, plain).NeedsCompaction(std::nullopt)); + + CheckpointDirtyAndReopen(); + EXPECT_TRUE( + graph_->get_edge_table(item, item, plain).NeedsCompaction(std::nullopt)); + + graph_->get_edge_table(item, item, plain).Compact(std::nullopt); + EXPECT_FALSE( + graph_->get_edge_table(item, item, plain).NeedsCompaction(std::nullopt)); + + auto edges = std::make_shared(); + edges->set(0, MakeValueColumn(std::vector{1})); + edges->set(1, MakeValueColumn(std::vector{3})); + auto add_edges = ap_->BatchAddEdges( + item, item, plain, + std::make_shared( + std::vector>{std::move(edges)})); + ASSERT_TRUE(add_edges.ok()) << add_edges.ToString(); + EXPECT_FALSE( + graph_->get_edge_table(item, item, plain).NeedsCompaction(std::nullopt)); + + CheckpointDirtyAndReopen(); + EXPECT_FALSE( + graph_->get_edge_table(item, item, plain).NeedsCompaction(std::nullopt)); +} + +TEST_F(APIndexTest, BulkFinalizeSkipsCleanPlainEdgeCsrCompaction) { + CreateItemTable(); + const auto item = graph_->schema().get_vertex_label_id("Item"); + auto vertices = + ap_->BatchAddVertices(item, MakeItemSupplier({{1, 10}, {2, 20}})); + ASSERT_TRUE(vertices) << vertices.error().ToString(); + + CreateEdgeTypeParamBuilder edge_builder; + auto create_edge = ap_->CreateEdgeType(edge_builder.SrcLabel("Item") + .DstLabel("Item") + .EdgeLabel("plain") + .Build()); + ASSERT_TRUE(create_edge.ok()) << create_edge.ToString(); + const auto plain = graph_->schema().get_edge_label_id("plain"); + const auto edge_triplet = + graph_->schema().generate_edge_label(item, item, plain); + auto& edge_table = graph_->get_edge_table(item, item, plain); + + auto compact_count = std::make_shared(0); + const auto make_counting_csr = [&] { + auto csr = std::make_unique(compact_count); + csr->Open(*checkpoint_mgr_.Current(), ModuleDescriptor{}, + MemoryLevel::kInMemory); + csr->resize(2); + return csr; + }; + edge_table.SetOutCsr(make_counting_csr()); + edge_table.SetInCsr(make_counting_csr()); + workspace_->MarkBulkEdgeTableForCheckpoint(edge_triplet); + + workspace_->FinalizeBulkTablesForCheckpoint(); + EXPECT_EQ(*compact_count, 0); + + const void* property = nullptr; + int32_t offset = 0; + ASSERT_TRUE(graph_ + ->AddEdge(item, 0, item, 1, plain, {}, 7, allocator_, offset, + property, false) + .ok()); + EXPECT_TRUE(edge_table.NeedsCompaction(std::nullopt)); + + workspace_->FinalizeBulkTablesForCheckpoint(); + EXPECT_EQ(*compact_count, 2); + EXPECT_FALSE(edge_table.NeedsCompaction(std::nullopt)); +} + TEST_F(APIndexTest, PartialBatchFailureIsDiscardedWithPrivateWorkspace) { CreateItemTable(); ResetStorageAdapter(); diff --git a/tests/storage/test_edge_table.cc b/tests/storage/test_edge_table.cc index a4087f57e..7c89b0187 100644 --- a/tests/storage/test_edge_table.cc +++ b/tests/storage/test_edge_table.cc @@ -27,6 +27,7 @@ #include "neug/storages/csr/csr_view_utils.h" #include "neug/storages/graph/edge_table.h" #include "neug/storages/loader/loader_utils.h" +#include "neug/storages/module/module_broker.h" #include "neug/storages/module_descriptor.h" #include "unittest/utils.h" @@ -995,6 +996,8 @@ TEST_F(EdgeTableTest, TestEdgeTableCompaction) { this->edge_table->AddEdge(src_lids[i], dst_lids[i], edge_data[i], 0, allocator, false); } + EXPECT_FALSE(this->edge_table->NeedsCompaction(std::nullopt)); + EXPECT_TRUE(this->edge_table->NeedsCompaction(std::string("data"))); this->ExpectBundledStats(edge_num); auto oe_view = this->edge_table->get_outgoing_view(neug::MAX_TIMESTAMP); auto ie_view = this->edge_table->get_incoming_view(neug::MAX_TIMESTAMP); @@ -1020,8 +1023,10 @@ TEST_F(EdgeTableTest, TestEdgeTableCompaction) { delete_count++; } } + EXPECT_TRUE(this->edge_table->NeedsCompaction(std::nullopt)); this->ExpectBundledStats(edge_num - delete_count); this->edge_table->Compact(std::nullopt); + EXPECT_FALSE(this->edge_table->NeedsCompaction(std::nullopt)); this->ExpectBundledStats(edge_num - delete_count); size_t edge_count = 0; for (size_t i = 0; i < dst_lids.size(); ++i) { @@ -1033,6 +1038,64 @@ TEST_F(EdgeTableTest, TestEdgeTableCompaction) { EXPECT_EQ(edge_count, edge_num - delete_count); } +TEST_F(EdgeTableTest, LegacyCheckpointUsesConservativeCompactionFallback) { + auto ckp = make_checkpoint(workspace()); + auto edge_schema = + schema_.get_edge_schema(src_label_, dst_label_, edge_label_empty_); + + const auto needs_compaction_without_scalar = [&](uint64_t base_timestamp) { + EdgeTable table(edge_schema); + table.Init(ckp, MemoryLevel::kInMemory); + + ModuleBroker store; + store.SetModule(EdgeTable::KeyOutCsr("person", "create0", "comment"), + table.TakeOutCsr()); + store.SetModule(EdgeTable::KeyInCsr("person", "create0", "comment"), + table.TakeInCsr()); + + CheckpointManifest legacy_manifest(base_timestamp); + auto reopened = EdgeTable::OpenFrom( + ckp, edge_schema, store, legacy_manifest, MemoryLevel::kInMemory); + return reopened.NeedsCompaction(std::nullopt); + }; + + EXPECT_FALSE(needs_compaction_without_scalar(0)); + EXPECT_TRUE(needs_compaction_without_scalar(7)); +} + +TEST_F(EdgeTableTest, CompactionStateTracksMutationEntrypoints) { + auto ckp = make_checkpoint(workspace()); + this->InitIndexers(*ckp, 2, 2); + this->ConstructEdgeTable(src_label_, dst_label_, edge_label_int_); + this->OpenEdgeTableInMemory(ckp, CheckpointManifest(), 2, 2); + Allocator allocator(MemoryLevel::kInMemory, allocator_dir_); + + auto inserted = + this->edge_table->AddEdge(0, 1, {Value::INT32(1)}, 0, allocator, false); + ASSERT_EQ(inserted.first, 0); + EXPECT_FALSE(this->edge_table->NeedsCompaction(std::nullopt)); + + this->edge_table->UpdateEdgeProperty(0, 1, 0, 0, 0, Value::INT32(2), 0); + EXPECT_FALSE(this->edge_table->NeedsCompaction(std::nullopt)); + this->edge_table->UpdateEdgeProperty(0, 1, 0, 0, 0, Value::INT32(3), 7); + EXPECT_TRUE(this->edge_table->NeedsCompaction(std::nullopt)); + this->edge_table->Compact(std::nullopt); + + this->edge_table->BatchDeleteEdges(std::vector{0}, + std::vector{1}); + EXPECT_TRUE(this->edge_table->NeedsCompaction(std::nullopt)); + this->edge_table->Compact(std::nullopt); + + this->edge_table->BatchDeleteEdges( + std::vector>{{0, 0}}, + std::vector>{{1, 0}}); + EXPECT_TRUE(this->edge_table->NeedsCompaction(std::nullopt)); + this->edge_table->Compact(std::nullopt); + + this->edge_table->BatchDeleteVertices(std::set{0}, {}); + EXPECT_TRUE(this->edge_table->NeedsCompaction(std::nullopt)); +} + TEST_F(EdgeTableTest, TestUpdateEdgeData) { auto ckp = make_checkpoint(workspace()); int64_t src_num = 10; diff --git a/tools/python_bind/tests/test_persistence.py b/tools/python_bind/tests/test_persistence.py index b42c4ccb7..6eef1e79f 100644 --- a/tools/python_bind/tests/test_persistence.py +++ b/tools/python_bind/tests/test_persistence.py @@ -472,6 +472,34 @@ def test_copy_from_edge_finalizes_sort_key_before_checkpoint(tmp_path): db.close() +def test_copy_from_plain_edge_survives_checkpoint_reopen(tmp_path): + db_dir = tmp_path / "copy_plain_edge" + people_csv = tmp_path / "plain_people.csv" + edges_csv = tmp_path / "plain_edges.csv" + people_csv.write_text("1\n2\n3\n4\n") + edges_csv.write_text("1,2\n1,3\n2,4\n") + + db = Database(db_path=str(db_dir), mode="w") + conn = db.connect() + conn.execute("CREATE NODE TABLE person(id INT64, PRIMARY KEY(id));") + conn.execute("CREATE REL TABLE follows(FROM person TO person, MANY_TO_MANY);") + conn.execute(f'COPY person FROM "{people_csv}" (header=false);') + conn.execute(f'COPY follows FROM "{edges_csv}" (header=false, delim=",");') + assert sorted( + list(conn.execute("MATCH (a:person)-[:follows]->(b:person) RETURN a.id, b.id;")) + ) == [[1, 2], [1, 3], [2, 4]] + conn.close() + db.close() + + db = Database(db_path=str(db_dir), mode="r") + conn = db.connect() + assert sorted( + list(conn.execute("MATCH (a:person)-[:follows]->(b:person) RETURN a.id, b.id;")) + ) == [[1, 2], [1, 3], [2, 4]] + conn.close() + db.close() + + @pytest.mark.skip(reason="TODO(zhanglei,lexiao): get view from invalid vid") def test_join_queries(modern_graph): conn = modern_graph From 69c3f5a2152ccb0c6e3828f00eca88c53ea745f9 Mon Sep 17 00:00:00 2001 From: "xiaolei.zl" Date: Mon, 14 Sep 2026 16:44:36 +0800 Subject: [PATCH 2/4] perf(storage): skip clean edge CSR compaction --- src/storages/graph/property_graph.cc | 5 +++- tests/storage/test_ap_index.cc | 45 ++++++++++++++++++++++++++++ 2 files changed, 49 insertions(+), 1 deletion(-) diff --git a/src/storages/graph/property_graph.cc b/src/storages/graph/property_graph.cc index f99cc8d73..60f232315 100644 --- a/src/storages/graph/property_graph.cc +++ b/src/storages/graph/property_graph.cc @@ -1024,7 +1024,10 @@ void PropertyGraph::Compact() { } const auto& sort_key_for_nbr = schema_.get_sort_key_for_nbr(src_label_i, dst_label_i, e_label_i); - edge_tables_.at(index).Compact(sort_key_for_nbr); + auto& edge_table = edge_tables_.at(index); + if (edge_table.NeedsCompaction(sort_key_for_nbr)) { + edge_table.Compact(sort_key_for_nbr); + } } } } diff --git a/tests/storage/test_ap_index.cc b/tests/storage/test_ap_index.cc index db97dfedf..4a4796943 100644 --- a/tests/storage/test_ap_index.cc +++ b/tests/storage/test_ap_index.cc @@ -1089,6 +1089,51 @@ TEST_F(APIndexTest, BulkFinalizeSkipsCleanPlainEdgeCsrCompaction) { EXPECT_FALSE(edge_table.NeedsCompaction(std::nullopt)); } +TEST_F(APIndexTest, PropertyGraphCompactSkipsCleanPlainEdgeCsrCompaction) { + CreateItemTable(); + const auto item = graph_->schema().get_vertex_label_id("Item"); + auto vertices = + ap_->BatchAddVertices(item, MakeItemSupplier({{1, 10}, {2, 20}})); + ASSERT_TRUE(vertices) << vertices.error().ToString(); + + CreateEdgeTypeParamBuilder edge_builder; + auto create_edge = ap_->CreateEdgeType(edge_builder.SrcLabel("Item") + .DstLabel("Item") + .EdgeLabel("plain") + .Build()); + ASSERT_TRUE(create_edge.ok()) << create_edge.ToString(); + const auto plain = graph_->schema().get_edge_label_id("plain"); + auto& edge_table = graph_->get_edge_table(item, item, plain); + + auto compact_count = std::make_shared(0); + const auto make_counting_csr = [&] { + auto csr = std::make_unique(compact_count); + csr->Open(*checkpoint_mgr_.Current(), ModuleDescriptor{}, + MemoryLevel::kInMemory); + csr->resize(2); + return csr; + }; + edge_table.SetOutCsr(make_counting_csr()); + edge_table.SetInCsr(make_counting_csr()); + graph_->MarkEdgeTableDirty(item, item, plain); + + graph_->Compact(); + EXPECT_EQ(*compact_count, 0); + + const void* property = nullptr; + int32_t offset = 0; + ASSERT_TRUE(graph_ + ->AddEdge(item, 0, item, 1, plain, {}, 7, allocator_, offset, + property, false) + .ok()); + graph_->MarkEdgeTableDirty(item, item, plain); + + graph_->Compact(); + EXPECT_EQ(*compact_count, 2); + EXPECT_FALSE( + graph_->get_edge_table(item, item, plain).NeedsCompaction(std::nullopt)); +} + TEST_F(APIndexTest, PartialBatchFailureIsDiscardedWithPrivateWorkspace) { CreateItemTable(); ResetStorageAdapter(); From fc50b7174ea97d1d662e0cb3d7816bd8022f44ac Mon Sep 17 00:00:00 2001 From: "xiaolei.zl" Date: Mon, 14 Sep 2026 17:51:59 +0800 Subject: [PATCH 3/4] refactor(storage): clarify CSR normalization state --- include/neug/storages/graph/edge_table.h | 14 +++++---- src/storages/graph/edge_table.cc | 37 ++++++++++++------------ 2 files changed, 27 insertions(+), 24 deletions(-) diff --git a/include/neug/storages/graph/edge_table.h b/include/neug/storages/graph/edge_table.h index cf16beead..50005f5c1 100644 --- a/include/neug/storages/graph/edge_table.h +++ b/include/neug/storages/graph/edge_table.h @@ -182,7 +182,7 @@ class NEUG_API EdgeTable { bool NeedsCompaction( const std::optional& sort_key_for_nbr) const noexcept { - return sort_key_for_nbr.has_value() || needs_csr_compaction_.load(); + return sort_key_for_nbr.has_value() || needs_csr_normalization_.load(); } void Compact(const std::optional& sort_key_for_nbr); @@ -213,11 +213,13 @@ class NEUG_API EdgeTable { std::unique_ptr
table_; std::atomic table_idx_{0}; std::atomic capacity_{0}; - // Batch COPY appends checkpoint-normalized CSR entries at timestamp zero. - // Ordinary writes and deletes require a later CSR scan to normalize - // timestamps or remove tombstones. Persist this bit across incremental - // checkpoint reopen so a later bulk load cannot skip required compaction. - std::atomic needs_csr_compaction_{false}; + // DirtyTracker answers whether this table must be checkpointed; this bit + // answers whether its CSR contents need an O(E) normalization scan. The + // states are independent: timestamp-zero COPY is dirty but normalized, while + // an incremental checkpoint clears dirtiness without removing MVCC + // timestamps or tombstones. Keep this state EdgeTable-local and persist it + // across reopen so COPY can skip redundant compaction safely. + std::atomic needs_csr_normalization_{false}; friend class PropertyGraph; friend class EdgeTableView; diff --git a/src/storages/graph/edge_table.cc b/src/storages/graph/edge_table.cc index 084868299..c520598a7 100644 --- a/src/storages/graph/edge_table.cc +++ b/src/storages/graph/edge_table.cc @@ -473,7 +473,7 @@ EdgeTable::EdgeTable(EdgeTable&& edge_table) table_ = std::move(edge_table.table_); table_idx_ = edge_table.table_idx_.load(); capacity_ = edge_table.capacity_.load(); - needs_csr_compaction_ = edge_table.needs_csr_compaction_.load(); + needs_csr_normalization_ = edge_table.needs_csr_normalization_.load(); } EdgeTable& EdgeTable::operator=(EdgeTable&& other) noexcept { @@ -486,7 +486,7 @@ EdgeTable& EdgeTable::operator=(EdgeTable&& other) noexcept { table_ = std::move(other.table_); table_idx_ = other.table_idx_.load(); capacity_ = other.capacity_.load(); - needs_csr_compaction_ = other.needs_csr_compaction_.load(); + needs_csr_normalization_ = other.needs_csr_normalization_.load(); } return *this; } @@ -504,9 +504,9 @@ void EdgeTable::Swap(EdgeTable& edge_table) { auto cap = capacity_.load(); capacity_.store(edge_table.capacity_.load()); edge_table.capacity_.store(cap); - auto needs_compaction = needs_csr_compaction_.load(); - needs_csr_compaction_.store(edge_table.needs_csr_compaction_.load()); - edge_table.needs_csr_compaction_.store(needs_compaction); + auto needs_normalization = needs_csr_normalization_.load(); + needs_csr_normalization_.store(edge_table.needs_csr_normalization_.load()); + edge_table.needs_csr_normalization_.store(needs_normalization); } EdgeTable EdgeTable::Clone() const { @@ -524,7 +524,7 @@ EdgeTable EdgeTable::Clone() const { cow_clone.table_idx_ = table_idx_.load(); cow_clone.capacity_ = capacity_.load(); - cow_clone.needs_csr_compaction_ = needs_csr_compaction_.load(); + cow_clone.needs_csr_normalization_ = needs_csr_normalization_.load(); return cow_clone; } @@ -565,7 +565,7 @@ void EdgeTable::SortByEdgeData(timestamp_t ts) { void EdgeTable::BatchDeleteVertices(const std::set& src_set, const std::set& dst_set) { if (!src_set.empty() || !dst_set.empty()) { - needs_csr_compaction_.store(true); + needs_csr_normalization_.store(true); } out_csr_->batch_delete_vertices(src_set, dst_set); in_csr_->batch_delete_vertices(dst_set, src_set); @@ -574,7 +574,7 @@ void EdgeTable::BatchDeleteVertices(const std::set& src_set, void EdgeTable::BatchDeleteEdges(const std::vector& src_list, const std::vector& dst_list) { if (!src_list.empty() || !dst_list.empty()) { - needs_csr_compaction_.store(true); + needs_csr_normalization_.store(true); } out_csr_->batch_delete_edges(src_list, dst_list); in_csr_->batch_delete_edges(dst_list, src_list); @@ -584,7 +584,7 @@ void EdgeTable::BatchDeleteEdges( const std::vector>& oe_edges, const std::vector>& ie_edges) { if (!oe_edges.empty() || !ie_edges.empty()) { - needs_csr_compaction_.store(true); + needs_csr_normalization_.store(true); } out_csr_->batch_delete_edges(oe_edges); in_csr_->batch_delete_edges(ie_edges); @@ -592,7 +592,7 @@ void EdgeTable::BatchDeleteEdges( void EdgeTable::DeleteEdge(vid_t src_lid, vid_t dst_lid, int32_t oe_offset, int32_t ie_offset, timestamp_t ts) { - needs_csr_compaction_.store(true); + needs_csr_normalization_.store(true); out_csr_->delete_edge(src_lid, oe_offset, ts); in_csr_->delete_edge(dst_lid, ie_offset, ts); } @@ -670,7 +670,7 @@ void EdgeTable::UpdateEdgeProperty(vid_t src_lid, vid_t dst_lid, int32_t col_id, const Value& prop, timestamp_t ts) { if (ts != 0) { - needs_csr_compaction_.store(true); + needs_csr_normalization_.store(true); } auto accessor = get_edge_data_accessor(col_id); auto oe_edges = out_csr_->get_generic_view(ts).get_edges(src_lid); @@ -833,7 +833,7 @@ std::pair EdgeTable::AddEdge( vid_t src_lid, vid_t dst_lid, const std::vector& edge_data, timestamp_t ts, Allocator& alloc, bool insert_safe) { if (ts != 0) { - needs_csr_compaction_.store(true); + needs_csr_normalization_.store(true); } return internal::insert_edge_into_csr_internal( *out_csr_, *in_csr_, *table_.get(), table_idx_, *meta_, src_lid, dst_lid, @@ -1021,7 +1021,7 @@ void EdgeTable::Compact(const std::optional& sort_key_for_nbr) { out_csr_->batch_sort_by_edge_data(1); in_csr_->batch_sort_by_edge_data(1); } - needs_csr_compaction_.store(false); + needs_csr_normalization_.store(false); } size_t EdgeTable::PropTableSize() const { @@ -1247,9 +1247,9 @@ EdgeTable EdgeTable::OpenFrom(std::shared_ptr ckp, et.SetCapacity( meta.GetScalarAs(ScalarKey(src, edge, dst, "capacity")) .value_or(0)); - et.needs_csr_compaction_.store( + et.needs_csr_normalization_.store( meta.GetScalarAs( - ScalarKey(src, edge, dst, "needs_csr_compaction")) + ScalarKey(src, edge, dst, "needs_csr_normalization")) .value_or(meta.base_timestamp() == 0 ? 0 : 1) != 0); return et; } @@ -1276,8 +1276,8 @@ void EdgeTable::DisassembleTo(ModuleBroker& store, CheckpointManifest& meta, } meta.SetScalar(ScalarKey(src, edge, dst, "capacity"), std::to_string(GetCapacity())); - meta.SetScalar(ScalarKey(src, edge, dst, "needs_csr_compaction"), - needs_csr_compaction_.load() ? "1" : "0"); + meta.SetScalar(ScalarKey(src, edge, dst, "needs_csr_normalization"), + needs_csr_normalization_.load() ? "1" : "0"); } void EdgeTable::ReuseCheckpointModules(Checkpoint& ckp, @@ -1300,7 +1300,8 @@ void EdgeTable::ReuseCheckpointModules(Checkpoint& ckp, meta.CopyScalarFrom(prev, ScalarKey(src, edge, dst, "table_idx")); } meta.CopyScalarFrom(prev, ScalarKey(src, edge, dst, "capacity")); - meta.CopyScalarFrom(prev, ScalarKey(src, edge, dst, "needs_csr_compaction")); + meta.CopyScalarFrom(prev, + ScalarKey(src, edge, dst, "needs_csr_normalization")); } } // namespace neug From 1ec3770c8478d66573fb7c0e79d52e03ed631079 Mon Sep 17 00:00:00 2001 From: "xiaolei.zl" Date: Mon, 14 Sep 2026 18:09:21 +0800 Subject: [PATCH 4/4] fix(storage): track GraphView edge normalization --- include/neug/storages/graph/graph_view.h | 2 ++ src/storages/graph/graph_view.cc | 4 ++++ tests/storage/test_graph_view.cc | 24 ++++++++++++++++++++++++ 3 files changed, 30 insertions(+) diff --git a/include/neug/storages/graph/graph_view.h b/include/neug/storages/graph/graph_view.h index 7ce0251af..8e8a143cf 100644 --- a/include/neug/storages/graph/graph_view.h +++ b/include/neug/storages/graph/graph_view.h @@ -14,6 +14,7 @@ */ #pragma once +#include #include #include #include @@ -105,6 +106,7 @@ class EdgeTableView { CsrBase* out_csr_{nullptr}; CsrBase* in_csr_{nullptr}; std::atomic* table_idx_{nullptr}; + std::atomic* needs_csr_normalization_{nullptr}; TableView view_; }; diff --git a/src/storages/graph/graph_view.cc b/src/storages/graph/graph_view.cc index 79f9e0b34..5e6c54f3d 100644 --- a/src/storages/graph/graph_view.cc +++ b/src/storages/graph/graph_view.cc @@ -145,6 +145,7 @@ EdgeTableView::EdgeTableView(EdgeTable& table) out_csr_(table.out_csr_.get()), in_csr_(table.in_csr_.get()), table_idx_(&table.table_idx_), + needs_csr_normalization_(&table.needs_csr_normalization_), view_(*table.table()) {} CsrView EdgeTableView::GetOutgoingView(timestamp_t ts) const { @@ -200,6 +201,9 @@ EdgeDataAccessor EdgeTableView::GetDataAccessor( std::pair EdgeTableView::AddEdge( vid_t src_lid, vid_t dst_lid, const std::vector& properties, timestamp_t ts, Allocator& alloc, bool insert_safe) { + if (ts != 0) { + needs_csr_normalization_->store(true); + } return internal::insert_edge_into_csr_internal( *out_csr_, *in_csr_, view_, *table_idx_, *meta_, src_lid, dst_lid, properties, ts, alloc, insert_safe); diff --git a/tests/storage/test_graph_view.cc b/tests/storage/test_graph_view.cc index 2da1d8481..977a0fbfd 100644 --- a/tests/storage/test_graph_view.cc +++ b/tests/storage/test_graph_view.cc @@ -276,6 +276,30 @@ TEST_F(GraphViewTest, EdgeIncomingMirrorsOutgoing) { EXPECT_EQ(out_count, 2u); } +TEST_F(GraphViewTest, AddEdgeTracksCsrNormalizationState) { + GraphView view(*graph_); + label_t person_label = view.schema().get_vertex_label_id("person"); + label_t knows_label = view.schema().get_edge_label_id("knows"); + auto& edge_table = + graph_->get_edge_table(person_label, person_label, knows_label); + + ASSERT_FALSE(edge_table.NeedsCompaction(std::nullopt)); + + int32_t edge_offset = 0; + const void* edge_property = nullptr; + ASSERT_TRUE(view.AddEdge(person_label, 2, person_label, 0, knows_label, + {neug::Value::DOUBLE(0.8)}, 0, *alloc_, edge_offset, + edge_property) + .ok()); + EXPECT_FALSE(edge_table.NeedsCompaction(std::nullopt)); + + ASSERT_TRUE(view.AddEdge(person_label, 0, person_label, 2, knows_label, + {neug::Value::DOUBLE(0.9)}, 7, *alloc_, edge_offset, + edge_property) + .ok()); + EXPECT_TRUE(edge_table.NeedsCompaction(std::nullopt)); +} + TEST_F(GraphViewTest, EdgeDataAccessorByIdAndName) { GraphView view(*graph_); label_t person_label = view.schema().get_vertex_label_id("person");