mirror of
https://github.com/qdrant/qdrant.git
synced 2026-07-23 11:11:00 -05:00
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 <andrey@vasnetsov.com> * Make new fields in API optional * Don't take range if no deferred point exist * Review remarks --------- Co-authored-by: Andrey Vasnetsov <andrey@vasnetsov.com>
This commit is contained in:
@@ -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",
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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<usize>;
|
||||
|
||||
@@ -346,6 +350,12 @@ pub trait NonAppendableSegmentEntry: SnapshotEntry {
|
||||
|
||||
/// Returns external IDs of all deferred points in the segment
|
||||
fn deferred_point_ids(&self) -> Vec<PointIdType>;
|
||||
|
||||
/// 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<bool>;
|
||||
|
||||
fn deferred_points_count(&self) -> usize;
|
||||
}
|
||||
|
||||
@@ -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<CardinalityEstimation> {
|
||||
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<PointIdType> {
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -94,9 +94,18 @@ pub struct Segment {
|
||||
pub error_status: Option<SegmentFailedState>,
|
||||
#[cfg(feature = "rocksdb")]
|
||||
pub database: Option<Arc<parking_lot::RwLock<DB>>>,
|
||||
pub(crate) deferred_point_status: Option<DeferredPointStatus>,
|
||||
}
|
||||
|
||||
#[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<PointOffsetType>,
|
||||
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 {
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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<PointIdType> {
|
||||
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()
|
||||
|
||||
@@ -69,7 +69,7 @@ impl Segment {
|
||||
hw_counter: &HardwareCounterCell,
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> Vec<PointIdType> {
|
||||
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<usize>,
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> Vec<PointIdType> {
|
||||
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<PointIdType> {
|
||||
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();
|
||||
|
||||
@@ -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<PointOffsetType> {
|
||||
self.deferred_point_status
|
||||
.as_ref()
|
||||
.map(|i| i.deferred_internal_id)
|
||||
}
|
||||
|
||||
pub(crate) fn deferred_deleted_count(&self) -> Option<usize> {
|
||||
self.deferred_point_status
|
||||
.as_ref()
|
||||
.map(|i| i.deferred_deleted_count)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<F, R, T>(
|
||||
}
|
||||
|
||||
// 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<F, R, T>(
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
@@ -472,6 +472,8 @@ pub struct SegmentInfo {
|
||||
pub segment_type: SegmentType,
|
||||
pub num_vectors: usize,
|
||||
pub num_points: usize,
|
||||
pub num_deferred_points: Option<usize>,
|
||||
pub num_deleted_deferred_points: Option<usize>,
|
||||
pub num_indexed_vectors: usize,
|
||||
pub num_deleted_vectors: usize,
|
||||
/// An ESTIMATION of effective amount of bytes used for vectors
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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<CardinalityEstimation> {
|
||||
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<bool> {
|
||||
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<PointIdType> {
|
||||
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()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user