use std::path::{Path, PathBuf};
use std::sync::Arc;
use itertools::Itertools;
use ordered_float::OrderedFloat;
use parking_lot::Mutex;
use crate::segment::common::operation_time_statistics::OperationDurationsAggregator;
use crate::segment::entry::ReadSegmentEntry;
use crate::segment::index::VectorIndexRead;
use crate::segment::segment::Segment;
use crate::segment::types::HnswGlobalConfig;
use crate::segment::vector_storage::VectorStorageRead;
use super::config::SegmentOptimizerConfig;
use super::segment_optimizer::{OptimizationPlanner, SegmentOptimizer};
use crate::shard::operations::optimization::OptimizerThresholds;
pub struct VacuumOptimizer {
deleted_threshold: f64,
min_vectors_number: usize,
thresholds_config: OptimizerThresholds,
segments_path: PathBuf,
temp_path: PathBuf,
segment_optimizer_config: SegmentOptimizerConfig,
hnsw_global_config: HnswGlobalConfig,
telemetry_durations_aggregator: Arc<Mutex<OperationDurationsAggregator>>,
}
impl VacuumOptimizer {
#[allow(clippy::too_many_arguments)]
pub fn new(
deleted_threshold: f64,
min_vectors_number: usize,
thresholds_config: OptimizerThresholds,
segments_path: PathBuf,
temp_path: PathBuf,
segment_optimizer_config: SegmentOptimizerConfig,
hnsw_global_config: HnswGlobalConfig,
) -> Self {
VacuumOptimizer {
deleted_threshold,
min_vectors_number,
thresholds_config,
segments_path,
temp_path,
segment_optimizer_config,
hnsw_global_config,
telemetry_durations_aggregator: OperationDurationsAggregator::new(),
}
}
fn littered_ratio_segment(&self, segment: &Segment) -> Option<f64> {
let littered_ratio =
segment.deleted_point_count() as f64 / segment.total_point_count() as f64;
let is_big = segment.total_point_count() >= self.min_vectors_number;
let is_littered = littered_ratio > self.deleted_threshold;
(is_big && is_littered).then_some(littered_ratio)
}
fn littered_vectors_index_ratio(&self, segment: &Segment) -> Option<f64> {
let segment_config = segment.config();
if !segment_config.is_any_vector_indexed() {
return None;
}
segment
.vector_data
.values()
.filter(|vector_data| vector_data.vector_index.borrow().is_index())
.filter_map(|vector_data| {
let vector_index = vector_data.vector_index.borrow();
let vector_storage = vector_data.vector_storage.borrow();
let indexed_vector_count = vector_index.indexed_vector_count();
let deleted_from_index =
indexed_vector_count.saturating_sub(vector_storage.available_vector_count());
let deleted_ratio = if indexed_vector_count != 0 {
deleted_from_index as f64 / indexed_vector_count as f64
} else {
0.0
};
let reached_minimum = deleted_from_index >= self.min_vectors_number;
let reached_ratio = deleted_ratio > self.deleted_threshold;
(reached_minimum && reached_ratio).then_some(deleted_ratio)
})
.max_by_key(|ratio| OrderedFloat(*ratio))
}
}
impl SegmentOptimizer for VacuumOptimizer {
fn name(&self) -> &'static str {
"vacuum"
}
fn segments_path(&self) -> &Path {
self.segments_path.as_path()
}
fn temp_path(&self) -> &Path {
self.temp_path.as_path()
}
fn segment_optimizer_config(&self) -> &SegmentOptimizerConfig {
&self.segment_optimizer_config
}
fn hnsw_global_config(&self) -> &HnswGlobalConfig {
&self.hnsw_global_config
}
fn threshold_config(&self) -> &OptimizerThresholds {
&self.thresholds_config
}
fn plan_optimizations(&self, planner: &mut OptimizationPlanner) {
let to_optimize = planner
.remaining()
.iter()
.filter_map(|(&segment_id, segment)| {
let segment = segment.read();
let littered_ratio_segment = self.littered_ratio_segment(&segment);
let littered_ratio_vectors = self.littered_vectors_index_ratio(&segment);
let worst_ratio = std::iter::chain(littered_ratio_segment, littered_ratio_vectors)
.max_by_key(|ratio| OrderedFloat(*ratio));
worst_ratio.map(|ratio| (segment_id, ratio))
})
.sorted_by_key(|(_, ratio)| OrderedFloat(-ratio))
.collect_vec();
for (segment_id, _) in to_optimize {
planner.plan(vec![segment_id]);
}
}
fn get_telemetry_counter(&self) -> &Mutex<OperationDurationsAggregator> {
&self.telemetry_durations_aggregator
}
}