From cf07d4ef10a286ed5b5fdae25245fd4d4faabe51 Mon Sep 17 00:00:00 2001 From: Ivan Pleshkov Date: Mon, 2 Mar 2026 09:57:19 +0100 Subject: [PATCH] Deferred threshold integration (#8246) * Deferred threshold integration * update deferred id * apply update_deferred_internal_id * fix segment inspector * use avaliable bytes count * renamings * review remarks * review remarks * move has_deferred_points * remove todo * update comments --- .../src/collection_manager/fixtures.rs | 2 + .../optimizers/config_mismatch_optimizer.rs | 3 ++ .../optimizers/indexing_optimizer.rs | 4 ++ .../optimizers/merge_optimizer.rs | 1 + .../optimizers/vacuum_optimizer.rs | 2 + lib/collection/src/optimizers_builder.rs | 9 ++++ lib/collection/src/shards/local_shard/mod.rs | 22 ++++++++- .../src/shards/local_shard/snapshot.rs | 15 +++++- .../src/shards/local_shard/snapshot_tests.rs | 1 + lib/collection/src/tests/mod.rs | 1 + .../src/update_workers/optimization_worker.rs | 1 + lib/edge/src/lib.rs | 17 ++++--- lib/segment/benches/multi_vector_search.rs | 2 +- lib/segment/src/index/struct_payload_index.rs | 8 +++- lib/segment/src/segment/mod.rs | 10 ++++ lib/segment/src/segment/segment_ops.rs | 48 +++++++++++++++++++ lib/segment/src/segment/tests.rs | 2 +- .../segment_constructor/segment_builder.rs | 11 ++++- .../segment_constructor_base.rs | 33 +++++++++++-- .../simple_segment_constructor.rs | 3 ++ .../integration/byte_storage_hnsw_test.rs | 2 +- .../byte_storage_quantization_test.rs | 2 +- .../tests/integration/fixtures/segment.rs | 3 ++ .../multivector_filtrable_hnsw_test.rs | 2 +- .../multivector_quantization_test.rs | 2 +- .../tests/integration/payload_index_test.rs | 4 +- .../tests/integration/segment_builder_test.rs | 1 + .../integration/segment_on_disk_snapshot.rs | 2 +- .../tests/integration/segment_tests.rs | 2 +- .../tests/integration/sparse_discover_test.rs | 6 +-- .../sparse_vector_index_search_tests.rs | 6 +-- lib/shard/src/operations/optimization.rs | 3 ++ lib/shard/src/optimize.rs | 7 +++ .../src/optimizers/indexing_optimizer.rs | 4 +- lib/shard/src/optimizers/segment_optimizer.rs | 2 + lib/shard/src/segment_holder/mod.rs | 11 ++++- lib/shard/src/segment_holder/snapshot.rs | 3 ++ lib/shard/src/segment_holder/tests.rs | 11 ++++- src/segment_inspector.rs | 2 +- 39 files changed, 233 insertions(+), 37 deletions(-) diff --git a/lib/collection/src/collection_manager/fixtures.rs b/lib/collection/src/collection_manager/fixtures.rs index b566907e07..79f087e93d 100644 --- a/lib/collection/src/collection_manager/fixtures.rs +++ b/lib/collection/src/collection_manager/fixtures.rs @@ -243,6 +243,7 @@ pub(crate) fn get_merge_optimizer( max_segment_size_kb: 100_000, memmap_threshold_kb: 1_000_000, indexing_threshold_kb: 1_000_000, + deferred_points_threshold_bytes: None, }), segment_path.to_owned(), collection_temp_dir.to_owned(), @@ -271,6 +272,7 @@ pub(crate) fn get_indexing_optimizer( max_segment_size_kb: 100_000, memmap_threshold_kb: 100, indexing_threshold_kb: 100, + deferred_points_threshold_bytes: None, }, segment_path.to_owned(), collection_temp_dir.to_owned(), diff --git a/lib/collection/src/collection_manager/optimizers/config_mismatch_optimizer.rs b/lib/collection/src/collection_manager/optimizers/config_mismatch_optimizer.rs index fd9923b4be..23e7ca0521 100644 --- a/lib/collection/src/collection_manager/optimizers/config_mismatch_optimizer.rs +++ b/lib/collection/src/collection_manager/optimizers/config_mismatch_optimizer.rs @@ -105,6 +105,7 @@ mod tests { max_segment_size_kb: usize::MAX, memmap_threshold_kb: usize::MAX, indexing_threshold_kb: 10, + deferred_points_threshold_bytes: None, }; // Base segment @@ -245,6 +246,7 @@ mod tests { max_segment_size_kb: usize::MAX, memmap_threshold_kb: usize::MAX, indexing_threshold_kb: 10, + deferred_points_threshold_bytes: None, }; // Base segment @@ -410,6 +412,7 @@ mod tests { max_segment_size_kb: usize::MAX, memmap_threshold_kb: usize::MAX, indexing_threshold_kb: 10, + deferred_points_threshold_bytes: None, }; let quantization_config_vector1 = QuantizationConfig::Scalar(segment::types::ScalarQuantization { diff --git a/lib/collection/src/collection_manager/optimizers/indexing_optimizer.rs b/lib/collection/src/collection_manager/optimizers/indexing_optimizer.rs index 9f28fdddd9..d0ff1c38ba 100644 --- a/lib/collection/src/collection_manager/optimizers/indexing_optimizer.rs +++ b/lib/collection/src/collection_manager/optimizers/indexing_optimizer.rs @@ -139,6 +139,7 @@ mod tests { max_segment_size_kb: 300, memmap_threshold_kb: 1000, indexing_threshold_kb: 1000, + deferred_points_threshold_bytes: None, }, segments_dir.path().to_owned(), segments_temp_dir.path().to_owned(), @@ -231,6 +232,7 @@ mod tests { max_segment_size_kb: 300, memmap_threshold_kb: 1000, indexing_threshold_kb: 1000, + deferred_points_threshold_bytes: None, }, segments_dir.path().to_owned(), segments_temp_dir.path().to_owned(), @@ -532,6 +534,7 @@ mod tests { max_segment_size_kb: 1000, memmap_threshold_kb: 1000, indexing_threshold_kb: 10, // Always optimize + deferred_points_threshold_bytes: None, }, segments_dir.path().to_owned(), segments_temp_dir.path().to_owned(), @@ -598,6 +601,7 @@ mod tests { max_segment_size_kb: usize::MAX, memmap_threshold_kb: 10, indexing_threshold_kb: usize::MAX, + deferred_points_threshold_bytes: None, }; let mut collection_params = CollectionParams { vectors: VectorsConfig::Single( diff --git a/lib/collection/src/collection_manager/optimizers/merge_optimizer.rs b/lib/collection/src/collection_manager/optimizers/merge_optimizer.rs index 154052a6bc..f6b7ea2b13 100644 --- a/lib/collection/src/collection_manager/optimizers/merge_optimizer.rs +++ b/lib/collection/src/collection_manager/optimizers/merge_optimizer.rs @@ -85,6 +85,7 @@ mod tests { max_segment_size_kb: 1000, memmap_threshold_kb: 100, indexing_threshold_kb: 50, + deferred_points_threshold_bytes: None, }), segment_path.to_owned(), collection_temp_dir.to_owned(), diff --git a/lib/collection/src/collection_manager/optimizers/vacuum_optimizer.rs b/lib/collection/src/collection_manager/optimizers/vacuum_optimizer.rs index 12e247b9d5..70ea65d13e 100644 --- a/lib/collection/src/collection_manager/optimizers/vacuum_optimizer.rs +++ b/lib/collection/src/collection_manager/optimizers/vacuum_optimizer.rs @@ -176,6 +176,7 @@ mod tests { max_segment_size_kb: 1000000, memmap_threshold_kb: 1000000, indexing_threshold_kb: 1000000, + deferred_points_threshold_bytes: None, }, dir.path().to_owned(), temp_dir.path().to_owned(), @@ -257,6 +258,7 @@ mod tests { max_segment_size_kb: usize::MAX, memmap_threshold_kb: usize::MAX, indexing_threshold_kb: 10, + deferred_points_threshold_bytes: None, }; let collection_params = CollectionParams { vectors: VectorsConfig::Multi(BTreeMap::from([ diff --git a/lib/collection/src/optimizers_builder.rs b/lib/collection/src/optimizers_builder.rs index e1d3c68a7b..9901b0808c 100644 --- a/lib/collection/src/optimizers_builder.rs +++ b/lib/collection/src/optimizers_builder.rs @@ -1,8 +1,10 @@ +use std::num::NonZeroUsize; use std::path::Path; use std::sync::Arc; use fs_err as fs; use schemars::JsonSchema; +use segment::common::BYTES_IN_KB; use segment::common::anonymize::Anonymize; use segment::index::hnsw_index::num_rayon_threads; use segment::types::{HnswConfig, HnswGlobalConfig, QuantizationConfig}; @@ -148,6 +150,7 @@ impl OptimizersConfig { memmap_threshold_kb, indexing_threshold_kb, max_segment_size_kb: self.get_max_segment_size_in_kilobytes(num_indexing_threads), + deferred_points_threshold_bytes: self.get_deferred_points_threshold_bytes(), } } @@ -158,6 +161,12 @@ impl OptimizersConfig { num_indexing_threads.saturating_mul(DEFAULT_MAX_SEGMENT_PER_CPU_KB) } } + + pub fn get_deferred_points_threshold_bytes(&self) -> Option { + (self.prevent_unoptimized == Some(true)) + .then(|| self.get_indexing_threshold_kb().saturating_mul(BYTES_IN_KB)) + .and_then(NonZeroUsize::new) + } } pub fn clear_temp_segments(shard_path: &Path) { diff --git a/lib/collection/src/shards/local_shard/mod.rs b/lib/collection/src/shards/local_shard/mod.rs index 4f111fe909..5a35787407 100644 --- a/lib/collection/src/shards/local_shard/mod.rs +++ b/lib/collection/src/shards/local_shard/mod.rs @@ -351,6 +351,9 @@ impl LocalShard { let wal_path = Self::wal_path(shard_path); let segments_path = Self::segments_path(shard_path); + let deferred_points_threshold_bytes = + effective_optimizers_config.get_deferred_points_threshold_bytes(); + let wal: SerdeWal = SerdeWal::new(&wal_path, (&collection_config_read.wal_config).into()) .map_err(|e| CollectionError::service_error(format!("Wal error: {e}")))?; @@ -406,7 +409,12 @@ impl LocalShard { let Some((segment_path, uuid)) = normalize_segment_dir(&segment_path)? else { return CollectionResult::Ok(None); }; - let mut segment = load_segment(&segment_path, uuid, &AtomicBool::new(false))?; + let mut segment = load_segment( + &segment_path, + uuid, + deferred_points_threshold_bytes, + &AtomicBool::new(false), + )?; segment.check_consistency_and_repair()?; @@ -489,6 +497,7 @@ impl LocalShard { &segments_path, segment_config, payload_index_schema.clone(), + deferred_points_threshold_bytes, )?; } @@ -600,6 +609,8 @@ impl LocalShard { .to_base_vector_data(config.quantization_config.as_ref()); let sparse_vector_params = config.params.to_sparse_vector_data(); let segment_number = config.optimizer_config.get_number_segments(); + let deferred_points_threshold_bytes = + effective_optimizers_config.get_deferred_points_threshold_bytes(); for _sid in 0..segment_number { let path_clone = segments_path.clone(); @@ -610,7 +621,14 @@ impl LocalShard { }; let segment = thread::Builder::new() .name(format!("shard-build-{collection_id}-{id}")) - .spawn(move || build_segment(&path_clone, &segment_config, true)) + .spawn(move || { + build_segment( + &path_clone, + &segment_config, + deferred_points_threshold_bytes, + true, + ) + }) .unwrap(); build_handlers.push(segment); } diff --git a/lib/collection/src/shards/local_shard/snapshot.rs b/lib/collection/src/shards/local_shard/snapshot.rs index 6305548588..8ec3e337d6 100644 --- a/lib/collection/src/shards/local_shard/snapshot.rs +++ b/lib/collection/src/shards/local_shard/snapshot.rs @@ -61,7 +61,15 @@ impl LocalShard { let shard_path = self.path.clone(); let segments_path = Self::segments_path(&self.path); - let segment_config = self.collection_config.read().await.to_base_segment_config(); + let (segment_config, deferred_points_threshold_bytes) = { + let collection_config = self.collection_config.read().await; + ( + collection_config.to_base_segment_config(), + collection_config + .optimizer_config + .get_deferred_points_threshold_bytes(), + ) + }; let applied_seq_path = self.applied_seq_handler.path().to_path_buf(); @@ -89,6 +97,7 @@ impl LocalShard { &segments_path, Some(segment_config), payload_index_schema, + deferred_points_threshold_bytes, &temp_path, &tar.descend(Path::new(SEGMENTS_PATH))?, format, @@ -244,6 +253,7 @@ pub fn snapshot_all_segments( segments_path: &Path, segment_config: Option, payload_index_schema: Arc>, + deferred_points_threshold_bytes: Option, temp_dir: &Path, tar: &tar_ext::BuilderExt, format: SnapshotFormat, @@ -257,6 +267,7 @@ pub fn snapshot_all_segments( segments_path, segment_config, payload_index_schema, + deferred_points_threshold_bytes, |segment| { let read_segment = segment.read(); let request_segment_manifest = if let Some(manifest) = manifest { @@ -307,6 +318,7 @@ pub fn proxy_all_segments_and_apply( segments_path: &Path, segment_config: Option, payload_index_schema: Arc>, + deferred_points_threshold_bytes: Option, mut operation: F, ) -> OperationResult<()> where @@ -322,6 +334,7 @@ where segments_path, segment_config, payload_index_schema, + deferred_points_threshold_bytes, )?; // Flush all pending changes of each segment, now wrapped segments won't change anymore diff --git a/lib/collection/src/shards/local_shard/snapshot_tests.rs b/lib/collection/src/shards/local_shard/snapshot_tests.rs index cd6af935a1..0fc167c7fe 100644 --- a/lib/collection/src/shards/local_shard/snapshot_tests.rs +++ b/lib/collection/src/shards/local_shard/snapshot_tests.rs @@ -47,6 +47,7 @@ fn test_snapshot_all() { segments_dir.path(), None, schema, + None, temp_dir.path(), &tar, SnapshotFormat::Regular, diff --git a/lib/collection/src/tests/mod.rs b/lib/collection/src/tests/mod.rs index b062a16a98..5def8c8414 100644 --- a/lib/collection/src/tests/mod.rs +++ b/lib/collection/src/tests/mod.rs @@ -242,6 +242,7 @@ async fn test_new_segment_when_all_over_capacity() { max_segment_size_kb: 1, memmap_threshold_kb: 1_000_000, indexing_threshold_kb: 1_000_000, + deferred_points_threshold_bytes: None, }; let hnsw_config = Default::default(); let segment_config = diff --git a/lib/collection/src/update_workers/optimization_worker.rs b/lib/collection/src/update_workers/optimization_worker.rs index a9d15de3f2..c16fe9d5c3 100644 --- a/lib/collection/src/update_workers/optimization_worker.rs +++ b/lib/collection/src/update_workers/optimization_worker.rs @@ -471,6 +471,7 @@ impl UpdateWorkers { segments_path, Some(segment_config.base_segment_config()), payload_index_schema, + thresholds_config.deferred_points_threshold_bytes, true, )?; let mut write_guard = parking_lot::RwLockUpgradableReadGuard::upgrade(segments_guard); diff --git a/lib/edge/src/lib.rs b/lib/edge/src/lib.rs index 892b012706..78b554e239 100644 --- a/lib/edge/src/lib.rs +++ b/lib/edge/src/lib.rs @@ -104,13 +104,15 @@ impl EdgeShard { continue; }; - let mut segment = load_segment(&segment_path, segment_uuid, &AtomicBool::new(false)) - .map_err(|err| { - OperationError::service_error(format!( - "failed to load segment {}: {err}", - segment_path.display(), - )) - })?; + let mut segment = + load_segment(&segment_path, segment_uuid, None, &AtomicBool::new(false)).map_err( + |err| { + OperationError::service_error(format!( + "failed to load segment {}: {err}", + segment_path.display(), + )) + }, + )?; if let Some(config) = &config { if !config.is_compatible(segment.config()) { @@ -156,6 +158,7 @@ impl EdgeShard { &segments_path, config.clone(), Arc::new(payload_index_schema), + None, )?; debug_assert!(segments.has_appendable_segment()); diff --git a/lib/segment/benches/multi_vector_search.rs b/lib/segment/benches/multi_vector_search.rs index 3a7da200db..9d01ebe7df 100644 --- a/lib/segment/benches/multi_vector_search.rs +++ b/lib/segment/benches/multi_vector_search.rs @@ -91,7 +91,7 @@ fn make_segment_index(rng: &mut R, distance: Distance) -> HNSWI let hw_counter = HardwareCounterCell::new(); - let mut segment = build_segment(segment_dir.path(), &segment_config, true).unwrap(); + let mut segment = build_segment(segment_dir.path(), &segment_config, None, true).unwrap(); for n in 0..NUM_POINTS { let idx = (n as u64).into(); let multi_vec = random_multi_vector(rng, VECTOR_DIM, NUM_VECTORS_PER_POINT); diff --git a/lib/segment/src/index/struct_payload_index.rs b/lib/segment/src/index/struct_payload_index.rs index 9abf48f250..5170d0f51f 100644 --- a/lib/segment/src/index/struct_payload_index.rs +++ b/lib/segment/src/index/struct_payload_index.rs @@ -1245,7 +1245,13 @@ mod tests { drop(payload_config); // Load once and drop. - load_segment(&full_segment_path, Uuid::nil(), &AtomicBool::new(false)).unwrap(); + load_segment( + &full_segment_path, + Uuid::nil(), + None, + &AtomicBool::new(false), + ) + .unwrap(); // Check that index type has been written to disk again. // Proves we'll always persist the exact index type if it wasn't known yet at that time diff --git a/lib/segment/src/segment/mod.rs b/lib/segment/src/segment/mod.rs index 67787f9ec1..48ab1d2444 100644 --- a/lib/segment/src/segment/mod.rs +++ b/lib/segment/src/segment/mod.rs @@ -16,12 +16,14 @@ mod vectors; use std::collections::HashMap; use std::fmt; +use std::num::NonZeroUsize; use std::path::PathBuf; use std::sync::Arc; use atomic_refcell::AtomicRefCell; use common::is_alive_lock::IsAliveLock; use common::storage_version::StorageVersion; +use common::types::PointOffsetType; use parking_lot::Mutex; #[cfg(feature = "rocksdb")] use rocksdb::DB; @@ -93,6 +95,14 @@ pub struct Segment { pub error_status: Option, #[cfg(feature = "rocksdb")] pub database: Option>>, + /// Deferred-points threshold in bytes propagated from optimizer config. + /// If `None`, deferred-points behavior is disabled for this segment. + pub(crate) deferred_points_threshold_bytes: Option, + /// Cached deferred internal ID for fast visibility checks. + /// Points with internal id >= this value are hidden from reads. + /// It's `None` if deferred_points_threshold_bytes is `None` or if there are no deferred points. + /// Also `None` for non-appendable segments, as they don't accept new points and thus don't have deferred points. + pub(crate) deferred_internal_id: Option, } pub struct VectorData { diff --git a/lib/segment/src/segment/segment_ops.rs b/lib/segment/src/segment/segment_ops.rs index cbe5c8b2b1..8c6a51734f 100644 --- a/lib/segment/src/segment/segment_ops.rs +++ b/lib/segment/src/segment/segment_ops.rs @@ -1,5 +1,6 @@ use std::cmp::max; use std::collections::HashMap; +use std::num::NonZeroUsize; use std::path::Path; use bitvec::prelude::BitVec; @@ -55,6 +56,7 @@ impl Segment { vector_index.update_vector(internal_id, vector, hw_counter)?; self.version_tracker.set_vector(vector_name, Some(op_num)); } + self.update_deferred_internal_id(); Ok(()) } @@ -87,6 +89,7 @@ impl Segment { vector_index.update_vector(internal_id, Some(new_vector.as_vec_ref()), hw_counter)?; self.version_tracker.set_vector(&vector_name, Some(op_num)); } + self.update_deferred_internal_id(); Ok(()) } @@ -112,6 +115,7 @@ impl Segment { self.version_tracker.set_vector(vector_name, Some(op_num)); } self.id_tracker.borrow_mut().set_link(point_id, new_index)?; + self.update_deferred_internal_id(); Ok(new_index) } @@ -652,6 +656,50 @@ impl Segment { pub fn fix_id_tracker_inconsistencies(&mut self) -> OperationResult> { self.id_tracker.borrow_mut().fix_inconsistencies() } + + pub fn has_deferred_points(&self) -> bool { + // Point is deferred if his internal ID >= deferred_internal_id + self.deferred_internal_id.is_some() + } + + pub(crate) fn update_deferred_internal_id(&mut self) { + if self.deferred_internal_id.is_some() + || !self.is_appendable() + || self.config().vector_data.is_empty() + { + return; + } + + if let Some(deferred_points_threshold_bytes) = self.deferred_points_threshold_bytes { + self.deferred_internal_id = self + .vector_data + .iter() + .filter(|(vector_name, _)| { + // Only consider vectors without multivector config + self.segment_config + .vector_data + .get(vector_name.as_str()) + .is_some_and(|config| config.multivector_config.is_none()) + }) + .filter_map(|(_, vector_data)| -> Option { + let storage = vector_data.vector_storage.borrow(); + let vector_size: NonZeroUsize = storage + .get_vector_layout() + .ok() + .map(|layout| layout.size()) + .and_then(NonZeroUsize::new)?; + let deferred_internal_id = deferred_points_threshold_bytes + .get() + .div_ceil(vector_size.get()); + if deferred_internal_id < storage.total_vector_count() { + Some(deferred_internal_id as PointOffsetType) + } else { + None + } + }) + .min(); + } + } } fn restore_snapshot_in_place(snapshot_path: &Path) -> OperationResult<()> { diff --git a/lib/segment/src/segment/tests.rs b/lib/segment/src/segment/tests.rs index fa9b134e2c..5388c220df 100644 --- a/lib/segment/src/segment/tests.rs +++ b/lib/segment/src/segment/tests.rs @@ -249,7 +249,7 @@ fn test_snapshot(#[case] format: SnapshotFormat) { assert_eq!(entry.file_name(), segment_id); let restored_segment = - load_segment(&entry.path(), Uuid::nil(), &AtomicBool::new(false)).unwrap(); + load_segment(&entry.path(), Uuid::nil(), None, &AtomicBool::new(false)).unwrap(); // validate restored snapshot is the same as original segment assert_eq!( diff --git a/lib/segment/src/segment_constructor/segment_builder.rs b/lib/segment/src/segment_constructor/segment_builder.rs index 59c529143c..e9c0005553 100644 --- a/lib/segment/src/segment_constructor/segment_builder.rs +++ b/lib/segment/src/segment_constructor/segment_builder.rs @@ -1,6 +1,7 @@ use std::cmp; use std::collections::HashMap; use std::hash::{Hash, Hasher}; +use std::num::NonZeroUsize; use std::ops::Deref; use std::path::Path; use std::sync::Arc; @@ -466,6 +467,7 @@ impl SegmentBuilder { self.build( segments_path, Uuid::new_v4(), + None, ResourcePermit::dummy(num_rayon_threads(0) as u32), &AtomicBool::new(false), &mut rand::rng(), @@ -480,6 +482,7 @@ impl SegmentBuilder { self, segments_path: &Path, segment_uuid: Uuid, + deferred_points_threshold_bytes: Option, permit: ResourcePermit, stopped: &AtomicBool, rng: &mut R, @@ -723,7 +726,13 @@ impl SegmentBuilder { let destination_path = segments_path.join(segment_uuid.to_string()); fs::rename(temp_dir.keep(), &destination_path) .describe("Moving segment data after optimization")?; - load_segment(&destination_path, segment_uuid, stopped) + + load_segment( + &destination_path, + segment_uuid, + deferred_points_threshold_bytes, + stopped, + ) } fn update_quantization( diff --git a/lib/segment/src/segment_constructor/segment_constructor_base.rs b/lib/segment/src/segment_constructor/segment_constructor_base.rs index 641c90fe4e..5a1cd1e33d 100644 --- a/lib/segment/src/segment_constructor/segment_constructor_base.rs +++ b/lib/segment/src/segment_constructor/segment_constructor_base.rs @@ -1,5 +1,6 @@ use std::collections::HashMap; use std::io::Read; +use std::num::NonZeroUsize; use std::path::{Path, PathBuf}; use std::sync::Arc; use std::sync::atomic::AtomicBool; @@ -452,11 +453,13 @@ pub(crate) fn create_sparse_vector_storage( } } +#[allow(clippy::too_many_arguments)] fn create_segment( initial_version: Option, version: Option, segment_path: &Path, uuid: Uuid, + deferred_points_threshold_bytes: Option, config: &SegmentConfig, stopped: &AtomicBool, create: bool, @@ -626,7 +629,7 @@ fn create_segment( SegmentType::Plain }; - Ok(Segment { + let mut segment = Segment { uuid, initial_version, version, @@ -644,7 +647,13 @@ fn create_segment( error_status: None, #[cfg(feature = "rocksdb")] database: db_builder.build(), - }) + deferred_points_threshold_bytes, + deferred_internal_id: None, + }; + + segment.update_deferred_internal_id(); + + Ok(segment) } fn create_segment_id_tracker( @@ -768,7 +777,12 @@ pub fn normalize_segment_dir(path: &Path) -> OperationResult OperationResult { +pub fn load_segment( + path: &Path, + uuid: Uuid, + deferred_points_threshold_bytes: Option, + stopped: &AtomicBool, +) -> OperationResult { let stored_version = SegmentVersion::load(path)?.ok_or_else(|| { OperationError::service_error(format!( "Segment version file not found in segment: {}", @@ -814,6 +828,7 @@ pub fn load_segment(path: &Path, uuid: Uuid, stopped: &AtomicBool) -> OperationR segment_state.version, path, uuid, + deferred_points_threshold_bytes, &segment_state.config, stopped, false, @@ -849,6 +864,7 @@ pub fn load_segment(path: &Path, uuid: Uuid, stopped: &AtomicBool) -> OperationR pub fn build_segment( segments_path: &Path, config: &SegmentConfig, + deferred_points_threshold_bytes: Option, ready: bool, ) -> OperationResult { let uuid = Uuid::new_v4(); @@ -856,7 +872,16 @@ pub fn build_segment( let stopped = AtomicBool::new(false); fs::create_dir_all(&segment_path)?; - let segment = create_segment(None, None, &segment_path, uuid, config, &stopped, true)?; + let segment = create_segment( + None, + None, + &segment_path, + uuid, + deferred_points_threshold_bytes, + config, + &stopped, + true, + )?; segment.save_current_state()?; // Version is the last file to save, as it will be used to check if segment was built correctly. diff --git a/lib/segment/src/segment_constructor/simple_segment_constructor.rs b/lib/segment/src/segment_constructor/simple_segment_constructor.rs index 44a707a4a1..952e952df2 100644 --- a/lib/segment/src/segment_constructor/simple_segment_constructor.rs +++ b/lib/segment/src/segment_constructor/simple_segment_constructor.rs @@ -42,6 +42,7 @@ pub fn build_simple_segment( sparse_vector_data: Default::default(), payload_storage_type: Default::default(), }, + None, true, ) } @@ -70,6 +71,7 @@ pub fn build_simple_segment_with_payload_storage( sparse_vector_data: Default::default(), payload_storage_type, }, + None, true, ) } @@ -113,6 +115,7 @@ pub fn build_multivec_segment( sparse_vector_data: Default::default(), payload_storage_type: Default::default(), }, + None, true, ) } diff --git a/lib/segment/tests/integration/byte_storage_hnsw_test.rs b/lib/segment/tests/integration/byte_storage_hnsw_test.rs index 3cf7358c8b..99ed78efb4 100644 --- a/lib/segment/tests/integration/byte_storage_hnsw_test.rs +++ b/lib/segment/tests/integration/byte_storage_hnsw_test.rs @@ -99,7 +99,7 @@ fn test_byte_storage_hnsw( let int_key = "int"; let mut segment_float = build_simple_segment(dir_float.path(), dim, distance).unwrap(); - let mut segment_byte = build_segment(dir_byte.path(), &config_byte, true).unwrap(); + let mut segment_byte = build_segment(dir_byte.path(), &config_byte, None, true).unwrap(); // check that `segment_byte` uses byte or half storage { let borrowed_storage = segment_byte.vector_data[DEFAULT_VECTOR_NAME] diff --git a/lib/segment/tests/integration/byte_storage_quantization_test.rs b/lib/segment/tests/integration/byte_storage_quantization_test.rs index f582cb8ae5..4e331492da 100644 --- a/lib/segment/tests/integration/byte_storage_quantization_test.rs +++ b/lib/segment/tests/integration/byte_storage_quantization_test.rs @@ -237,7 +237,7 @@ fn test_byte_storage_binary_quantization_hnsw( let int_key = "int"; - let mut segment_byte = build_segment(dir_byte.path(), &config_byte, true).unwrap(); + let mut segment_byte = build_segment(dir_byte.path(), &config_byte, None, true).unwrap(); // check that `segment_byte` uses byte or half storage { let borrowed_storage = segment_byte.vector_data[DEFAULT_VECTOR_NAME] diff --git a/lib/segment/tests/integration/fixtures/segment.rs b/lib/segment/tests/integration/fixtures/segment.rs index 8fd6964cfb..7cb860128f 100644 --- a/lib/segment/tests/integration/fixtures/segment.rs +++ b/lib/segment/tests/integration/fixtures/segment.rs @@ -172,6 +172,7 @@ pub fn build_segment_3(path: &Path) -> Segment { sparse_vector_data: Default::default(), payload_storage_type: Default::default(), }, + None, true, ) .unwrap(); @@ -268,6 +269,7 @@ pub fn build_segment_sparse_1(path: &Path) -> Segment { )]), payload_storage_type: Default::default(), }, + None, true, ) .unwrap(); @@ -361,6 +363,7 @@ pub fn build_segment_sparse_2(path: &Path) -> Segment { )]), payload_storage_type: Default::default(), }, + None, true, ) .unwrap(); diff --git a/lib/segment/tests/integration/multivector_filtrable_hnsw_test.rs b/lib/segment/tests/integration/multivector_filtrable_hnsw_test.rs index a3e25e73a1..08fcbfa31c 100644 --- a/lib/segment/tests/integration/multivector_filtrable_hnsw_test.rs +++ b/lib/segment/tests/integration/multivector_filtrable_hnsw_test.rs @@ -83,7 +83,7 @@ fn test_multi_filterable_hnsw( let hw_counter = HardwareCounterCell::new(); - let mut segment = build_segment(dir.path(), &config, true).unwrap(); + let mut segment = build_segment(dir.path(), &config, None, true).unwrap(); for n in 0..num_points { let idx = n.into(); // Random number of vectors per multivec point diff --git a/lib/segment/tests/integration/multivector_quantization_test.rs b/lib/segment/tests/integration/multivector_quantization_test.rs index cdaa7b7974..0dcaa95360 100644 --- a/lib/segment/tests/integration/multivector_quantization_test.rs +++ b/lib/segment/tests/integration/multivector_quantization_test.rs @@ -230,7 +230,7 @@ fn test_multivector_quantization_hnsw( let int_key = "int"; - let mut segment = build_segment(dir.path(), &config, true).unwrap(); + let mut segment = build_segment(dir.path(), &config, None, true).unwrap(); let hw_counter = HardwareCounterCell::new(); diff --git a/lib/segment/tests/integration/payload_index_test.rs b/lib/segment/tests/integration/payload_index_test.rs index 05e3592ec3..c9f3bc2cec 100644 --- a/lib/segment/tests/integration/payload_index_test.rs +++ b/lib/segment/tests/integration/payload_index_test.rs @@ -86,9 +86,9 @@ impl TestSegments { let config = Self::make_simple_config(true); let mut plain_segment = - build_segment(&base_dir.path().join("plain"), &config, true).unwrap(); + build_segment(&base_dir.path().join("plain"), &config, None, true).unwrap(); let mut struct_segment = - build_segment(&base_dir.path().join("struct"), &config, true).unwrap(); + build_segment(&base_dir.path().join("struct"), &config, None, true).unwrap(); let num_points = 3000; let points_to_delete = 500; diff --git a/lib/segment/tests/integration/segment_builder_test.rs b/lib/segment/tests/integration/segment_builder_test.rs index 8682aef6e0..adfe2572f8 100644 --- a/lib/segment/tests/integration/segment_builder_test.rs +++ b/lib/segment/tests/integration/segment_builder_test.rs @@ -358,6 +358,7 @@ fn estimate_build_time(segment: &Segment, stop_delay_millis: Option) -> (u6 let res = builder.build( dir.path(), Uuid::new_v4(), + None, permit, &stopped, &mut rng, diff --git a/lib/segment/tests/integration/segment_on_disk_snapshot.rs b/lib/segment/tests/integration/segment_on_disk_snapshot.rs index 4a78a67d3e..9cceabf715 100644 --- a/lib/segment/tests/integration/segment_on_disk_snapshot.rs +++ b/lib/segment/tests/integration/segment_on_disk_snapshot.rs @@ -196,7 +196,7 @@ fn test_on_disk_segment_snapshot(#[case] format: SnapshotFormat) { assert_eq!(entry.file_name(), segment_id); let restored_segment = - load_segment(&entry.path(), Uuid::nil(), &AtomicBool::new(false)).unwrap(); + load_segment(&entry.path(), Uuid::nil(), None, &AtomicBool::new(false)).unwrap(); // validate restored snapshot is the same as original segment assert_eq!( diff --git a/lib/segment/tests/integration/segment_tests.rs b/lib/segment/tests/integration/segment_tests.rs index 0c62cd165c..f9e8c36425 100644 --- a/lib/segment/tests/integration/segment_tests.rs +++ b/lib/segment/tests/integration/segment_tests.rs @@ -206,7 +206,7 @@ fn ordered_deletion_test() { segment.segment_path.clone() }; - let segment = load_segment(&path, Uuid::nil(), &AtomicBool::new(false)).unwrap(); + let segment = load_segment(&path, Uuid::nil(), None, &AtomicBool::new(false)).unwrap(); let query_vector = [1.0, 1.0, 1.0, 1.0].into(); let res = segment diff --git a/lib/segment/tests/integration/sparse_discover_test.rs b/lib/segment/tests/integration/sparse_discover_test.rs index 7feceac54e..74fafa43da 100644 --- a/lib/segment/tests/integration/sparse_discover_test.rs +++ b/lib/segment/tests/integration/sparse_discover_test.rs @@ -152,8 +152,8 @@ fn sparse_index_discover_test() { sparse_vector_data: Default::default(), }; - let mut sparse_segment = build_segment(dir.path(), &sparse_config, true).unwrap(); - let mut dense_segment = build_segment(dir.path(), &dense_config, true).unwrap(); + let mut sparse_segment = build_segment(dir.path(), &sparse_config, None, true).unwrap(); + let mut dense_segment = build_segment(dir.path(), &dense_config, None, true).unwrap(); let hw_counter = HardwareCounterCell::new(); @@ -277,7 +277,7 @@ fn sparse_index_hardware_measurement_test() { payload_storage_type: Default::default(), }; - let mut sparse_segment = build_segment(dir.path(), &sparse_config, true).unwrap(); + let mut sparse_segment = build_segment(dir.path(), &sparse_config, None, true).unwrap(); let hw_counter = HardwareCounterCell::new(); diff --git a/lib/segment/tests/integration/sparse_vector_index_search_tests.rs b/lib/segment/tests/integration/sparse_vector_index_search_tests.rs index 3e97985f70..6003a187a5 100644 --- a/lib/segment/tests/integration/sparse_vector_index_search_tests.rs +++ b/lib/segment/tests/integration/sparse_vector_index_search_tests.rs @@ -599,7 +599,7 @@ fn sparse_vector_index_persistence_test() { )]), payload_storage_type: Default::default(), }; - let mut segment = build_segment(dir.path(), &config, true).unwrap(); + let mut segment = build_segment(dir.path(), &config, None, true).unwrap(); let hw_counter = HardwareCounterCell::new(); @@ -636,7 +636,7 @@ fn sparse_vector_index_persistence_test() { // persistence using rebuild of inverted index // for appendable segment vector index has to be rebuilt - let segment = load_segment(&path, Uuid::nil(), &stopped).unwrap(); + let segment = load_segment(&path, Uuid::nil(), None, &stopped).unwrap(); let search_after_reload_result = segment .search( SPARSE_VECTOR_NAME, @@ -771,7 +771,7 @@ fn sparse_vector_test_large_index() { )]), payload_storage_type: Default::default(), }; - let mut segment = build_segment(dir.path(), &config, true).unwrap(); + let mut segment = build_segment(dir.path(), &config, None, true).unwrap(); let hw_counter = HardwareCounterCell::new(); diff --git a/lib/shard/src/operations/optimization.rs b/lib/shard/src/operations/optimization.rs index 0ecb6a0f3d..7e919ecebc 100644 --- a/lib/shard/src/operations/optimization.rs +++ b/lib/shard/src/operations/optimization.rs @@ -1,3 +1,5 @@ +use std::num::NonZeroUsize; + use common::progress_tracker::ProgressTree; use schemars::JsonSchema; use serde::{Deserialize, Serialize}; @@ -126,4 +128,5 @@ pub struct OptimizerThresholds { pub max_segment_size_kb: usize, pub memmap_threshold_kb: usize, pub indexing_threshold_kb: usize, + pub deferred_points_threshold_bytes: Option, } diff --git a/lib/shard/src/optimize.rs b/lib/shard/src/optimize.rs index 10d543c2c1..ba2f113515 100644 --- a/lib/shard/src/optimize.rs +++ b/lib/shard/src/optimize.rs @@ -4,6 +4,7 @@ //! The collection layer provides the strategy via `OptimizationStrategy`. use std::collections::HashSet; +use std::num::NonZeroUsize; use std::ops::Deref; use std::path::{Path, PathBuf}; use std::sync::atomic::AtomicBool; @@ -153,6 +154,7 @@ fn build_new_segment( factory: &F, input_segments: &[LockedSegment], // Segments to optimize/merge into one output_segment_uuid: Uuid, // The UUID of the resulting optimized segment + deferred_points_threshold_bytes: Option, proxies: &[LockedSegment], permit: ResourcePermit, // IO resources for copying data resource_budget: ResourceBudget, @@ -283,6 +285,7 @@ fn build_new_segment( let mut optimized_segment = segment_builder.build( segments_path, output_segment_uuid, + deferred_points_threshold_bytes, indexing_permit, stopped, &mut rng, @@ -348,6 +351,7 @@ fn optimize_segment_propagate_changes( factory: &F, optimizing_segments: Vec, output_segment_uuid: Uuid, + deferred_points_threshold_bytes: Option, proxies: &[LockedSegment], permit: ResourcePermit, // IO resources for copying data resource_budget: ResourceBudget, @@ -364,6 +368,7 @@ fn optimize_segment_propagate_changes( factory, &optimizing_segments, output_segment_uuid, + deferred_points_threshold_bytes, proxies, permit, resource_budget, @@ -590,6 +595,7 @@ pub fn execute_optimization( segment_holder: LockedSegmentHolder, input_segment_ids: Vec, output_segment_uuid: Uuid, + deferred_points_threshold_bytes: Option, paths: &OptimizationPaths, permit: ResourcePermit, resource_budget: ResourceBudget, @@ -712,6 +718,7 @@ pub fn execute_optimization( factory, input_segments, output_segment_uuid, + deferred_points_threshold_bytes, &locked_proxies, permit, resource_budget, diff --git a/lib/shard/src/optimizers/indexing_optimizer.rs b/lib/shard/src/optimizers/indexing_optimizer.rs index 057222ef6d..baad029371 100644 --- a/lib/shard/src/optimizers/indexing_optimizer.rs +++ b/lib/shard/src/optimizers/indexing_optimizer.rs @@ -65,6 +65,8 @@ impl IndexingOptimizer { .memmap_threshold_kb .saturating_mul(BYTES_IN_KB); + let has_deferred_points = segment.has_deferred_points(); + for (vector_name, vector_cfg) in &self.segment_config.dense_vector { if let Some(vector_data) = segment_data_config.vector_data.get(vector_name) { let is_indexed = vector_data.index.is_indexed(); @@ -83,7 +85,7 @@ impl IndexingOptimizer { is_big_for_mmap && !is_on_disk }; - if optimize_for_index || optimize_for_mmap { + if optimize_for_index || optimize_for_mmap || has_deferred_points { return true; } } diff --git a/lib/shard/src/optimizers/segment_optimizer.rs b/lib/shard/src/optimizers/segment_optimizer.rs index b807295b9a..30b0437680 100644 --- a/lib/shard/src/optimizers/segment_optimizer.rs +++ b/lib/shard/src/optimizers/segment_optimizer.rs @@ -123,6 +123,7 @@ pub trait SegmentOptimizer: Sync { Ok(LockedSegment::new(build_segment( self.segments_path(), &config, + self.threshold_config().deferred_points_threshold_bytes, save_version, )?)) } @@ -327,6 +328,7 @@ pub trait SegmentOptimizer: Sync { segment_holder, input_segment_ids, output_segment_uuid, + self.threshold_config().deferred_points_threshold_bytes, &paths, permit, resource_budget, diff --git a/lib/shard/src/segment_holder/mod.rs b/lib/shard/src/segment_holder/mod.rs index 7f57fa619f..b21d8a4494 100644 --- a/lib/shard/src/segment_holder/mod.rs +++ b/lib/shard/src/segment_holder/mod.rs @@ -8,6 +8,7 @@ mod tests; use std::cmp::{max, min}; use std::collections::hash_map::Entry; use std::collections::{BTreeMap, BTreeSet, BinaryHeap, HashSet}; +use std::num::NonZeroUsize; use std::ops::Deref; use std::path::Path; use std::sync::Arc; @@ -756,11 +757,13 @@ impl SegmentHolder { segments_path: &Path, segment_config: SegmentConfig, payload_index_schema: Arc>, + deferred_points_threshold_bytes: Option, ) -> OperationResult { let segment = self.build_tmp_segment( segments_path, Some(segment_config), payload_index_schema, + deferred_points_threshold_bytes, true, )?; self.add_new_locked(segment.clone()); @@ -789,6 +792,7 @@ impl SegmentHolder { segments_path: &Path, segment_config: Option, payload_index_schema: Arc>, + deferred_points_threshold_bytes: Option, save_version: bool, ) -> OperationResult { let config = match segment_config { @@ -809,7 +813,12 @@ impl SegmentHolder { .clone(), }; - let mut segment = build_segment(segments_path, &config, save_version)?; + let mut segment = build_segment( + segments_path, + &config, + deferred_points_threshold_bytes, + save_version, + )?; // Internal operation. let hw_counter = HardwareCounterCell::disposable(); diff --git a/lib/shard/src/segment_holder/snapshot.rs b/lib/shard/src/segment_holder/snapshot.rs index 34e6b47bbd..8a44cba376 100644 --- a/lib/shard/src/segment_holder/snapshot.rs +++ b/lib/shard/src/segment_holder/snapshot.rs @@ -1,3 +1,4 @@ +use std::num::NonZeroUsize; use std::path::Path; use std::sync::Arc; @@ -36,6 +37,7 @@ impl SegmentHolder { segments_path: &Path, segment_config: Option, payload_index_schema: Arc>, + deferred_points_threshold_bytes: Option, ) -> OperationResult<( Vec<(SegmentId, LockedSegment)>, SegmentId, @@ -50,6 +52,7 @@ impl SegmentHolder { segments_path, segment_config, payload_index_schema, + deferred_points_threshold_bytes, false, )?; diff --git a/lib/shard/src/segment_holder/tests.rs b/lib/shard/src/segment_holder/tests.rs index 581db47188..fff0b5ece8 100644 --- a/lib/shard/src/segment_holder/tests.rs +++ b/lib/shard/src/segment_holder/tests.rs @@ -548,6 +548,7 @@ fn test_double_proxies() { segments_dir.path(), None, schema.clone(), + None, ) .unwrap(); @@ -568,8 +569,14 @@ fn test_double_proxies() { .unwrap(); let (outer_proxies, outer_tmp_segment, outer_segments_lock) = - SegmentHolder::proxy_all_segments(inner_segments_lock, segments_dir.path(), None, schema) - .unwrap(); + SegmentHolder::proxy_all_segments( + inner_segments_lock, + segments_dir.path(), + None, + schema, + None, + ) + .unwrap(); let mut has_point = false; for (_proxy_id, proxy) in &outer_proxies { diff --git a/src/segment_inspector.rs b/src/segment_inspector.rs index 335ed10515..5169948b3b 100644 --- a/src/segment_inspector.rs +++ b/src/segment_inspector.rs @@ -48,7 +48,7 @@ fn main() { .and_then(|s| Uuid::try_parse(s.to_str()?).ok()) .unwrap_or(Uuid::nil()); - let segment = load_segment(path, segment_uuid, &AtomicBool::new(false)).unwrap(); + let segment = load_segment(path, segment_uuid, None, &AtomicBool::new(false)).unwrap(); eprintln!( "path = {:#?}, size-points = {}",