From cab59feb50d3064533127862c429704459f0edf6 Mon Sep 17 00:00:00 2001 From: Gareth Lloyd Date: Thu, 2 Nov 2023 21:57:34 +0000 Subject: [PATCH] More proactive deletion --- CMakeLists.txt | 2 +- src/memory/CMakeLists.txt | 8 +-- src/storage/v2/inmemory/storage.cpp | 79 +++++++++++++++++++---------- src/storage/v2/inmemory/storage.hpp | 5 +- 4 files changed, 59 insertions(+), 35 deletions(-) diff --git a/CMakeLists.txt b/CMakeLists.txt index 4bfd4dfca..8937e9ad9 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -296,7 +296,7 @@ if (MG_ENTERPRISE) add_definitions(-DMG_ENTERPRISE) endif() -set(ENABLE_JEMALLOC ON) +option(ENABLE_JEMALLOC "Use jemalloc" ON) if (ASAN) message(WARNING "Disabling jemalloc as it doesn't work well with ASAN") diff --git a/src/memory/CMakeLists.txt b/src/memory/CMakeLists.txt index aadbbe23c..e975c1d5c 100644 --- a/src/memory/CMakeLists.txt +++ b/src/memory/CMakeLists.txt @@ -3,14 +3,14 @@ set(memory_src_files global_memory_control.cpp query_memory_control.cpp) - - -find_package(jemalloc REQUIRED) - add_library(mg-memory STATIC ${memory_src_files}) target_link_libraries(mg-memory mg-utils fmt) +message(STATUS "ENABLE_JEMALLOC: ${ENABLE_JEMALLOC}") if (ENABLE_JEMALLOC) + find_package(jemalloc REQUIRED) target_link_libraries(mg-memory Jemalloc::Jemalloc ${CMAKE_DL_LIBS}) target_compile_definitions(mg-memory PRIVATE USE_JEMALLOC=1) +else() + target_compile_definitions(mg-memory PRIVATE USE_JEMALLOC=0) endif() diff --git a/src/storage/v2/inmemory/storage.cpp b/src/storage/v2/inmemory/storage.cpp index 2bbb1888a..89f81e5ef 100644 --- a/src/storage/v2/inmemory/storage.cpp +++ b/src/storage/v2/inmemory/storage.cpp @@ -1039,10 +1039,14 @@ void InMemoryStorage::InMemoryAccessor::FinalizeTransaction() { if (!transaction_.deltas.use().empty()) { // Only hand over delta to be GC'ed if there was any deltas mem_storage->committed_transactions_.WithLock([&](auto &committed_transactions) { - // using mark of 0 as GC will assign a mark_timestamp after unlinking - committed_transactions.emplace_back(0, std::move(transaction_.deltas), + // using commit timestamp as mark, set here reduces atomic contention in GC, + // after unlinking, an actual mark_timestamp will be assigned + committed_transactions.emplace_back(*commit_timestamp_, std::move(transaction_.deltas), std::move(transaction_.commit_timestamp)); }); + } else { + // proactively clear commit_timestamp, no deltas refer to it + transaction_.commit_timestamp.reset(); } commit_timestamp_.reset(); } @@ -1245,7 +1249,7 @@ void InMemoryStorage::CollectGarbage(std::unique_lock main_ main_lock_.lock_shared(); } } else { - MG_ASSERT(main_guard.mutex() == std::addressof(main_lock_), "main_guard should be only for the main_lock_"); + DMG_ASSERT(main_guard.mutex() == std::addressof(main_lock_), "main_guard should be only for the main_lock_"); } utils::OnScopeExit lock_releaser{[&] { @@ -1272,19 +1276,20 @@ void InMemoryStorage::CollectGarbage(std::unique_lock main_ return; } - uint64_t oldest_active_start_timestamp = commit_log_->OldestActive(); - // Deltas from previous GC runs or from aborts can be cleaned up here garbage_undo_buffers_.WithLock([&](auto &garbage_undo_buffers) { - if constexpr (force) { - // if force is set to true we can simply delete all the leftover undos because - // no transaction is active - garbage_undo_buffers.clear(); - } else { - // garbage_undo_buffers is ordered, pop until we can't - while (!garbage_undo_buffers.empty() && - garbage_undo_buffers.front().mark_timestamp_ <= oldest_active_start_timestamp) { - garbage_undo_buffers.pop_front(); + if (!garbage_undo_buffers.empty()) { + if constexpr (force) { + // if force is set to true we can simply delete all the leftover undos because + // no transaction is active + garbage_undo_buffers.clear(); + } else { + uint64_t oldest_active_start_timestamp = commit_log_->OldestActive(); + // garbage_undo_buffers is ordered, pop until we can't + while (!garbage_undo_buffers.empty() && + garbage_undo_buffers.front().mark_timestamp_ <= oldest_active_start_timestamp) { + garbage_undo_buffers.pop_front(); + } } } }); @@ -1294,6 +1299,10 @@ void InMemoryStorage::CollectGarbage(std::unique_lock main_ // garbage_undo_buffers lock. std::list unlinked_undo_buffers{}; + // same as unlinked_undo_buffers, but it is known none of these deltas are accessible + // (because no unlinking was required) + std::list unlinked_undo_buffers_delete_ASAP{}; + // We will only free vertices deleted up until now in this GC cycle, and we // will do it after cleaning-up the indices. That way we are sure that all // vertices that appear in an index also exist in main storage. @@ -1310,20 +1319,13 @@ void InMemoryStorage::CollectGarbage(std::unique_lock main_ committed_transactions_.WithLock( [&](auto &committed_transactions) { committed_transactions.swap(linked_undo_buffers); }); - // Flag that will be used to determine whether the Index GC should be run. It - // should be run when there were any items that were cleaned up (there were - // updates between this run of the GC and the previous run of the GC). This - // eliminates high CPU usage when the GC doesn't have to clean up anything. - bool run_index_cleanup = !linked_undo_buffers.empty() || !garbage_undo_buffers_->empty() || need_full_scan_vertices || - need_full_scan_edges; - auto const end_linked_undo_buffers = linked_undo_buffers.end(); + uint64_t oldest_active_start_timestamp = commit_log_->OldestActive(); for (auto linked_entry = linked_undo_buffers.begin(); linked_entry != end_linked_undo_buffers;) { auto const *const commit_timestamp_ptr = linked_entry->commit_timestamp_.get(); - auto const commit_timestamp = commit_timestamp_ptr->load(std::memory_order_acquire); - // only process those that are no longer active - if (commit_timestamp >= oldest_active_start_timestamp) { + // only process those that are no longer active (here mark_timestamp_ == *commit_timestamp_) + if (linked_entry->mark_timestamp_ >= oldest_active_start_timestamp) { ++linked_entry; // can not process, skip continue; // must continue to next transaction, because committed_transactions_ was not ordered } @@ -1359,6 +1361,7 @@ void InMemoryStorage::CollectGarbage(std::unique_lock main_ // chain in a broken state. // The chain can be only read without taking any locks. + bool any_unlinks = false; for (Delta &delta : linked_entry->deltas_.use()) { while (true) { auto prev = delta.prev.Get(); @@ -1372,6 +1375,7 @@ void InMemoryStorage::CollectGarbage(std::unique_lock main_ continue; } vertex->delta = nullptr; + any_unlinks = true; if (vertex->deleted) { DMG_ASSERT(delta.action == memgraph::storage::Delta::Action::RECREATE_OBJECT); current_deleted_vertices.push_back(vertex->gid); @@ -1387,6 +1391,7 @@ void InMemoryStorage::CollectGarbage(std::unique_lock main_ continue; } edge->delta = nullptr; + any_unlinks = true; if (edge->deleted) { DMG_ASSERT(delta.action == memgraph::storage::Delta::Action::RECREATE_OBJECT); current_deleted_edges.push_back(edge->gid); @@ -1443,8 +1448,8 @@ void InMemoryStorage::CollectGarbage(std::unique_lock main_ // chain. continue; } - Delta *prev_delta = prev.delta; - prev_delta->next.store(nullptr, std::memory_order_release); + prev.delta->next.store(nullptr, std::memory_order_release); + any_unlinks = true; break; } case PreviousPtr::Type::NULLPTR: { @@ -1458,9 +1463,20 @@ void InMemoryStorage::CollectGarbage(std::unique_lock main_ // Now unlinked, move to unlinked_undo_buffers auto const to_move = linked_entry; ++linked_entry; // advanced to next before we move the list node - unlinked_undo_buffers.splice(unlinked_undo_buffers.end(), linked_undo_buffers, to_move); + if (any_unlinks) { + // At least one delta in this collection maybe accessible (move to unlinked_undo_buffers) + unlinked_undo_buffers.splice(unlinked_undo_buffers.end(), linked_undo_buffers, to_move); + } else { + // None of the deltas in this collection are accessible (move to unlinked_undo_buffers_delete_ASAP) + // Can't delete here because maybe another linked_entry unlinking way try to read these deltas + unlinked_undo_buffers_delete_ASAP.splice(unlinked_undo_buffers_delete_ASAP.end(), linked_undo_buffers, to_move); + } } + // Now we have unlinked, we can clear the inaccessible deltas + // they are not accessible to any active transaction or any future unlinking + unlinked_undo_buffers_delete_ASAP.clear(); + if (!linked_undo_buffers.empty()) { // some were not able to be collected, add them back to committed_transactions_ for the next GC run committed_transactions_.WithLock([&linked_undo_buffers](auto &committed_transactions) { @@ -1468,6 +1484,13 @@ void InMemoryStorage::CollectGarbage(std::unique_lock main_ }); } + // Flag that will be used to determine whether the Index GC should be run. It + // should be run when there were any items that were cleaned up (there were + // updates between this run of the GC and the previous run of the GC). This + // eliminates high CPU usage when the GC doesn't have to clean up anything. + bool run_index_cleanup = !current_deleted_vertices.empty() || !current_deleted_edges.empty() || + need_full_scan_vertices || need_full_scan_edges; + // After unlinking deltas from vertices, we refresh the indices. That way // we're sure that none of the vertices from `current_deleted_vertices` // appears in an index, and we can safely remove the from the main storage @@ -1480,7 +1503,7 @@ void InMemoryStorage::CollectGarbage(std::unique_lock main_ mem_unique_constraints->RemoveObsoleteEntries(oldest_active_start_timestamp); } - { + if (!unlinked_undo_buffers.empty()) { std::unique_lock guard(engine_lock_); uint64_t mark_timestamp = timestamp_; // a timestamp no active transaction can currently have diff --git a/src/storage/v2/inmemory/storage.hpp b/src/storage/v2/inmemory/storage.hpp index 73b3bfc5c..06745b170 100644 --- a/src/storage/v2/inmemory/storage.hpp +++ b/src/storage/v2/inmemory/storage.hpp @@ -422,8 +422,9 @@ class InMemoryStorage final : public Storage { GCDeltas(GCDeltas &&) = default; GCDeltas &operator=(GCDeltas &&) = default; - uint64_t mark_timestamp_{}; //!< a timestamp no active transaction currently has - BondPmrLd deltas_; //!< the deltas that need cleaning + uint64_t mark_timestamp_{}; //!< linked: non-atomic copy of commit_timestamp_, unlinked: a timestamp no active + //!< transaction currently had + BondPmrLd deltas_; //!< the deltas that need cleaning std::unique_ptr> commit_timestamp_{}; //!< the timestamp the deltas are pointing at };