mirror of
https://github.com/qdrant/qdrant.git
synced 2026-10-03 03:17:43 -05:00
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
This commit is contained in:
@@ -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(),
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -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([
|
||||
|
||||
@@ -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<NonZeroUsize> {
|
||||
(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) {
|
||||
|
||||
@@ -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<OperationWithClockTag> =
|
||||
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);
|
||||
}
|
||||
|
||||
@@ -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<SegmentConfig>,
|
||||
payload_index_schema: Arc<SaveOnDisk<PayloadIndexSchema>>,
|
||||
deferred_points_threshold_bytes: Option<NonZeroUsize>,
|
||||
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<F>(
|
||||
segments_path: &Path,
|
||||
segment_config: Option<SegmentConfig>,
|
||||
payload_index_schema: Arc<SaveOnDisk<PayloadIndexSchema>>,
|
||||
deferred_points_threshold_bytes: Option<NonZeroUsize>,
|
||||
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
|
||||
|
||||
@@ -47,6 +47,7 @@ fn test_snapshot_all() {
|
||||
segments_dir.path(),
|
||||
None,
|
||||
schema,
|
||||
None,
|
||||
temp_dir.path(),
|
||||
&tar,
|
||||
SnapshotFormat::Regular,
|
||||
|
||||
@@ -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 =
|
||||
|
||||
@@ -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);
|
||||
|
||||
+10
-7
@@ -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());
|
||||
|
||||
@@ -91,7 +91,7 @@ fn make_segment_index<R: Rng + ?Sized>(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);
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<SegmentFailedState>,
|
||||
#[cfg(feature = "rocksdb")]
|
||||
pub database: Option<Arc<parking_lot::RwLock<DB>>>,
|
||||
/// 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<NonZeroUsize>,
|
||||
/// 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<PointOffsetType>,
|
||||
}
|
||||
|
||||
pub struct VectorData {
|
||||
|
||||
@@ -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<Vec<PointOffsetType>> {
|
||||
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<PointOffsetType> {
|
||||
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<()> {
|
||||
|
||||
@@ -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!(
|
||||
|
||||
@@ -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<NonZeroUsize>,
|
||||
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(
|
||||
|
||||
@@ -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<SeqNumberType>,
|
||||
version: Option<SeqNumberType>,
|
||||
segment_path: &Path,
|
||||
uuid: Uuid,
|
||||
deferred_points_threshold_bytes: Option<NonZeroUsize>,
|
||||
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<Option<(PathBuf, Uu
|
||||
/// Preferably, the `uuid` should match the last component of `path`.
|
||||
/// In production use [`normalize_segment_dir`] to obtain correct path and UUID.
|
||||
/// In tests it is acceptable to pass an arbitrary UUID, e.g., [`Uuid::nil()`].
|
||||
pub fn load_segment(path: &Path, uuid: Uuid, stopped: &AtomicBool) -> OperationResult<Segment> {
|
||||
pub fn load_segment(
|
||||
path: &Path,
|
||||
uuid: Uuid,
|
||||
deferred_points_threshold_bytes: Option<NonZeroUsize>,
|
||||
stopped: &AtomicBool,
|
||||
) -> OperationResult<Segment> {
|
||||
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<NonZeroUsize>,
|
||||
ready: bool,
|
||||
) -> OperationResult<Segment> {
|
||||
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.
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
}
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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();
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -358,6 +358,7 @@ fn estimate_build_time(segment: &Segment, stop_delay_millis: Option<u64>) -> (u6
|
||||
let res = builder.build(
|
||||
dir.path(),
|
||||
Uuid::new_v4(),
|
||||
None,
|
||||
permit,
|
||||
&stopped,
|
||||
&mut rng,
|
||||
|
||||
@@ -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!(
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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();
|
||||
|
||||
|
||||
@@ -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();
|
||||
|
||||
|
||||
@@ -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<NonZeroUsize>,
|
||||
}
|
||||
|
||||
@@ -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<F: ?Sized + OptimizationStrategy>(
|
||||
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<NonZeroUsize>,
|
||||
proxies: &[LockedSegment],
|
||||
permit: ResourcePermit, // IO resources for copying data
|
||||
resource_budget: ResourceBudget,
|
||||
@@ -283,6 +285,7 @@ fn build_new_segment<F: ?Sized + OptimizationStrategy>(
|
||||
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<F: ?Sized + OptimizationStrategy>(
|
||||
factory: &F,
|
||||
optimizing_segments: Vec<LockedSegment>,
|
||||
output_segment_uuid: Uuid,
|
||||
deferred_points_threshold_bytes: Option<NonZeroUsize>,
|
||||
proxies: &[LockedSegment],
|
||||
permit: ResourcePermit, // IO resources for copying data
|
||||
resource_budget: ResourceBudget,
|
||||
@@ -364,6 +368,7 @@ fn optimize_segment_propagate_changes<F: ?Sized + OptimizationStrategy>(
|
||||
factory,
|
||||
&optimizing_segments,
|
||||
output_segment_uuid,
|
||||
deferred_points_threshold_bytes,
|
||||
proxies,
|
||||
permit,
|
||||
resource_budget,
|
||||
@@ -590,6 +595,7 @@ pub fn execute_optimization<F: ?Sized + OptimizationStrategy>(
|
||||
segment_holder: LockedSegmentHolder,
|
||||
input_segment_ids: Vec<SegmentId>,
|
||||
output_segment_uuid: Uuid,
|
||||
deferred_points_threshold_bytes: Option<NonZeroUsize>,
|
||||
paths: &OptimizationPaths,
|
||||
permit: ResourcePermit,
|
||||
resource_budget: ResourceBudget,
|
||||
@@ -712,6 +718,7 @@ pub fn execute_optimization<F: ?Sized + OptimizationStrategy>(
|
||||
factory,
|
||||
input_segments,
|
||||
output_segment_uuid,
|
||||
deferred_points_threshold_bytes,
|
||||
&locked_proxies,
|
||||
permit,
|
||||
resource_budget,
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<SaveOnDisk<PayloadIndexSchema>>,
|
||||
deferred_points_threshold_bytes: Option<NonZeroUsize>,
|
||||
) -> OperationResult<LockedSegment> {
|
||||
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<SegmentConfig>,
|
||||
payload_index_schema: Arc<SaveOnDisk<PayloadIndexSchema>>,
|
||||
deferred_points_threshold_bytes: Option<NonZeroUsize>,
|
||||
save_version: bool,
|
||||
) -> OperationResult<LockedSegment> {
|
||||
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();
|
||||
|
||||
@@ -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<SegmentConfig>,
|
||||
payload_index_schema: Arc<SaveOnDisk<PayloadIndexSchema>>,
|
||||
deferred_points_threshold_bytes: Option<NonZeroUsize>,
|
||||
) -> OperationResult<(
|
||||
Vec<(SegmentId, LockedSegment)>,
|
||||
SegmentId,
|
||||
@@ -50,6 +52,7 @@ impl SegmentHolder {
|
||||
segments_path,
|
||||
segment_config,
|
||||
payload_index_schema,
|
||||
deferred_points_threshold_bytes,
|
||||
false,
|
||||
)?;
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 = {}",
|
||||
|
||||
Reference in New Issue
Block a user