From cc8d1696a96eac38b7dd58466e1bf7f952bf4a21 Mon Sep 17 00:00:00 2001 From: Andrey Vasnetsov Date: Wed, 3 Jun 2026 12:43:51 +0200 Subject: [PATCH] feat/readonly map index with live reload (#9264) * live reload function for map index * fmt --- lib/gridstore/src/gridstore/mod.rs | 5 +- lib/gridstore/src/gridstore/tests.rs | 2 +- lib/gridstore/src/gridstore/view.rs | 5 +- .../map_index/mutable_map_index/inner.rs | 26 ++++++++++ .../map_index/mutable_map_index/lifecycle.rs | 16 +----- .../read_only/live_reload.rs | 50 +++++++++++++++++++ .../mutable_map_index/read_only/mod.rs | 1 + .../universal_map_index/lifecycle.rs | 24 ++++----- .../universal_map_index/live_reload.rs | 31 ++++++++++++ .../map_index/universal_map_index/mod.rs | 1 + .../sparse/mmap_sparse_vector_storage.rs | 2 +- 11 files changed, 130 insertions(+), 33 deletions(-) create mode 100644 lib/segment/src/index/field_index/map_index/mutable_map_index/read_only/live_reload.rs create mode 100644 lib/segment/src/index/field_index/map_index/universal_map_index/live_reload.rs diff --git a/lib/gridstore/src/gridstore/mod.rs b/lib/gridstore/src/gridstore/mod.rs index 82c896327c..daa9f377f2 100644 --- a/lib/gridstore/src/gridstore/mod.rs +++ b/lib/gridstore/src/gridstore/mod.rs @@ -8,6 +8,7 @@ use std::path::PathBuf; use std::sync::Arc; use ahash::AHashMap; +use common::counter::counter_cell::CounterCell; use common::counter::hardware_counter::HardwareCounterCell; use common::counter::referenced_counter::HwMetricRefCounter; use common::fs::atomic_save_json; @@ -380,14 +381,14 @@ impl Gridstore { &self, offsets: &[PointOffset], callback: F, - hw_counter: &HardwareCounterCell, + hw_counter_cell: &CounterCell, ) -> std::result::Result<(), E> where P: AccessPattern, F: FnMut(usize, Option) -> std::result::Result<(), E>, E: From, { - self.with_view(|view| view.for_each_in_batch::(offsets, callback, hw_counter)) + self.with_view(|view| view.for_each_in_batch::(offsets, callback, hw_counter_cell)) } #[cfg(test)] diff --git a/lib/gridstore/src/gridstore/tests.rs b/lib/gridstore/src/gridstore/tests.rs index 2143cca287..099a805f54 100644 --- a/lib/gridstore/src/gridstore/tests.rs +++ b/lib/gridstore/src/gridstore/tests.rs @@ -1380,7 +1380,7 @@ fn test_for_each_in_batch_congruent_with_get_value() { batch_results[idx] = value; Ok(()) }, - &hw_counter, + hw_counter.payload_io_read_counter(), ) .unwrap(); diff --git a/lib/gridstore/src/gridstore/view.rs b/lib/gridstore/src/gridstore/view.rs index 8ee2c5f534..5b6c45a7ef 100644 --- a/lib/gridstore/src/gridstore/view.rs +++ b/lib/gridstore/src/gridstore/view.rs @@ -1,6 +1,7 @@ use std::borrow::Cow; use std::ops::ControlFlow; +use common::counter::counter_cell::CounterCell; use common::counter::hardware_counter::HardwareCounterCell; use common::counter::referenced_counter::HwMetricRefCounter; use common::generic_consts::{AccessPattern, Sequential}; @@ -112,7 +113,7 @@ impl<'a, V: Blob, S: UniversalRead> GridstoreView<'a, V, S> { &self, point_offsets: &[PointOffset], mut callback: F, - hw_counter: &HardwareCounterCell, + hw_counter_cell: &CounterCell, ) -> std::result::Result<(), E> where P: AccessPattern, @@ -129,7 +130,7 @@ impl<'a, V: Blob, S: UniversalRead> GridstoreView<'a, V, S> { self.pages .read_batch_from_pages::(pointers, self.config, |idx, raw_opt| { let value = raw_opt.map(|raw| { - hw_counter.payload_io_read_counter().incr_delta(raw.len()); + hw_counter_cell.incr_delta(raw.len()); let decompressed = self.decompress(raw); V::from_bytes(&decompressed) }); diff --git a/lib/segment/src/index/field_index/map_index/mutable_map_index/inner.rs b/lib/segment/src/index/field_index/map_index/mutable_map_index/inner.rs index 94573759de..70c9215fef 100644 --- a/lib/segment/src/index/field_index/map_index/mutable_map_index/inner.rs +++ b/lib/segment/src/index/field_index/map_index/mutable_map_index/inner.rs @@ -68,6 +68,32 @@ where self.map.entry(value).or_default().insert(idx); } + /// Remove a point from the in-memory state, updating the value map and + /// counters accordingly. + /// + /// Returns `true` if the point was within the indexed range; callers backed + /// by a store should delete the on-disk entry only in that case. + pub fn remove_point(&mut self, idx: PointOffsetType) -> bool { + if self.point_to_values.len() <= idx as usize { + return false; + } + + let removed_values = std::mem::take(&mut self.point_to_values[idx as usize]); + + if !removed_values.is_empty() { + self.indexed_points -= 1; + } + self.values_count -= removed_values.len(); + + for value in &removed_values { + if let Some(vals) = self.map.get_mut(value.borrow()) { + vals.remove(idx); + } + } + + true + } + pub(in crate::index::field_index::map_index) fn for_points_values( &self, points: impl Iterator, diff --git a/lib/segment/src/index/field_index/map_index/mutable_map_index/lifecycle.rs b/lib/segment/src/index/field_index/map_index/mutable_map_index/lifecycle.rs index 16284469d2..9e77021da2 100644 --- a/lib/segment/src/index/field_index/map_index/mutable_map_index/lifecycle.rs +++ b/lib/segment/src/index/field_index/map_index/mutable_map_index/lifecycle.rs @@ -1,4 +1,3 @@ -use std::borrow::Borrow; use std::path::PathBuf; use common::counter::hardware_counter::HardwareCounterCell; @@ -120,23 +119,10 @@ where } pub fn remove_point(&mut self, idx: PointOffsetType) -> OperationResult<()> { - if self.inner.point_to_values.len() <= idx as usize { + if !self.inner.remove_point(idx) { return Ok(()); } - let removed_values = std::mem::take(&mut self.inner.point_to_values[idx as usize]); - - if !removed_values.is_empty() { - self.inner.indexed_points -= 1; - } - self.inner.values_count -= removed_values.len(); - - for value in &removed_values { - if let Some(vals) = self.inner.map.get_mut(value.borrow()) { - vals.remove(idx); - } - } - self.storage.delete_value(idx)?; Ok(()) diff --git a/lib/segment/src/index/field_index/map_index/mutable_map_index/read_only/live_reload.rs b/lib/segment/src/index/field_index/map_index/mutable_map_index/read_only/live_reload.rs new file mode 100644 index 0000000000..7eb0154b1d --- /dev/null +++ b/lib/segment/src/index/field_index/map_index/mutable_map_index/read_only/live_reload.rs @@ -0,0 +1,50 @@ +use common::counter::hardware_counter::HardwareCounterCell; +use common::generic_consts::Random; +use common::types::PointOffsetType; +use common::universal_io::UniversalRead; +use gridstore::Blob; +use gridstore::error::GridstoreError; + +use crate::common::operation_error::OperationResult; +use crate::index::field_index::map_index::MapIndexKey; +use crate::index::field_index::map_index::mutable_map_index::read_only::ReadOnlyAppendableMapIndex; + +impl ReadOnlyAppendableMapIndex +where + Vec<::Owned>: Blob + Send + Sync, +{ + pub fn live_reload( + &mut self, + fs: &S::Fs, + deleted_points: &[PointOffsetType], + new_points: &[PointOffsetType], + hw_counter: &HardwareCounterCell, + ) -> OperationResult<()> { + self.storage.live_reload(fs)?; + + let in_memory_storage = &mut self.inner; + + for deleted_point in deleted_points { + in_memory_storage.remove_point(*deleted_point); + } + + self.storage + .view() + .for_each_in_batch::( + new_points, + |idx, maybe_values: Option>| { + let Some(values) = maybe_values else { + return Ok(()); + }; + let point_offset = new_points[idx]; + for value in values { + in_memory_storage.ingest(point_offset, value); + } + Ok(()) + }, + hw_counter.payload_index_io_read_counter(), + )?; + + Ok(()) + } +} diff --git a/lib/segment/src/index/field_index/map_index/mutable_map_index/read_only/mod.rs b/lib/segment/src/index/field_index/map_index/mutable_map_index/read_only/mod.rs index 120c7b7104..2b8499feaa 100644 --- a/lib/segment/src/index/field_index/map_index/mutable_map_index/read_only/mod.rs +++ b/lib/segment/src/index/field_index/map_index/mutable_map_index/read_only/mod.rs @@ -5,6 +5,7 @@ use super::super::MapIndexKey; use super::inner::MutableMapIndexInner; mod lifecycle; +mod live_reload; mod read_ops; /// Read-only counterpart to [`super::MutableMapIndex`]. diff --git a/lib/segment/src/index/field_index/map_index/universal_map_index/lifecycle.rs b/lib/segment/src/index/field_index/map_index/universal_map_index/lifecycle.rs index 77d82242a9..0dd9572f89 100644 --- a/lib/segment/src/index/field_index/map_index/universal_map_index/lifecycle.rs +++ b/lib/segment/src/index/field_index/map_index/universal_map_index/lifecycle.rs @@ -99,6 +99,18 @@ where is_on_disk, })) } + + /// Marks `idx` as deleted in the in-memory deletion bitvec. + /// + /// Not persisted: on reopen, deletions must be re-supplied via the + /// `deleted_points` argument to [`Self::open`]. + pub fn remove_point(&mut self, idx: PointOffsetType) { + let idx = idx as usize; + if idx < self.storage.deleted.len() && !self.storage.deleted.get_bit(idx).unwrap_or(true) { + self.storage.deleted.set(idx, true); + self.deleted_count += 1; + } + } } impl UniversalMapIndex @@ -219,18 +231,6 @@ where files } - /// Marks `idx` as deleted in the in-memory deletion bitvec. - /// - /// Not persisted: on reopen, deletions must be re-supplied via the - /// `deleted_points` argument to [`Self::open`]. - pub fn remove_point(&mut self, idx: PointOffsetType) { - let idx = idx as usize; - if idx < self.storage.deleted.len() && !self.storage.deleted.get_bit(idx).unwrap_or(true) { - self.storage.deleted.set(idx, true); - self.deleted_count += 1; - } - } - /// Populate all pages in the mmap. /// Block until all pages are populated. pub fn populate(&self) -> OperationResult<()> { diff --git a/lib/segment/src/index/field_index/map_index/universal_map_index/live_reload.rs b/lib/segment/src/index/field_index/map_index/universal_map_index/live_reload.rs new file mode 100644 index 0000000000..bf00bdaa44 --- /dev/null +++ b/lib/segment/src/index/field_index/map_index/universal_map_index/live_reload.rs @@ -0,0 +1,31 @@ +use common::counter::hardware_counter::HardwareCounterCell; +use common::persisted_hashmap::Key; +use common::types::PointOffsetType; +use common::universal_io::UniversalRead; + +use crate::common::operation_error::OperationResult; +use crate::index::field_index::map_index::MapIndexKey; +use crate::index::field_index::map_index::universal_map_index::UniversalMapIndex; + +impl UniversalMapIndex +where + N: MapIndexKey + Key + ?Sized, + S: UniversalRead, +{ + pub fn live_reload( + &mut self, + _fs: &S::Fs, + deleted_points: &[PointOffsetType], + _new_points: &[PointOffsetType], + _hw_counter: &HardwareCounterCell, + ) -> OperationResult<()> { + // No on-disk state is changing when we live-reload, as + // this UniversalMapIndex is not mutable. + // We only patch in-memory deleted bitslice representation. + for deleted_point in deleted_points { + self.remove_point(*deleted_point) + } + + Ok(()) + } +} diff --git a/lib/segment/src/index/field_index/map_index/universal_map_index/mod.rs b/lib/segment/src/index/field_index/map_index/universal_map_index/mod.rs index 142c1b7eb0..2ad0842976 100644 --- a/lib/segment/src/index/field_index/map_index/universal_map_index/mod.rs +++ b/lib/segment/src/index/field_index/map_index/universal_map_index/mod.rs @@ -10,6 +10,7 @@ use super::MapIndexKey; use crate::index::field_index::stored_point_to_values::StoredPointToValues; mod lifecycle; +mod live_reload; mod read_ops; pub(super) const DELETED_PATH: &str = "deleted.bin"; diff --git a/lib/segment/src/vector_storage/sparse/mmap_sparse_vector_storage.rs b/lib/segment/src/vector_storage/sparse/mmap_sparse_vector_storage.rs index 99df01c786..b0dbfef706 100644 --- a/lib/segment/src/vector_storage/sparse/mmap_sparse_vector_storage.rs +++ b/lib/segment/src/vector_storage/sparse/mmap_sparse_vector_storage.rs @@ -218,7 +218,7 @@ impl SparseVectorStorage for MmapSparseVectorStorage { } Ok(()) }, - &HardwareCounterCell::disposable(), + HardwareCounterCell::disposable().vector_io_read(), ) } }