use std::collections::BTreeMap;
use std::collections::Bound::{Excluded, Included, Unbounded};
use std::fmt::Debug;
use std::ops::Bound;
use std::path::{Path, PathBuf};
use crate::common::fs::{atomic_save_bin, atomic_save_json};
use crate::common::types::PointOffsetType;
use crate::common::universal_io::{CachedReadFs, UniversalReadFs, read_bin_via, read_json_via};
use itertools::Itertools;
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
use crate::segment::common::operation_error::OperationResult;
use crate::segment::index::field_index::numeric_point::{Numericable, Point};
use crate::segment::index::field_index::utils::check_boundaries;
const MIN_BUCKET_SIZE: usize = 10;
const CONFIG_PATH: &str = "histogram_config.json";
const BORDERS_PATH: &str = "histogram_borders.bin";
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct Counts {
pub left: usize,
pub right: usize,
}
#[derive(Default, Debug, PartialEq)]
pub struct Histogram<T: Numericable + Serialize + DeserializeOwned> {
max_bucket_size: usize,
precision: f64,
total_count: usize,
borders: BTreeMap<Point<T>, Counts>,
}
#[derive(Debug, Serialize, Deserialize)]
struct HistogramConfig {
max_bucket_size: usize,
precision: f64,
total_count: usize,
}
impl<T: Numericable + Serialize + DeserializeOwned> Histogram<T> {
pub fn new(max_bucket_size: usize, precision: f64) -> Self {
assert!(precision < 1.0);
assert!(precision > 0.0);
Self {
max_bucket_size,
precision,
total_count: 0,
borders: BTreeMap::default(),
}
}
pub fn preopen(fs: &impl CachedReadFs, path: &Path) -> OperationResult<()> {
fs.schedule_prefetch(&path.join(CONFIG_PATH), None, None)?;
fs.schedule_prefetch(&path.join(BORDERS_PATH), None, None)?;
Ok(())
}
pub fn open<Fs: UniversalReadFs>(fs: &Fs, path: &Path) -> OperationResult<Self> {
let config_path = path.join(CONFIG_PATH);
let borders_path = path.join(BORDERS_PATH);
let histogram_config: HistogramConfig = read_json_via(fs, &config_path)?;
let histogram_buckets: Vec<(Point<T>, Counts)> = read_bin_via(fs, &borders_path)?;
Ok(Self {
max_bucket_size: histogram_config.max_bucket_size,
precision: histogram_config.precision,
total_count: histogram_config.total_count,
borders: histogram_buckets.into_iter().collect(),
})
}
pub fn save(&self, path: &Path) -> OperationResult<()> {
let config_path = path.join(CONFIG_PATH);
let borders_path = path.join(BORDERS_PATH);
atomic_save_json(
&config_path,
&HistogramConfig {
max_bucket_size: self.max_bucket_size,
precision: self.precision,
total_count: self.total_count,
},
)?;
let borders: Vec<(Point<T>, Counts)> =
self.borders.iter().map(|(k, v)| (*k, v.clone())).collect();
atomic_save_bin(&borders_path, &borders)?;
Ok(())
}
pub fn files(path: &Path) -> Vec<PathBuf> {
vec![path.join(CONFIG_PATH), path.join(BORDERS_PATH)]
}
pub fn immutable_files(path: &Path) -> Vec<PathBuf> {
vec![path.join(CONFIG_PATH)]
}
#[cfg(test)]
pub fn total_count(&self) -> usize {
self.total_count
}
#[cfg(test)]
pub fn borders(&self) -> &BTreeMap<Point<T>, Counts> {
&self.borders
}
pub fn current_bucket_size(&self) -> usize {
let bucket_size = (self.total_count as f64 * self.precision) as usize;
bucket_size.clamp(MIN_BUCKET_SIZE, self.max_bucket_size)
}
pub fn get_total_count(&self) -> usize {
self.total_count
}
pub fn get_range_by_size(&self, from: Bound<T>, range_size: usize) -> Bound<T> {
let from_ = match from {
Included(val) => Included(Point::new(val, PointOffsetType::MIN)),
Excluded(val) => Excluded(Point::new(val, PointOffsetType::MAX)),
Unbounded => Unbounded,
};
let mut reached_count = 0;
for (border, counts) in self.borders.range((from_, Unbounded)) {
if reached_count + counts.left > range_size {
return Included(border.val);
} else {
reached_count += counts.left;
}
}
Unbounded
}
pub fn estimate(&self, from: Bound<T>, to: Bound<T>) -> (usize, usize, usize) {
let from_ = match &from {
Included(val) => Included(Point::new(*val, PointOffsetType::MIN)),
Excluded(val) => Excluded(Point::new(*val, PointOffsetType::MAX)),
Unbounded => Unbounded,
};
let to_ = match &to {
Included(val) => Included(Point::new(*val, PointOffsetType::MAX)),
Excluded(val) => Excluded(Point::new(*val, PointOffsetType::MIN)),
Unbounded => Unbounded,
};
let from_val = match from {
Included(val) => val,
Excluded(val) => val,
Unbounded => T::min_value(),
};
let to_val = match to {
Included(val) => val,
Excluded(val) => val,
Unbounded => T::max_value(),
};
let left_border = {
if matches!(from_, Unbounded) {
None
} else {
self.borders.range((Unbounded, from_)).next_back()
}
};
let right_border = {
if matches!(to_, Unbounded) {
None
} else {
self.borders.range((to_, Unbounded)).next()
}
};
if !check_boundaries(&from_, &to_) {
return (0, 0, 0);
}
left_border
.into_iter()
.chain(self.borders.range((from_, to_)))
.chain(right_border)
.tuple_windows()
.map(
|((a, a_count), (b, _b_count)): ((&Point<T>, &Counts), (&Point<T>, _))| {
let val_range = (b.val - a.val).to_f64();
if val_range == 0. {
let estimates = a_count.right + 1;
return (estimates, estimates, estimates);
}
if a_count.right == 0 {
return (1, 1, 1);
}
let cover_range = (to_val.min(b.val) - from_val.max(a.val)).to_f64();
let covered_frac = cover_range / val_range;
let estimate = (a_count.right as f64 * covered_frac).round() as usize + 1;
let min_estimate = if cover_range == val_range {
a_count.right + 1
} else {
0
};
let max_estimate = a_count.right + 1;
(min_estimate, estimate, max_estimate)
},
)
.reduce(|a, b| (a.0 + b.0, a.1 + b.1, a.2 + b.2))
.unwrap_or((0, 0, 0))
}
pub fn remove<F, G>(&mut self, val: &Point<T>, left_neighbour: F, right_neighbour: G)
where
F: Fn(&Point<T>) -> Option<Point<T>>,
G: Fn(&Point<T>) -> Option<Point<T>>,
{
let (mut close_neighbors, (mut far_left_neighbor, mut far_right_neighbor)) = {
let mut left_iterator = self
.borders
.range((Unbounded, Included(val)))
.map(|(k, v)| (*k, v.clone()));
let mut right_iterator = self
.borders
.range((Excluded(val), Unbounded))
.map(|(k, v)| (*k, v.clone()));
(
(left_iterator.next_back(), right_iterator.next()),
(left_iterator.next_back(), right_iterator.next()),
)
};
let (to_remove, to_create, removed) = match &mut close_neighbors {
(None, None) => (None, None, false), (Some((left_border, left_border_count)), None) => {
if left_border == val {
if left_border_count.left == 0 {
(Some(*left_border), None, true)
} else {
if let Some((_fln, fln_count)) = &mut far_left_neighbor {
fln_count.right -= 1
}
let (new_border, new_border_count) = (
left_neighbour(left_border).unwrap(),
Counts {
left: left_border_count.left - 1,
right: 0,
},
);
(
Some(*left_border),
Some((new_border, new_border_count)),
true,
)
}
} else {
(None, None, false)
}
}
(None, Some((right_border, right_border_count))) => {
if right_border == val {
if right_border_count.right == 0 {
(Some(*right_border), None, true)
} else {
if let Some((_frn, frn_count)) = &mut far_right_neighbor {
frn_count.left -= 1
}
let (new_border, new_border_count) = (
right_neighbour(right_border).unwrap(),
Counts {
left: 0,
right: right_border_count.right - 1,
},
);
(
Some(*right_border),
Some((new_border, new_border_count)),
true,
)
}
} else {
(None, None, false)
}
}
(Some((left_border, left_border_count)), Some((right_border, right_border_count))) => {
if left_border == val {
if left_border_count.right == 0 {
right_border_count.left = left_border_count.left;
(Some(*left_border), None, true)
} else if right_border_count.left + left_border_count.left
<= self.current_bucket_size()
&& far_left_neighbor.is_some()
{
if let Some((_fln, fln_count)) = &mut far_left_neighbor {
fln_count.right += right_border_count.left;
right_border_count.left = fln_count.right;
}
(Some(*left_border), None, true)
} else {
right_border_count.left -= 1;
let (new_border, new_border_count) = (
right_neighbour(left_border).unwrap(),
Counts {
left: left_border_count.left,
right: left_border_count.right - 1,
},
);
(
Some(*left_border),
Some((new_border, new_border_count)),
true,
)
}
} else if right_border == val {
if right_border_count.left == 0 {
left_border_count.right = right_border_count.left;
(Some(*right_border), None, true)
} else if left_border_count.right + right_border_count.right
<= self.current_bucket_size()
&& far_right_neighbor.is_some()
{
if let Some((_frn, frn_count)) = &mut far_right_neighbor {
frn_count.left += left_border_count.right;
left_border_count.right = frn_count.left;
}
(Some(*right_border), None, true)
} else {
left_border_count.right -= 1;
let (new_border, new_border_count) = (
left_neighbour(right_border).unwrap(),
Counts {
left: right_border_count.right,
right: right_border_count.left - 1,
},
);
(
Some(*right_border),
Some((new_border, new_border_count)),
true,
)
}
} else if right_border_count.left == 0 {
(None, None, false)
} else {
right_border_count.left -= 1;
left_border_count.right -= 1;
(None, None, true)
}
}
};
if removed {
self.total_count -= 1;
}
let (left_border_opt, right_border_opt) = close_neighbors;
if let Some((k, v)) = left_border_opt {
self.borders.insert(k, v);
}
if let Some((k, v)) = right_border_opt {
self.borders.insert(k, v);
}
if let Some((k, v)) = far_left_neighbor {
self.borders.insert(k, v);
}
if let Some((k, v)) = far_right_neighbor {
self.borders.insert(k, v);
}
if let Some(remove_border) = to_remove {
self.borders.remove(&remove_border);
}
if let Some((new_border, new_border_count)) = to_create {
self.borders.insert(new_border, new_border_count);
}
}
pub fn insert<F, G>(&mut self, val: Point<T>, left_neighbour: F, right_neighbour: G)
where
F: Fn(&Point<T>) -> Option<Point<T>>,
G: Fn(&Point<T>) -> Option<Point<T>>,
{
self.total_count += 1;
if self.borders.len() < 2 {
self.borders.insert(val, Counts { left: 0, right: 0 });
return;
}
let (mut close_neighbors, (mut far_left_neighbor, mut far_right_neighbor)) = {
let mut left_iterator = self
.borders
.range((Unbounded, Included(val)))
.map(|(k, v)| (*k, v.clone()));
let mut right_iterator = self
.borders
.range((Excluded(val), Unbounded))
.map(|(k, v)| (*k, v.clone()));
(
(left_iterator.next_back(), right_iterator.next()),
(left_iterator.next_back(), right_iterator.next()),
)
};
let (to_remove, to_create) = match &mut close_neighbors {
(None, Some((right_border, right_border_count))) => {
let new_count = right_border_count.right + 1;
let (new_border, mut new_border_count) = (
val,
Counts {
left: 0,
right: new_count,
},
);
if new_count > self.current_bucket_size() {
new_border_count.right = 0;
(None, Some((new_border, new_border_count)))
} else {
if let Some((_frn, frn_count)) = &mut far_right_neighbor {
frn_count.left = new_count;
}
(Some(*right_border), Some((new_border, new_border_count)))
}
}
(Some((left_border, left_border_count)), None) => {
let new_count = left_border_count.left + 1;
let (new_border, mut new_border_count) = (
val,
Counts {
left: new_count,
right: 0,
},
);
if new_count > self.current_bucket_size() {
new_border_count.left = 0;
(None, Some((new_border, new_border_count)))
} else {
if let Some((_fln, fln_count)) = &mut far_left_neighbor {
fln_count.right = new_count
}
(Some(*left_border), Some((new_border, new_border_count)))
}
}
(Some((left_border, left_border_count)), Some((right_border, right_border_count))) => {
assert_eq!(left_border_count.right, right_border_count.left);
let new_count = left_border_count.right + 1;
if new_count > self.current_bucket_size() {
let left_dist = val.val.abs_diff(left_border.val);
let right_dist = val.val.abs_diff(right_border.val);
if left_dist < right_dist {
let (new_border, mut new_border_count) = (
right_neighbour(left_border).unwrap(),
Counts {
left: left_border_count.left + 1,
right: left_border_count.right,
},
);
if left_border_count.left < self.current_bucket_size()
&& far_left_neighbor.is_some()
{
if let Some((_fln, fln_count)) = &mut far_left_neighbor {
fln_count.right = new_border_count.left
}
(Some(*left_border), Some((new_border, new_border_count)))
} else {
new_border_count.left = 0;
left_border_count.right = 0;
(None, Some((new_border, new_border_count)))
}
} else {
let (new_border, mut new_border_count) = (
left_neighbour(right_border).unwrap(),
Counts {
left: right_border_count.left,
right: right_border_count.right + 1,
},
);
if right_border_count.right < self.current_bucket_size()
&& far_right_neighbor.is_some()
{
if let Some((_frn, frn_count)) = &mut far_right_neighbor {
frn_count.left = new_border_count.right
}
(Some(*right_border), Some((new_border, new_border_count)))
} else {
new_border_count.right = 0;
right_border_count.left = 0;
(None, Some((new_border, new_border_count)))
}
}
} else {
left_border_count.right = new_count;
right_border_count.left = new_count;
(None, None)
}
}
(None, None) => unreachable!(),
};
let (left_border_opt, right_border_opt) = close_neighbors;
if let Some((k, v)) = left_border_opt {
self.borders.insert(k, v);
}
if let Some((k, v)) = right_border_opt {
self.borders.insert(k, v);
}
if let Some((k, v)) = far_left_neighbor {
self.borders.insert(k, v);
}
if let Some((k, v)) = far_right_neighbor {
self.borders.insert(k, v);
}
if let Some(remove_border) = to_remove {
self.borders.remove(&remove_border);
}
if let Some((new_border, new_border_count)) = to_create {
self.borders.insert(new_border, new_border_count);
}
}
pub fn ram_usage_bytes(&self) -> usize {
let Self {
max_bucket_size: _,
precision: _,
total_count: _,
borders,
} = self;
let btree_entry_overhead = std::mem::size_of::<usize>() * 3;
borders.len()
* (std::mem::size_of::<Point<T>>()
+ std::mem::size_of::<Counts>()
+ btree_entry_overhead)
}
}