From 7b196be23191ca4367d7f8ea6adf2f15075e22fd Mon Sep 17 00:00:00 2001 From: Jojii <15957865+JojiiOfficial@users.noreply.github.com> Date: Thu, 19 Mar 2026 11:33:09 +0100 Subject: [PATCH] Correct calculation of deferred point counts (#8366) * Don't account for deferred points in some places # Conflicts: # lib/collection/src/shards/local_shard/scroll.rs * Add in QueryContext * Cover more places * Coderabbit review remarks * Properly count amount of deleted deferred points (#8386) * Properly count amount of deleted deferred points * Prevent double-counting of the same point * Remove hints to estimations * Properly handle counts in ProxySegment * Add tests and fix deleted point count issue * Adjusts tests + fix issues * Separte fields for deferred points in telemetry (SegmentInfo) * Remove deferred_points_count() * openapi * Fix test by manually calculating visible points * Adjust ProxySegment test to revertion of SegmentInfo * Throw error if collection was not found in telemetry (e2e Test) * Update lib/segment/src/segment/segment_ops.rs Co-authored-by: Andrey Vasnetsov * Make new fields in API optional * Don't take range if no deferred point exist * Review remarks --------- Co-authored-by: Andrey Vasnetsov --- docs/redoc/master/openapi.json | 12 +++ .../collection_manager/segments_searcher.rs | 2 +- .../src/shards/local_shard/indexed_only.rs | 2 +- .../src/shards/local_shard/scroll.rs | 2 +- .../src/update_workers/update_worker.rs | 2 +- lib/edge/src/scroll.rs | 2 +- lib/segment/src/data_types/query_context.rs | 3 +- lib/segment/src/entry/entry_point.rs | 12 ++- lib/segment/src/segment/entry.rs | 51 ++++++------ lib/segment/src/segment/facet.rs | 11 +-- lib/segment/src/segment/mod.rs | 11 ++- lib/segment/src/segment/order_by.rs | 4 +- lib/segment/src/segment/sampling.rs | 6 +- lib/segment/src/segment/scroll.rs | 6 +- lib/segment/src/segment/segment_ops.rs | 70 ++++++++++++---- lib/segment/src/segment/tests.rs | 82 ++++++++++++++++--- .../segment_constructor_base.rs | 16 +++- lib/segment/src/types.rs | 2 + .../src/optimizers/indexing_optimizer.rs | 4 +- lib/shard/src/proxy_segment/mod.rs | 3 + lib/shard/src/proxy_segment/segment_entry.rs | 71 ++++++++++++++-- lib/shard/src/proxy_segment/tests.rs | 64 +++++++++++++++ tests/openapi/test_deferred_points.py | 54 ++++++++++-- 23 files changed, 392 insertions(+), 100 deletions(-) diff --git a/docs/redoc/master/openapi.json b/docs/redoc/master/openapi.json index be45b2a64e..38bfcb1eae 100644 --- a/docs/redoc/master/openapi.json +++ b/docs/redoc/master/openapi.json @@ -12778,6 +12778,18 @@ "format": "uint", "minimum": 0 }, + "num_deferred_points": { + "type": "integer", + "format": "uint", + "minimum": 0, + "nullable": true + }, + "num_deleted_deferred_points": { + "type": "integer", + "format": "uint", + "minimum": 0, + "nullable": true + }, "num_indexed_vectors": { "type": "integer", "format": "uint", diff --git a/lib/collection/src/collection_manager/segments_searcher.rs b/lib/collection/src/collection_manager/segments_searcher.rs index 7a86004781..d1261b76b7 100644 --- a/lib/collection/src/collection_manager/segments_searcher.rs +++ b/lib/collection/src/collection_manager/segments_searcher.rs @@ -707,7 +707,7 @@ fn execute_batch_search( return Err(CollectionError::timeout(timeout, "batch search")); }; - let segment_points = read_segment.available_point_count(); + let segment_points = read_segment.available_point_count_without_deferred(); let segment_config = read_segment.config(); let top = if use_sampling { diff --git a/lib/collection/src/shards/local_shard/indexed_only.rs b/lib/collection/src/shards/local_shard/indexed_only.rs index 4c3afc7733..36a3bf9e15 100644 --- a/lib/collection/src/shards/local_shard/indexed_only.rs +++ b/lib/collection/src/shards/local_shard/indexed_only.rs @@ -39,7 +39,7 @@ pub fn get_index_only_excluded_vectors( .filter_map(move |vector_name| { let segment_config = segment_guard.config().vector_data.get(&vector_name)?; - let points = segment_guard.available_point_count(); + let points = segment_guard.available_point_count_without_deferred(); // Skip segments that have an index. if segment_config.index.is_indexed() { diff --git a/lib/collection/src/shards/local_shard/scroll.rs b/lib/collection/src/shards/local_shard/scroll.rs index 966c637d53..ce7b4cf6ab 100644 --- a/lib/collection/src/shards/local_shard/scroll.rs +++ b/lib/collection/src/shards/local_shard/scroll.rs @@ -382,7 +382,7 @@ impl LocalShard { let read_segment = get_segment.read(); Ok(( - read_segment.available_point_count(), + read_segment.available_point_count_without_deferred(), read_segment.read_random_filtered( limit, filter.as_ref(), diff --git a/lib/collection/src/update_workers/update_worker.rs b/lib/collection/src/update_workers/update_worker.rs index 5ecca59c69..b1c13a3494 100644 --- a/lib/collection/src/update_workers/update_worker.rs +++ b/lib/collection/src/update_workers/update_worker.rs @@ -229,7 +229,7 @@ impl UpdateWorkers { let segments = locked_segments.read(); segments.iter().any(|(_, segment)| { let segment_guard = segment.get().read(); - segment_guard.deferred_points_count() > 0 + segment_guard.has_deferred_points() }) })) .await diff --git a/lib/edge/src/scroll.rs b/lib/edge/src/scroll.rs index 152081b372..eda0602269 100644 --- a/lib/edge/src/scroll.rs +++ b/lib/edge/src/scroll.rs @@ -253,7 +253,7 @@ impl EdgeShard { let segment = segment.get(); let segment = segment.read(); - let point_count = segment.available_point_count(); + let point_count = segment.available_point_count_without_deferred(); let point_ids = segment.read_random_filtered( limit, filter, diff --git a/lib/segment/src/data_types/query_context.rs b/lib/segment/src/data_types/query_context.rs index f6efd7ce1e..a3f5218d8b 100644 --- a/lib/segment/src/data_types/query_context.rs +++ b/lib/segment/src/data_types/query_context.rs @@ -26,7 +26,7 @@ pub struct QueryIdfStats { #[derive(Debug)] pub struct QueryContext { - /// Total amount of available points in the segment. + /// Total amount of available (and visible) points in the segment. available_point_count: usize, /// Parameter, which defines how big a plain segment can be to be considered @@ -70,6 +70,7 @@ impl QueryContext { self } + /// Returns the amount of available (and visible) points. pub fn available_point_count(&self) -> usize { self.available_point_count } diff --git a/lib/segment/src/entry/entry_point.rs b/lib/segment/src/entry/entry_point.rs index 53498a588d..47b11f5d4b 100644 --- a/lib/segment/src/entry/entry_point.rs +++ b/lib/segment/src/entry/entry_point.rs @@ -191,11 +191,15 @@ pub trait NonAppendableSegmentEntry: SnapshotEntry { /// Number of available points /// /// - excludes soft deleted points + /// - includes deferred points. fn available_point_count(&self) -> usize; /// Number of deleted points fn deleted_point_count(&self) -> usize; + /// Similar to `available_point_count()` but excludes all deferred points. + fn available_point_count_without_deferred(&self) -> usize; + /// Size of all available vectors in storage fn available_vectors_size_in_bytes(&self, vector_name: &VectorName) -> OperationResult; @@ -346,6 +350,12 @@ pub trait NonAppendableSegmentEntry: SnapshotEntry { /// Returns external IDs of all deferred points in the segment fn deferred_point_ids(&self) -> Vec; + + /// Returns `true` if there is at least one point that is hidden (deferred). + /// Non-appendable segments always return `false` as they can't have deferred points. + /// + /// Note: the deferred point can be deleted and this function would still return `true`. + fn has_deferred_points(&self) -> bool; } /// Define mutable operations which can be performed with Segment or Segment-like entity. @@ -407,6 +417,4 @@ pub trait SegmentEntry: NonAppendableSegmentEntry { point_id: PointIdType, hw_counter: &HardwareCounterCell, ) -> OperationResult; - - fn deferred_points_count(&self) -> usize; } diff --git a/lib/segment/src/segment/entry.rs b/lib/segment/src/segment/entry.rs index 9c6f2924f5..f24a17b096 100644 --- a/lib/segment/src/segment/entry.rs +++ b/lib/segment/src/segment/entry.rs @@ -77,7 +77,7 @@ impl NonAppendableSegmentEntry for Segment { .get(vector_name) .ok_or_else(|| OperationError::vector_name_not_exists(vector_name))?; let vector_query_context = - query_context.get_vector_context(vector_name, self.deferred_internal_id); + query_context.get_vector_context(vector_name, self.deferred_internal_id()); let internal_results = vector_data.vector_index.borrow().search( query_vectors, filter, @@ -183,7 +183,7 @@ impl NonAppendableSegmentEntry for Segment { // Filter out deferred points. This is done in two stages to prevent cloning `point_ids` and iterating more that needed // but still satisfy rusts ownership constraints. let behavior_allows_filtering = !deferred_behavior.include_all_points(); - let filter_deferred = self.deferred_points_count() > 0 && behavior_allows_filtering; + let filter_deferred = self.has_deferred_points() && behavior_allows_filtering; let filtered_point_ids = filter_deferred.then(|| { point_ids .iter() @@ -415,7 +415,7 @@ impl NonAppendableSegmentEntry for Segment { ) -> OperationResult { Ok(match filter { None => { - let available = self.non_deferred_point_count_estimated(); + let available = self.available_point_count_without_deferred(); CardinalityEstimation { primary_clauses: vec![], min: available, @@ -428,7 +428,7 @@ impl NonAppendableSegmentEntry for Segment { let cardinality = payload_index.estimate_cardinality(filter, hw_counter); let total_points = self.id_tracker.borrow().available_point_count(); - let available_points = self.non_deferred_point_count_estimated(); + let available_points = self.available_point_count_without_deferred(); adjust_for_deferred_points(cardinality, available_points, total_points) } }) @@ -507,9 +507,7 @@ impl NonAppendableSegmentEntry for Segment { 0 }; - let num_points = self.available_point_count(); - - let vectors_size_bytes = total_average_vectors_size_bytes * num_points; + let vectors_size_bytes = total_average_vectors_size_bytes * self.available_point_count(); // Unwrap and default to 0 here because the RocksDB storage is the only faillible one, and we will remove it eventually. let payloads_size_bytes = self @@ -523,7 +521,9 @@ impl NonAppendableSegmentEntry for Segment { segment_type: self.segment_type, num_vectors, num_indexed_vectors, - num_points: self.non_deferred_point_count_estimated(), + num_points: self.available_point_count(), + num_deferred_points: Some(self.deferred_point_count()), + num_deleted_deferred_points: Some(self.deferred_deleted_count().unwrap_or_default()), num_deleted_vectors: self.deleted_point_count(), vectors_size_bytes, // Considers vector storage, but not indices payloads_size_bytes, // Considers payload storage, but not indices @@ -532,7 +532,7 @@ impl NonAppendableSegmentEntry for Segment { is_appendable: self.appendable_flag, index_schema: HashMap::new(), vector_data: vector_data_info, - deferred_internal_id: self.deferred_internal_id, + deferred_internal_id: self.deferred_internal_id(), } } @@ -891,7 +891,7 @@ impl NonAppendableSegmentEntry for Segment { } fn fill_query_context(&self, query_context: &mut QueryContext) { - query_context.add_available_point_count(self.available_point_count()); + query_context.add_available_point_count(self.available_point_count_without_deferred()); let hw_acc = query_context.hardware_usage_accumulator(); let hw_counter = hw_acc.get_counter_cell(); @@ -918,7 +918,7 @@ impl NonAppendableSegmentEntry for Segment { } fn point_is_deferred(&self, point_id: PointIdType) -> bool { - if let Some(deferred_from) = self.deferred_internal_id + if let Some(deferred_from) = self.deferred_internal_id() && let Some(internal_id) = self.id_tracker.borrow().internal_id(point_id) { return self.is_appendable() && internal_id >= deferred_from; @@ -927,9 +927,13 @@ impl NonAppendableSegmentEntry for Segment { } fn deferred_point_ids(&self) -> Vec { - let Some(deferred_from) = self.deferred_internal_id else { + let Some(deferred_from) = self.deferred_internal_id() else { return vec![]; }; + if self.deferred_point_count() == 0 { + return vec![]; + } + let id_tracker = self.id_tracker.borrow(); id_tracker .iter_internal() @@ -937,6 +941,18 @@ impl NonAppendableSegmentEntry for Segment { .filter_map(|internal_id| id_tracker.external_id(internal_id)) .collect() } + + fn available_point_count_without_deferred(&self) -> usize { + self.id_tracker + .borrow() + .available_point_count() + .saturating_sub(self.deferred_point_count()) + } + + fn has_deferred_points(&self) -> bool { + self.deferred_internal_id() + .is_some_and(|deferred_from| self.total_point_count() > deferred_from as usize) + } } impl SegmentEntry for Segment { @@ -1117,15 +1133,4 @@ impl SegmentEntry for Segment { }), }) } - - fn deferred_points_count(&self) -> usize { - if let Some(deferred_from) = self.deferred_internal_id - && self.is_appendable() - { - return self - .total_point_count() - .saturating_sub(deferred_from as usize); - } - 0 - } } diff --git a/lib/segment/src/segment/facet.rs b/lib/segment/src/segment/facet.rs index 7136d78adc..a7f66fdf24 100644 --- a/lib/segment/src/segment/facet.rs +++ b/lib/segment/src/segment/facet.rs @@ -58,7 +58,7 @@ impl Segment { &filter_cardinality, hw_counter, is_stopped, - self.deferred_internal_id, + self.deferred_internal_id(), ) .filter(|point_id| !id_tracker.is_deleted_point(*point_id)) .fold(HashMap::new(), |mut map, point_id| { @@ -98,7 +98,8 @@ impl Segment { let count = iter .dedup() .take_while(|&point_id| { - point_id < self.deferred_internal_id.unwrap_or(PointOffsetType::MAX) + point_id + < self.deferred_internal_id().unwrap_or(PointOffsetType::MAX) }) .filter(|&point_id| context.check(point_id)) .count(); @@ -112,7 +113,7 @@ impl Segment { } else { // just count how many points each value has let iter = facet_index - .iter_counts_per_value(self.deferred_internal_id) + .iter_counts_per_value(self.deferred_internal_id()) .stop_if(is_stopped) .filter(|hit| hit.count > 0); @@ -152,7 +153,7 @@ impl Segment { &filter_cardinality, hw_counter, is_stopped, - self.deferred_internal_id, + self.deferred_internal_id(), ) .filter(|point_id| !id_tracker.is_deleted_point(*point_id)) .fold(BTreeSet::new(), |mut set, point_id| { @@ -164,7 +165,7 @@ impl Segment { .collect() } else { facet_index - .iter_values(hw_counter, self.deferred_internal_id) + .iter_values(hw_counter, self.deferred_internal_id()) .stop_if(is_stopped) .map(|value_ref| value_ref.to_owned()) .collect() diff --git a/lib/segment/src/segment/mod.rs b/lib/segment/src/segment/mod.rs index 91feec3dba..713f288317 100644 --- a/lib/segment/src/segment/mod.rs +++ b/lib/segment/src/segment/mod.rs @@ -94,9 +94,18 @@ pub struct Segment { pub error_status: Option, #[cfg(feature = "rocksdb")] pub database: Option>>, + pub(crate) deferred_point_status: Option, +} + +#[derive(Debug)] +pub struct DeferredPointStatus { /// Points with internal id >= this value are hidden from reads. /// Available for appendable segments only. - pub(crate) deferred_internal_id: Option, + pub(crate) deferred_internal_id: PointOffsetType, + + /// Amount of deleted deferred points. Must kept track of properly to be able + /// to calculate the amount of available deferred and visible points. + pub(crate) deferred_deleted_count: usize, } pub struct VectorData { diff --git a/lib/segment/src/segment/order_by.rs b/lib/segment/src/segment/order_by.rs index e4c1e0031a..2cb7cc5b25 100644 --- a/lib/segment/src/segment/order_by.rs +++ b/lib/segment/src/segment/order_by.rs @@ -40,7 +40,7 @@ impl Segment { let start_from = order_by.start_from(); - let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id); + let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id()); let values_ids_iterator = payload_index .iter_filtered_points( @@ -114,7 +114,7 @@ impl Segment { // We can't early stop the iterator for deferred points because the items are sorted lexicographically by type `(T, internalID)`. .filter(|&(_, internal_id)| { deferred_behavior.include_all_points() - || internal_id < self.deferred_internal_id.unwrap_or(PointOffsetType::MAX) + || internal_id < self.deferred_internal_id().unwrap_or(PointOffsetType::MAX) }); let directed_range_iter = match order_by.direction() { diff --git a/lib/segment/src/segment/sampling.rs b/lib/segment/src/segment/sampling.rs index 7dfcb96ec3..c54ee8364a 100644 --- a/lib/segment/src/segment/sampling.rs +++ b/lib/segment/src/segment/sampling.rs @@ -28,7 +28,7 @@ impl Segment { &cardinality_estimation, hw_counter, is_stopped, - self.deferred_internal_id, + self.deferred_internal_id(), ) .filter_map(|internal_id| id_tracker.external_id(internal_id)); @@ -50,7 +50,7 @@ impl Segment { let filter_context = payload_index.filter_context(condition, hw_counter); self.id_tracker .borrow() - .iter_random_visible(self.deferred_internal_id) + .iter_random_visible(self.deferred_internal_id()) .stop_if(is_stopped) .filter(move |(_, internal_id)| filter_context.check(*internal_id)) .map(|(external_id, _)| external_id) @@ -61,7 +61,7 @@ impl Segment { pub(super) fn read_by_random_id(&self, limit: usize) -> Vec { self.id_tracker .borrow() - .iter_random_visible(self.deferred_internal_id) + .iter_random_visible(self.deferred_internal_id()) .map(|x| x.0) .take(limit) .collect() diff --git a/lib/segment/src/segment/scroll.rs b/lib/segment/src/segment/scroll.rs index ed00f63140..8bde3213d7 100644 --- a/lib/segment/src/segment/scroll.rs +++ b/lib/segment/src/segment/scroll.rs @@ -69,7 +69,7 @@ impl Segment { hw_counter: &HardwareCounterCell, deferred_behavior: DeferredBehavior, ) -> Vec { - let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id); + let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id()); let payload_index = self.payload_index.borrow(); let filter_context = payload_index.filter_context(condition, hw_counter); @@ -89,7 +89,7 @@ impl Segment { limit: Option, deferred_behavior: DeferredBehavior, ) -> Vec { - let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id); + let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id()); self.id_tracker .borrow() @@ -108,7 +108,7 @@ impl Segment { hw_counter: &HardwareCounterCell, deferred_behavior: DeferredBehavior, ) -> Vec { - let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id); + let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id()); let payload_index = self.payload_index.borrow(); let id_tracker = self.id_tracker.borrow(); diff --git a/lib/segment/src/segment/segment_ops.rs b/lib/segment/src/segment/segment_ops.rs index b7c5189ff7..fa8fc8d181 100644 --- a/lib/segment/src/segment/segment_ops.rs +++ b/lib/segment/src/segment/segment_ops.rs @@ -309,7 +309,22 @@ impl Segment { .borrow_mut() .clear_payload(internal_id, hw_counter)?; - self.id_tracker.borrow_mut().drop_internal(internal_id)?; + let mut id_tracker = self.id_tracker.borrow_mut(); + + let is_point_already_deleted = id_tracker.is_deleted_point(internal_id); + + id_tracker.drop_internal(internal_id)?; + + let deferred_point_status = self.deferred_point_status.as_mut(); + + // Increase counter for deleted points. + if let Some(deferred_point_status) = deferred_point_status + && internal_id >= deferred_point_status.deferred_internal_id + // Don't count the deletion of the same point twice + && !is_point_already_deleted + { + deferred_point_status.deferred_deleted_count += 1; + } // Before, we propagated point deletions to also delete its vectors. This turns // out to be problematic because this sometimes makes us lose vector data @@ -653,27 +668,46 @@ impl Segment { self.id_tracker.borrow_mut().fix_inconsistencies() } - /// Returns the (estimated) amount of deferred points. - /// - /// This value is an estimation because it does not account for deferred points - /// that have been deleted before becoming visible. - pub fn deferred_point_count_estimated(&self) -> usize { - match self.deferred_internal_id { - Some(internal_id) => { - let id_tracker = self.id_tracker.borrow(); - let max_id = id_tracker.total_point_count(); - max_id.saturating_sub(internal_id as usize) - } + /// Returns the amount of non-deleted deferred points. + pub fn deferred_point_count(&self) -> usize { + match self.deferred_internal_id() { + Some(internal_id) => self + .id_tracker + .borrow() + .total_point_count() + .saturating_sub(internal_id as usize) + .saturating_sub(self.deferred_deleted_count().unwrap_or_default()), None => 0, } } - /// Returns the amount of points that are not deferred. - pub fn non_deferred_point_count_estimated(&self) -> usize { - self.id_tracker - .borrow() - .available_point_count() - .saturating_sub(self.deferred_point_count_estimated()) + /// Calculates the amount of deleted deferred points by iterating over all points in the ID tracker. Therefore this operation + /// can be expensive and should only be run once at segment creation. + pub(crate) fn calculate_deleted_deferred_point_count(&self) -> usize { + let Some(deferred_from) = self.deferred_internal_id() else { + return 0; + }; + + let id_tracker = self.id_tracker.borrow(); + let total_points = id_tracker.total_point_count(); + + if total_points < deferred_from as usize { + return 0; + } + + id_tracker.deleted_point_bitslice()[deferred_from as usize..total_points].count_ones() + } + + pub(crate) fn deferred_internal_id(&self) -> Option { + self.deferred_point_status + .as_ref() + .map(|i| i.deferred_internal_id) + } + + pub(crate) fn deferred_deleted_count(&self) -> Option { + self.deferred_point_status + .as_ref() + .map(|i| i.deferred_deleted_count) } } diff --git a/lib/segment/src/segment/tests.rs b/lib/segment/src/segment/tests.rs index 8791031cb2..1885e48b96 100644 --- a/lib/segment/src/segment/tests.rs +++ b/lib/segment/src/segment/tests.rs @@ -806,8 +806,7 @@ fn create_deferred_segment( .unwrap(); // Initially, no deferred points (empty segment) - - assert_eq!(segment.deferred_points_count(), 0); + assert!(!segment.has_deferred_points()); for i in 0..n_vectors { assert!( !segment.point_is_deferred(PointIdType::from(i as u64)), @@ -896,9 +895,9 @@ fn create_deferred_segment( .unwrap(); // Now we should have deferred points + assert_eq!(segment.has_deferred_points(), n_deferred > 0); if n_deferred > 0 { - assert_eq!(segment.deferred_points_count(), n_deferred); - assert_eq!(segment.deferred_internal_id, Some(n_vectors as u32)); + assert_eq!(segment.deferred_internal_id(), Some(n_vectors as u32)); } // Points 1 to n_vectors should NOT be deferred @@ -941,8 +940,8 @@ fn create_deferred_segment( ); // Test deferred point count estimation. - assert_eq!(segment.deferred_point_count_estimated(), n_deferred); - assert_eq!(segment.non_deferred_point_count_estimated(), n_vectors); + assert_eq!(segment.deferred_point_count(), n_deferred); + assert_eq!(segment.available_point_count_without_deferred(), n_vectors); segment } @@ -991,11 +990,11 @@ fn test_dense_deferred_points() { // Deferred points should still be the same after reopening assert!( - segment.deferred_points_count() > 0, + segment.has_deferred_points(), "Segment should still have deferred points after reopening" ); assert_eq!( - segment.deferred_internal_id, + segment.deferred_internal_id(), Some(13), "Deferred internal ID should still be `DEFERRED_POINTS_ID` after reopening" ); @@ -1042,7 +1041,7 @@ fn test_deferred_point_estimation_with_filter() { // For consistency we also test that the same cardinality is estimated if no deferred points exist. if n_deferred == 0 { - assert_eq!(segment.deferred_internal_id, None); + assert_eq!(segment.deferred_internal_id(), None); let estimation = segment .estimate_point_count(Some(&filter), &hw_counter) .unwrap(); @@ -1273,14 +1272,14 @@ fn test_deferred_point_facets() { .facet(&request, &AtomicBool::new(false), &hw_counter) .unwrap(); - let old_deferred_id = segment.deferred_internal_id.take(); + let old_status = segment.deferred_point_status.take(); if n_deferred > 0 { - assert!(old_deferred_id.is_some()); + assert!(old_status.is_some()); } let facet_res = segment .facet(&request, &AtomicBool::new(false), &hw_counter) .unwrap(); - segment.deferred_internal_id = old_deferred_id; + segment.deferred_point_status = old_status; let expected_deferred = if filter.is_some() { n_deferred.div_ceil(3) @@ -1413,7 +1412,7 @@ fn assert_deferred_points_excluded( } // Disable deferred points and search again. - segment.deferred_internal_id = None; + segment.deferred_point_status = None; let search_res_normal = operation(&segment, filter_set.filter.as_ref()); assert_eq!( search_res_normal.len(), @@ -1425,3 +1424,60 @@ fn assert_deferred_points_excluded( } } } + +#[test] +fn test_deleted_deferred_point_count() { + let hw_counter = HardwareCounterCell::new(); + + for n_deferred in [0, 1, 10, 300] { + let dir = Builder::new().prefix("segment_dir").tempdir().unwrap(); + let mut segment = create_deferred_segment(&dir, 5, N_POINTS, n_deferred); + + assert_eq!(segment.deferred_point_count(), n_deferred); + + if n_deferred == 0 { + continue; + } + + assert_eq!(segment.available_point_count_without_deferred(), N_POINTS); + + for d in 0..n_deferred { + let delete_id = segment.deferred_internal_id().unwrap() + d as u32; + segment + .delete_point_internal(delete_id, &hw_counter) + .unwrap(); + + let deleted_count = d + 1; // The first index is 0 but this point is deleted, so count must be 1. + assert_eq!( + segment.deferred_point_count(), + n_deferred.checked_sub(deleted_count).unwrap() + ); + assert_eq!( + segment.calculate_deleted_deferred_point_count(), + deleted_count, + ); + + // Do the operation twice to test that we don't double count the same point. + segment + .delete_point_internal(delete_id, &hw_counter) + .unwrap(); + + assert_eq!( + segment.deferred_point_count(), + n_deferred.checked_sub(deleted_count).unwrap() + ); + + assert_eq!( + segment.calculate_deleted_deferred_point_count(), + deleted_count + ); + + assert_eq!(segment.available_point_count_without_deferred(), N_POINTS); + } + + // We delete all deferred points in the segment. + assert_eq!(segment.deferred_point_count(), 0); + assert_eq!(segment.calculate_deleted_deferred_point_count(), n_deferred); + assert_eq!(segment.available_point_count_without_deferred(), N_POINTS); + } +} diff --git a/lib/segment/src/segment_constructor/segment_constructor_base.rs b/lib/segment/src/segment_constructor/segment_constructor_base.rs index 64d28924a6..a2567b5284 100644 --- a/lib/segment/src/segment_constructor/segment_constructor_base.rs +++ b/lib/segment/src/segment_constructor/segment_constructor_base.rs @@ -52,7 +52,9 @@ use crate::payload_storage::on_disk_payload_storage::OnDiskPayloadStorage; use crate::payload_storage::payload_storage_enum::PayloadStorageEnum; #[cfg(feature = "rocksdb")] use crate::payload_storage::simple_payload_storage::SimplePayloadStorage; -use crate::segment::{SEGMENT_STATE_FILE, Segment, SegmentVersion, VectorData}; +use crate::segment::{ + DeferredPointStatus, SEGMENT_STATE_FILE, Segment, SegmentVersion, VectorData, +}; #[cfg(feature = "rocksdb")] use crate::types::MultiVectorConfig; use crate::types::{ @@ -684,12 +686,18 @@ fn create_segment( error_status: None, #[cfg(feature = "rocksdb")] database: db_builder.build(), - deferred_internal_id: None, + deferred_point_status: None, }; - if segment.is_appendable() { - segment.deferred_internal_id = deferred_internal_id; + if let Some(deferred_internal_id) = deferred_internal_id + && segment.is_appendable() + { + segment.deferred_point_status = Some(DeferredPointStatus { + deferred_internal_id, + deferred_deleted_count: segment.calculate_deleted_deferred_point_count(), + }); } + Ok(segment) } diff --git a/lib/segment/src/types.rs b/lib/segment/src/types.rs index ff4b62258e..23b62b071d 100644 --- a/lib/segment/src/types.rs +++ b/lib/segment/src/types.rs @@ -472,6 +472,8 @@ pub struct SegmentInfo { pub segment_type: SegmentType, pub num_vectors: usize, pub num_points: usize, + pub num_deferred_points: Option, + pub num_deleted_deferred_points: Option, pub num_indexed_vectors: usize, pub num_deleted_vectors: usize, /// An ESTIMATION of effective amount of bytes used for vectors diff --git a/lib/shard/src/optimizers/indexing_optimizer.rs b/lib/shard/src/optimizers/indexing_optimizer.rs index c20566e8d3..be802aaae3 100644 --- a/lib/shard/src/optimizers/indexing_optimizer.rs +++ b/lib/shard/src/optimizers/indexing_optimizer.rs @@ -4,7 +4,7 @@ use std::sync::Arc; use parking_lot::Mutex; use segment::common::operation_time_statistics::OperationDurationsAggregator; -use segment::entry::{NonAppendableSegmentEntry as _, SegmentEntry}; +use segment::entry::NonAppendableSegmentEntry as _; use segment::segment::Segment; use segment::types::HnswGlobalConfig; @@ -62,7 +62,7 @@ impl IndexingOptimizer { .memmap_threshold_kb .saturating_mul(BYTES_IN_KB); - let has_deferred_points = segment.deferred_points_count() > 0; + let has_deferred_points = segment.has_deferred_points(); for (vector_name, vector_cfg) in &self.segment_optimizer_config.dense_vector { if let Some(vector_data) = segment_data_config.vector_data.get(vector_name) { diff --git a/lib/shard/src/proxy_segment/mod.rs b/lib/shard/src/proxy_segment/mod.rs index 61ca197102..240e2a1671 100644 --- a/lib/shard/src/proxy_segment/mod.rs +++ b/lib/shard/src/proxy_segment/mod.rs @@ -32,6 +32,7 @@ pub struct ProxySegment { /// May contain points which are not in wrapped_segment, /// because the set is shared among all proxy segments deleted_points: DeletedPoints, + deleted_deferred_count: usize, wrapped_config: SegmentConfig, /// Version of the last change in this proxy, considering point deletes and payload index @@ -63,6 +64,7 @@ impl ProxySegment { deleted_mask, changed_indexes: ProxyIndexChanges::default(), deleted_points: AHashMap::new(), + deleted_deferred_count: 0, wrapped_config, version, } @@ -230,6 +232,7 @@ impl ProxySegment { OperationResult::Ok(()) })?; self.deleted_points.clear(); + self.deleted_deferred_count = 0; // Note: We do not clear the deleted mask here, as it provides // no performance advantage and does not affect the correctness of search. diff --git a/lib/shard/src/proxy_segment/segment_entry.rs b/lib/shard/src/proxy_segment/segment_entry.rs index 9403686ff6..80251b37df 100644 --- a/lib/shard/src/proxy_segment/segment_entry.rs +++ b/lib/shard/src/proxy_segment/segment_entry.rs @@ -407,6 +407,22 @@ impl NonAppendableSegmentEntry for ProxySegment { wrapped_segment_count.saturating_sub(deleted_points_count) } + fn available_point_count_without_deferred(&self) -> usize { + let wrapped_segment_visible_count = self + .wrapped_segment + .get() + .read() + .available_point_count_without_deferred(); + + // Amount of visible points that are deleted in this proxy. + let deleted_visible = self + .deleted_points + .len() + .saturating_sub(self.deleted_deferred_count); + + wrapped_segment_visible_count.saturating_sub(deleted_visible) + } + fn deleted_point_count(&self) -> usize { self.wrapped_segment.get().read().deleted_point_count() + self.deleted_points.len() } @@ -437,14 +453,17 @@ impl NonAppendableSegmentEntry for ProxySegment { filter: Option<&'a Filter>, hw_counter: &HardwareCounterCell, ) -> OperationResult { - let deleted_point_count = self.deleted_points.len(); + let deleted_point_count = self + .deleted_points + .len() + .saturating_sub(self.deleted_deferred_count); let (wrapped_segment_est, total_wrapped_size) = { let wrapped_segment = self.wrapped_segment.get(); let wrapped_segment_guard = wrapped_segment.read(); ( wrapped_segment_guard.estimate_point_count(filter, hw_counter)?, - wrapped_segment_guard.available_point_count(), + wrapped_segment_guard.available_point_count_without_deferred(), ) }; @@ -488,6 +507,7 @@ impl NonAppendableSegmentEntry for ProxySegment { let vector_name_count = self.config().vector_data.len() + self.config().sparse_vector_data.len(); + let deleted_points_count = self.deleted_points.len(); // This is a best estimate @@ -511,6 +531,14 @@ impl NonAppendableSegmentEntry for ProxySegment { num_vectors, num_indexed_vectors, num_points: self.available_point_count(), + num_deferred_points: wrapped_info.num_deferred_points.map(|num_deferred_points| { + num_deferred_points.saturating_sub(self.deleted_deferred_count) + }), + num_deleted_deferred_points: wrapped_info.num_deleted_deferred_points.map( + |num_deleted_deferred_points| { + num_deleted_deferred_points.saturating_add(self.deleted_deferred_count) + }, + ), num_deleted_vectors: wrapped_info.num_deleted_vectors + deleted_points_count * vector_name_count, vectors_size_bytes: wrapped_info.vectors_size_bytes, // + write_info.vectors_size_bytes, @@ -624,12 +652,21 @@ impl NonAppendableSegmentEntry for ProxySegment { _hw_counter: &HardwareCounterCell, ) -> OperationResult { let mut was_deleted = false; + let was_deferred_point; self.version = cmp::max(self.version, op_num); let point_offset = match &self.wrapped_segment { LockedSegment::Original(raw_segment) => { - let point_offset = raw_segment.read().get_internal_id(point_id); + let (point_offset, is_deferred) = { + let read_segment = raw_segment.read(); + ( + read_segment.get_internal_id(point_id), + read_segment.point_is_deferred(point_id), + ) + }; + was_deferred_point = is_deferred; + if point_offset.is_some() { let prev = self.deleted_points.insert( point_id, @@ -649,7 +686,16 @@ impl NonAppendableSegmentEntry for ProxySegment { point_offset } LockedSegment::Proxy(proxy) => { - if proxy.read().has_point(point_id) { + let (has_point, is_deferred) = { + let read_proxy = proxy.read(); + ( + read_proxy.has_point(point_id), + read_proxy.point_is_deferred(point_id), + ) + }; + was_deferred_point = is_deferred; + + if has_point { let prev = self.deleted_points.insert( point_id, ProxyDeletedPoint { @@ -671,6 +717,11 @@ impl NonAppendableSegmentEntry for ProxySegment { self.set_deleted_offset(point_offset); + // Increase delete counter for deferred point. + if was_deleted && was_deferred_point { + self.deleted_deferred_count += 1; + } + Ok(was_deleted) } @@ -729,9 +780,15 @@ impl NonAppendableSegmentEntry for ProxySegment { fn deferred_point_ids(&self) -> Vec { let mut ids = self.wrapped_segment.get().read().deferred_point_ids(); - ids.retain(|point_id| !self.deleted_points.contains_key(point_id)); + if self.deleted_deferred_count > 0 { + ids.retain(|point_id| !self.deleted_points.contains_key(point_id)); + } ids } + + fn has_deferred_points(&self) -> bool { + self.wrapped_segment.get().read().has_deferred_points() + } } impl SegmentEntry for ProxySegment { @@ -818,8 +875,4 @@ impl SegmentEntry for ProxySegment { "Clear payload is disabled for proxy segments: operation {op_num} on point {point_id}", ))) } - - fn deferred_points_count(&self) -> usize { - self.wrapped_segment.get().read().deferred_points_count() - } } diff --git a/lib/shard/src/proxy_segment/tests.rs b/lib/shard/src/proxy_segment/tests.rs index 9bd9e2bdff..bb0e62e2c4 100644 --- a/lib/shard/src/proxy_segment/tests.rs +++ b/lib/shard/src/proxy_segment/tests.rs @@ -532,3 +532,67 @@ fn test_proxy_segment_flush() { // So we have to keep WAL for deleted points. assert!(version_after_delete > flushed_version_2); } + +#[test] +fn test_proxy_deferred() { + let hw_counter = HardwareCounterCell::new(); + + let tmp_dir = tempfile::Builder::new() + .prefix("segment_dir") + .tempdir() + .unwrap(); + + let mut wrapped_segment = build_segment_with_deferred_1(tmp_dir.path()); + + let initial_estimation = wrapped_segment.estimate_point_count(None, &hw_counter); + + let initial_deferred_point_count = wrapped_segment.size_info().num_deferred_points.unwrap(); + + wrapped_segment + .delete_point_internal(3, &hw_counter) + .unwrap(); + + assert_eq!( + wrapped_segment.size_info().num_deferred_points.unwrap(), + initial_deferred_point_count - 1 + ); + + let mut proxy_segment = ProxySegment::new(LockedSegment::new(wrapped_segment)); + + assert_eq!( + proxy_segment.size_info().num_deferred_points.unwrap(), + initial_deferred_point_count - 1 + ); + + assert_eq!(proxy_segment.available_point_count_without_deferred(), 3); + + proxy_segment + .delete_point(7, 5.into(), &hw_counter) + .unwrap(); + + assert_eq!( + proxy_segment.size_info().num_deferred_points.unwrap(), + initial_deferred_point_count - 2 + ); + + assert_eq!(proxy_segment.available_point_count_without_deferred(), 3); + + // We didn't touch normal points so estimation should not change. + assert_eq!( + proxy_segment.estimate_point_count(None, &hw_counter), + initial_estimation + ); + + // Touch normal points + proxy_segment + .delete_point(6, 1.into(), &hw_counter) + .unwrap(); + + // Now we must see a difference in estimation. + assert_ne!( + proxy_segment.estimate_point_count(None, &hw_counter), + initial_estimation + ); + + assert_eq!(proxy_segment.available_point_count_without_deferred(), 2); +} diff --git a/tests/openapi/test_deferred_points.py b/tests/openapi/test_deferred_points.py index 00029c6668..54187a357b 100644 --- a/tests/openapi/test_deferred_points.py +++ b/tests/openapi/test_deferred_points.py @@ -55,6 +55,41 @@ def get_collection_info(): return response.json()['result'] +def get_point_count_excluding_deferred(): + """ + Manually calculate the amount of visible (non-deferred) points for now, + since we don't provide it in CollectionInfo yet. + """ + + response = request_with_validation( + api='/telemetry', + method="GET", + query_params={"details_level": 10} + ) + assert response.ok + collections = response.json()['result']['collections']['collections'] + + num_points = 0 + had_collection = False + for collection in collections: + if collection['id'] != COLLECTION_NAME: + continue + + had_collection = True + + for shard in collection['shards']: + if 'local' not in shard: + continue + + for segment in shard['local']['segments']: + info = segment['info'] + num_points += int(info['num_points']) - int(info['num_deferred_points']) + + assert had_collection, "Collection not found in telemetry!" + + return num_points + + def upsert_points_batch(points, wait=True): """Upsert a list of points in a single request.""" response = request_with_validation( @@ -223,9 +258,10 @@ def test_deferred_points(): visible_ids = {p['id'] for p in scrolled} # Point count from collection info should match scrolled count - info = get_collection_info() - assert info['points_count'] == visible_count, ( - f"points_count ({info['points_count']}) should match scrolled count ({visible_count})" + # info = get_collection_info() + point_count = get_point_count_excluding_deferred() + assert point_count == visible_count, ( + f"points_count ({point_count}) should match scrolled count ({visible_count})" ) # Retrieve all 2000 IDs: only non-deferred ones should be returned @@ -261,9 +297,9 @@ def test_deferred_points(): time.sleep(2) # These new points should also be deferred (added beyond the threshold) - info = get_collection_info() - assert info['points_count'] == visible_count, ( - f"Expected {visible_count} visible points (new points deferred), got {info['points_count']}" + point_count = get_point_count_excluding_deferred() + assert point_count == visible_count, ( + f"Expected {visible_count} visible points (new points deferred), got {point_count}" ) new_point_ids = set(range(new_points_start, new_points_start + new_points_count)) @@ -287,9 +323,9 @@ def test_deferred_points(): # After optimization, ALL points should be visible expected_total = total_upserted + new_points_count + 1 - info = get_collection_info() - assert info['points_count'] == expected_total, ( - f"After optimization, expected {expected_total} points, got {info['points_count']}" + point_count = get_point_count_excluding_deferred() + assert point_count == expected_total, ( + f"After optimization, expected {expected_total} points, got {point_count}" ) # The deferred set_payload should now be visible