use std::cmp::Reverse;
use std::fmt;
use std::ops::Range;
use itertools::Itertools;
use quickwit_doc_mapper::tag_pruning::{field_tag, match_tag_field_name};
use quickwit_metastore::SplitMetadata;
use tracing::debug;
use crate::new_split_id;
pub enum MergeOperation {
Merge {
merge_split_id: String,
splits: Vec<SplitMetadata>,
},
Demux {
demux_split_ids: Vec<String>,
splits: Vec<SplitMetadata>,
},
}
impl MergeOperation {
pub fn new_merge_operation(splits: Vec<SplitMetadata>) -> MergeOperation {
MergeOperation::Merge {
merge_split_id: new_split_id(),
splits,
}
}
pub fn splits(&self) -> &[SplitMetadata] {
match self {
MergeOperation::Merge { splits, .. } | MergeOperation::Demux { splits, .. } => {
splits.as_slice()
}
}
}
}
impl fmt::Debug for MergeOperation {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
MergeOperation::Merge {
merge_split_id: split_id,
splits,
} => {
write!(f, "Merge(merged_split_id={},splits=[", split_id)?;
for split in splits {
write!(f, "{},", split.split_id())?;
}
write!(f, "])")?;
}
MergeOperation::Demux {
demux_split_ids,
splits,
} => {
write!(f, "Merge(demux_splits=[")?;
for split_id in demux_split_ids {
write!(f, "{},", split_id)?;
}
write!(f, "], input_splits=[")?;
for split in splits {
write!(f, "{},", split.split_id())?;
}
write!(f, "])")?;
}
}
Ok(())
}
}
pub trait MergePolicy: Send + Sync + fmt::Debug {
fn operations(&self, splits: &mut Vec<SplitMetadata>) -> Vec<MergeOperation>;
fn is_mature(&self, split: &SplitMetadata) -> bool;
}
#[derive(Clone, Debug)]
pub struct StableMultitenantWithTimestampMergePolicy {
pub demux_enabled: bool,
pub demux_factor: usize,
pub demux_field_name: Option<String>,
pub min_level_num_docs: usize,
pub merge_enabled: bool,
pub merge_factor: usize,
pub max_merge_factor: usize,
pub split_num_docs_target: usize,
}
impl Default for StableMultitenantWithTimestampMergePolicy {
fn default() -> Self {
StableMultitenantWithTimestampMergePolicy {
demux_enabled: false,
demux_field_name: None,
demux_factor: 6,
min_level_num_docs: 100_000,
merge_enabled: true,
merge_factor: 10,
max_merge_factor: 12,
split_num_docs_target: 10_000_000,
}
}
}
fn remove_matching_items<T, Pred: Fn(&T) -> bool>(items: &mut Vec<T>, predicate: Pred) -> Vec<T> {
let mut matching_items = Vec::new();
let mut i = 0;
while i < items.len() {
if predicate(&items[i]) {
let matching_item = items.remove(i);
matching_items.push(matching_item);
} else {
i += 1;
}
}
matching_items
}
struct SplitShortDebug<'a>(&'a SplitMetadata);
impl<'a> fmt::Debug for SplitShortDebug<'a> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Split")
.field("split_id", &self.0.split_id())
.field("num_docs", &self.0.num_docs)
.finish()
}
}
fn splits_short_debug(splits: &[SplitMetadata]) -> Vec<SplitShortDebug> {
splits.iter().map(SplitShortDebug).collect()
}
impl MergePolicy for StableMultitenantWithTimestampMergePolicy {
fn operations(&self, splits: &mut Vec<SplitMetadata>) -> Vec<MergeOperation> {
let original_num_splits = splits.len();
let mut operations = self.merge_operations(splits);
operations.append(&mut self.demux_operations(splits));
debug_assert_eq!(
original_num_splits,
operations.iter().map(|op| op.splits().len()).sum::<usize>() + splits.len(),
"The merge policy is supposed to keep the number of splits."
);
operations
}
fn is_mature(&self, split: &SplitMetadata) -> bool {
self.is_mature_for_merge(split) && self.is_mature_for_demux(split)
}
}
#[derive(Clone, Copy, Eq, PartialEq)]
enum MergeCandidateSize {
TooSmall,
ValidSplit,
OneMoreSplitWouldBeTooBig,
}
impl StableMultitenantWithTimestampMergePolicy {
fn is_mature_for_merge(&self, split: &SplitMetadata) -> bool {
if !self.merge_enabled {
return true;
}
split.num_docs >= self.split_num_docs_target || split.demux_num_ops > 0
}
fn is_mature_for_demux(&self, split: &SplitMetadata) -> bool {
if !self.demux_enabled || self.demux_field_name.is_none() {
return true;
}
if split.num_docs >= self.demux_factor * self.split_num_docs_target {
return true;
}
let demux_field_name = self.demux_field_name.as_ref().unwrap();
if split.tags.contains(&field_tag(demux_field_name))
&& split
.tags
.iter()
.filter(|tag| match_tag_field_name(demux_field_name, tag))
.count()
< 2
{
return true;
};
split.num_docs < self.split_num_docs_target || split.demux_num_ops > 0
}
fn merge_operations(&self, splits: &mut Vec<SplitMetadata>) -> Vec<MergeOperation> {
if !self.merge_enabled || splits.is_empty() {
return Vec::new();
}
let splits_not_for_merge =
remove_matching_items(splits, |split| self.is_mature_for_merge(split));
let mut merge_operations: Vec<MergeOperation> = Vec::new();
splits.sort_by_key(|split| {
let time_end = split
.time_range
.as_ref()
.map(|time_range| Reverse(*time_range.end()));
(time_end, split.num_docs)
});
debug!(splits=?splits_short_debug(&splits[..]), "merge-policy-run");
let split_levels = self.build_split_levels(splits);
for split_range in split_levels.into_iter().rev() {
debug!(splits=?splits_short_debug(&splits[split_range.clone()]));
if let Some(merge_range) = self.merge_candidate_from_level(splits, split_range) {
debug!(merge_range=?merge_range, "merge-candidate");
let splits_in_merge: Vec<SplitMetadata> = splits.drain(merge_range).collect();
let merge_operation = MergeOperation::new_merge_operation(splits_in_merge);
merge_operations.push(merge_operation);
} else {
debug!("no-merge");
}
}
splits.extend(splits_not_for_merge);
merge_operations
}
fn demux_operations(&self, splits: &mut Vec<SplitMetadata>) -> Vec<MergeOperation> {
if !self.demux_enabled || self.demux_field_name.is_none() || splits.is_empty() {
return Vec::new();
}
let excluded_splits =
remove_matching_items(splits, |split| self.is_mature_for_demux(split));
let mut merge_operations: Vec<MergeOperation> = Vec::new();
splits.sort_by_key(|split| {
split
.time_range
.as_ref()
.map(|time_range| *time_range.end())
});
merge_operations.append(&mut self.build_first_demux_operation(splits));
splits.extend(excluded_splits);
merge_operations
}
pub(crate) fn build_first_demux_operation(
&self,
splits: &mut Vec<SplitMetadata>,
) -> Vec<MergeOperation> {
assert!(self.demux_factor > 1, "Demux factor must be > 1");
assert!(
splits.iter().all(|split| split.demux_num_ops == 0),
"All splits are expected to have never been demuxed."
);
assert!(
splits
.iter()
.all(|split| split.num_docs < self.split_num_docs_target * self.demux_factor),
"Each split size must satisfy `max_merge_docs <= size < demux_factor * max_merge_docs`"
);
let mut total_num_docs_left: usize = splits.iter().map(|split| split.num_docs).sum();
if splits.is_empty() || total_num_docs_left < self.demux_factor * self.split_num_docs_target
{
return Vec::new();
}
let mut operations = Vec::new();
while !splits.is_empty()
&& total_num_docs_left >= self.demux_factor * self.split_num_docs_target
{
let mut end_split_idx = 0;
let mut num_docs_to_demux = 0;
for (split_idx, split) in splits.iter().enumerate().take(self.demux_factor) {
num_docs_to_demux += split.num_docs;
if num_docs_to_demux >= self.demux_factor * self.split_num_docs_target {
end_split_idx = split_idx;
break;
}
}
assert!(
end_split_idx > 0,
"Impossible state with splits.len() > 0 and total docs > demux factor * max merge \
docs. This should never happened."
);
let splits_for_demux: Vec<SplitMetadata> = splits.drain(0..end_split_idx + 1).collect();
total_num_docs_left -= num_docs_to_demux;
let merge_operation = MergeOperation::Demux {
demux_split_ids: (0..self.demux_factor).map(|_| new_split_id()).collect_vec(),
splits: splits_for_demux,
};
operations.push(merge_operation);
}
operations
}
pub(crate) fn build_split_levels(&self, splits: &[SplitMetadata]) -> Vec<Range<usize>> {
assert!(
splits
.iter()
.all(|split| split.num_docs < self.split_num_docs_target),
"All splits are expected to be smaller than `max_merge_docs`."
);
if splits.is_empty() {
return Vec::new();
}
let mut split_levels: Vec<Range<usize>> = Vec::new();
let mut current_level_start_ord = 0;
let mut current_level_max_docs = (splits[0].num_docs * 3).max(self.min_level_num_docs);
let mut levels = vec![(0..current_level_max_docs)]; for (split_ord, split) in splits.iter().enumerate() {
if split.num_docs >= current_level_max_docs {
split_levels.push(current_level_start_ord..split_ord);
current_level_start_ord = split_ord;
current_level_max_docs = 3 * split.num_docs;
levels.push(split.num_docs..current_level_max_docs)
}
}
debug!(levels=?levels);
split_levels.push(current_level_start_ord..splits.len());
split_levels
}
fn merge_candidate_from_level(
&self,
splits: &[SplitMetadata],
level_range: Range<usize>,
) -> Option<Range<usize>> {
let merge_candidate_end = level_range.end;
let mut merge_candidate_start = merge_candidate_end;
for split_ord in level_range.rev() {
if self.merge_candidate_size(&splits[merge_candidate_start..merge_candidate_end])
== MergeCandidateSize::OneMoreSplitWouldBeTooBig
{
break;
}
merge_candidate_start = split_ord;
}
if self.merge_candidate_size(&splits[merge_candidate_start..merge_candidate_end])
== MergeCandidateSize::TooSmall
{
return None;
}
Some(merge_candidate_start..merge_candidate_end)
}
fn merge_candidate_size(&self, splits: &[SplitMetadata]) -> MergeCandidateSize {
if splits.len() <= 1 {
return MergeCandidateSize::TooSmall;
}
if splits.len() >= self.max_merge_factor {
return MergeCandidateSize::OneMoreSplitWouldBeTooBig;
}
let num_docs_in_merge: usize = splits.iter().map(|split| split.num_docs).sum();
if num_docs_in_merge >= self.split_num_docs_target {
return MergeCandidateSize::OneMoreSplitWouldBeTooBig;
}
if splits.len() < self.merge_factor {
return MergeCandidateSize::TooSmall;
}
MergeCandidateSize::ValidSplit
}
fn case_levels_given_growth_factor(&self, growth_factor: usize) -> Vec<usize> {
assert!(self.min_level_num_docs > 0);
assert!(self.merge_factor > 1);
assert!(self.max_merge_factor >= self.merge_factor);
assert!(self.split_num_docs_target > self.min_level_num_docs);
let mut levels_start_num_docs = vec![1];
let mut level_end_doc = self.min_level_num_docs;
while level_end_doc < self.split_num_docs_target {
levels_start_num_docs.push(level_end_doc);
level_end_doc *= growth_factor;
}
levels_start_num_docs.push(self.split_num_docs_target);
levels_start_num_docs
}
pub fn max_num_splits_ideal_case(&self, num_docs: u64) -> usize {
let levels = self.case_levels_given_growth_factor(self.merge_factor);
self.max_num_splits_knowning_levels(num_docs, &levels, true)
}
pub fn max_num_splits_worst_case(&self, num_docs: u64) -> usize {
let levels = self.case_levels_given_growth_factor(3);
self.max_num_splits_knowning_levels(num_docs, &levels, false)
}
fn max_num_splits_knowning_levels(
&self,
mut num_docs: u64,
levels: &[usize],
sorted: bool,
) -> usize {
assert!(is_sorted(levels));
if num_docs == 0 {
return 0;
}
let (&head, tail) = levels.split_first().unwrap();
if num_docs < head as u64 {
return 0;
}
let first_level_min_saturation_docs = if sorted {
head * (self.merge_factor - 1)
} else {
head + (self.merge_factor - 2)
};
if tail.is_empty() || num_docs <= first_level_min_saturation_docs as u64 {
return (num_docs as usize + head - 1) / head;
}
num_docs -= first_level_min_saturation_docs as u64;
self.merge_factor - 1 + self.max_num_splits_knowning_levels(num_docs, tail, sorted)
}
}
fn is_sorted(els: &[usize]) -> bool {
els.windows(2).all(|w| w[0] <= w[1])
}
#[cfg(test)]
mod tests {
use std::collections::BTreeSet;
use std::iter::FromIterator;
use std::ops::RangeInclusive;
use super::*;
fn create_splits(num_docs_vec: Vec<usize>) -> Vec<SplitMetadata> {
let num_docs_with_timestamp = num_docs_vec
.into_iter()
.map(|num_docs| (num_docs, (1630563067..=1630564067)))
.collect();
create_splits_with_timestamps(num_docs_with_timestamp)
}
fn create_splits_with_timestamps(
num_docs_vec: Vec<(usize, RangeInclusive<i64>)>,
) -> Vec<SplitMetadata> {
num_docs_vec
.into_iter()
.enumerate()
.map(|(split_ord, (num_docs, time_range))| SplitMetadata {
split_id: format!("split_{:02}", split_ord),
num_docs,
time_range: Some(time_range),
..Default::default()
})
.collect()
}
fn create_splits_with_tags(
num_docs_vec: Vec<usize>,
demux_field_name: &str,
tag_counts: &[usize],
) -> Vec<SplitMetadata> {
num_docs_vec
.into_iter()
.zip(tag_counts.iter())
.enumerate()
.map(|(split_ord, (num_docs, &tag_count))| SplitMetadata {
split_id: format!("split_{:02}", split_ord),
num_docs,
tags: (0..tag_count)
.into_iter()
.map(|i| format!("{}:{}", demux_field_name, i))
.collect::<BTreeSet<_>>(),
..Default::default()
})
.collect()
}
#[test]
fn test_split_is_mature_with_no_demux_field() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let mut split = create_splits(vec![9_000_000]).into_iter().next().unwrap();
assert!(!merge_policy.is_mature(&split));
let merge_policy_with_disabled_merge = StableMultitenantWithTimestampMergePolicy {
merge_enabled: false,
..Default::default()
};
assert!(merge_policy_with_disabled_merge.is_mature(&split));
split.demux_num_ops = 1;
assert!(merge_policy.is_mature(&split));
split.num_docs = merge_policy.split_num_docs_target + 1;
split.demux_num_ops = 0;
assert!(merge_policy.is_mature(&split));
split.num_docs = merge_policy.split_num_docs_target + 1;
split.demux_num_ops = 1;
assert!(merge_policy.is_mature(&split));
}
#[test]
fn test_split_is_mature_with_demux_field() {
let merge_policy = StableMultitenantWithTimestampMergePolicy {
demux_enabled: true,
demux_field_name: Some("demux_field".to_owned()),
..Default::default()
};
{
let mut split = create_splits(vec![9_000_000]).into_iter().next().unwrap();
assert!(!merge_policy.is_mature(&split));
let demux_tags = BTreeSet::from_iter(vec![
"demux_field!".to_string(), "demux_field:1".to_string(),
"demux_field:2".to_string(),
]);
split.tags = demux_tags;
assert!(!merge_policy.is_mature(&split));
split.num_docs = merge_policy.split_num_docs_target + 1;
assert!(!merge_policy.is_mature(&split));
}
{
let mut split = create_splits(vec![9_000_000]).into_iter().next().unwrap();
split.num_docs = merge_policy.split_num_docs_target + 1;
let demux_tags = BTreeSet::from_iter(vec![
"demux_field!".to_string(), "demux_field:1".to_string(),
]);
split.tags = demux_tags;
assert!(merge_policy.is_mature(&split));
split.num_docs = merge_policy.demux_factor * merge_policy.split_num_docs_target;
split.demux_num_ops = 0;
assert!(merge_policy.is_mature(&split));
split.num_docs = 100;
split.demux_num_ops = 1;
assert!(merge_policy.is_mature(&split));
let other_tags = BTreeSet::from_iter(vec![
"demux_field!".to_string(), "other_field:1".to_string(),
"other_field:2".to_string(),
]);
split.tags = other_tags;
split.num_docs = merge_policy.split_num_docs_target + 1;
split.demux_num_ops = 0;
assert!(merge_policy.is_mature(&split));
}
{
let merge_policy_with_disabled_demux = StableMultitenantWithTimestampMergePolicy {
demux_enabled: false,
demux_field_name: Some("demux_field".to_owned()),
..Default::default()
};
let mut split = create_splits(vec![9_000_000]).into_iter().next().unwrap();
split.num_docs = merge_policy.split_num_docs_target + 1;
assert!(merge_policy_with_disabled_demux.is_mature(&split));
split.num_docs = 100;
split.demux_num_ops = 1;
assert!(merge_policy.is_mature(&split));
split.num_docs = merge_policy.split_num_docs_target + 1;
split.demux_num_ops = 1;
assert!(merge_policy.is_mature(&split));
}
}
#[test]
fn test_build_split_levels() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let splits = Vec::new();
let split_groups = merge_policy.build_split_levels(&splits);
assert!(split_groups.is_empty());
}
#[test]
fn test_stable_multitenant_merge_policy_build_split_simple() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let splits = create_splits(vec![100_000, 100_000, 100_000, 800_000, 900_000]);
let split_groups = merge_policy.build_split_levels(&splits);
assert_eq!(&split_groups, &[0..3, 3..5]);
}
#[test]
fn test_stable_multitenant_merge_policy_build_split_perfect_world() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let splits = create_splits(vec![
100_000, 100_000, 100_000, 100_000, 100_000, 100_000, 100_000, 100_000, 800_000,
1_600_000,
]);
let split_groups = merge_policy.build_split_levels(&splits);
assert_eq!(&split_groups, &[0..8, 8..10]);
}
#[test]
fn test_stable_multitenant_merge_policy_build_split_decreasing() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let splits = create_splits(vec![
100_000, 100_000, 100_000, 100_000, 100_000, 100_000, 100_000, 100_000, 800_000,
100_000, 1_600_000,
]);
let split_groups = merge_policy.build_split_levels(&splits);
assert_eq!(&split_groups, &[0..8, 8..11]);
}
#[test]
#[should_panic(expected = "All splits are expected to be smaller than `max_merge_docs`.")]
fn test_stable_multitenant_merge_policy_build_split_panics_if_exceeding_max_merge_docs() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let splits = create_splits(vec![11_000_000]);
merge_policy.build_split_levels(&splits);
}
#[test]
fn test_stable_multitenant_merge_policy_not_enough_splits() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let mut splits = create_splits(vec![100; 7]);
assert_eq!(splits.len(), 7);
assert!(merge_policy.operations(&mut splits).is_empty());
}
#[test]
fn test_stable_multitenant_merge_policy_just_enough_splits_for_a_merge() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let mut splits = create_splits(vec![100; 10]);
let mut merge_ops = merge_policy.operations(&mut splits);
assert!(splits.is_empty());
assert_eq!(merge_ops.len(), 1);
let merge_op = merge_ops.pop().unwrap();
let mut merge_segment_ids: Vec<String> = merge_op
.splits()
.iter()
.map(|split| split.split_id().to_string())
.collect();
merge_segment_ids.sort();
assert!(matches!(merge_op, MergeOperation::Merge { .. }));
assert_eq!(
merge_segment_ids,
&[
"split_00", "split_01", "split_02", "split_03", "split_04", "split_05", "split_06",
"split_07", "split_08", "split_09"
]
);
}
#[test]
fn test_stable_multitenant_merge_policy_many_splits_on_same_level() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let mut splits = create_splits(vec![100; 13]);
let mut merge_ops = merge_policy.operations(&mut splits);
assert_eq!(splits.len(), 1);
assert_eq!(splits[0].split_id(), "split_00");
assert_eq!(merge_ops.len(), 1);
let merge_op = merge_ops.pop().unwrap();
let mut merge_split_ids: Vec<String> = merge_op
.splits()
.iter()
.map(|split| split.split_id().to_string())
.collect();
merge_split_ids.sort();
assert_eq!(
merge_split_ids,
&[
"split_01", "split_02", "split_03", "split_04", "split_05", "split_06", "split_07",
"split_08", "split_09", "split_10", "split_11", "split_12"
]
);
assert!(matches!(merge_op, MergeOperation::Merge { .. }));
}
#[test]
fn test_stable_multitenant_merge_policy_splits_below_min_level() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let mut splits = create_splits(vec![
100, 1000, 10_000, 10_000, 10_000, 10_000, 10_000, 40_000, 40_000, 40_000,
]);
let mut merge_ops = merge_policy.operations(&mut splits);
assert_eq!(splits.len(), 0);
assert_eq!(merge_ops.len(), 1);
let merge_op = merge_ops.pop().unwrap();
let mut merge_split_ids: Vec<String> = merge_op
.splits()
.iter()
.map(|split| split.split_id().to_string())
.collect();
merge_split_ids.sort();
assert_eq!(
merge_split_ids,
&[
"split_00", "split_01", "split_02", "split_03", "split_04", "split_05", "split_06",
"split_07", "split_08", "split_09"
]
);
assert!(matches!(merge_op, MergeOperation::Merge { .. }));
}
#[test]
fn test_stable_multitenant_merge_policy_splits_above_min_level() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let mut splits = create_splits(vec![
100_000, 1_000_000, 1_000_000, 1_000_000, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
]);
let merge_ops = merge_policy.operations(&mut splits);
assert_eq!(splits.len(), 8);
assert_eq!(merge_ops.len(), 0);
}
#[test]
fn test_stable_multitenant_merge_policy_above_max_merge_docs_is_ignored() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let mut splits = create_splits(vec![
100_000, 100_000, 100_000, 100_000, 100_000,
10_000_000, 100_000, 100_000, 100_000, 100_000, 100_000,
]);
let merge_ops = merge_policy.operations(&mut splits);
assert_eq!(splits.len(), 1);
assert_eq!(splits[0].num_docs, 10_000_000);
assert_eq!(merge_ops.len(), 1);
}
#[test]
fn test_merge_policy_splits_too_large_are_ignored() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let mut splits = create_splits(vec![9_999_999, 10_000_000]);
let merge_ops = merge_policy.operations(&mut splits);
assert_eq!(splits.len(), 2);
assert_eq!(splits[0].num_docs, 9_999_999);
assert_eq!(splits[1].num_docs, 10_000_000);
assert!(merge_ops.is_empty());
}
#[test]
fn test_merge_policy_splits_entire_level_reach_merge_max_doc() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let mut splits = create_splits(vec![5_000_000, 5_000_000]);
let merge_ops = merge_policy.operations(&mut splits);
assert!(splits.is_empty());
assert_eq!(merge_ops.len(), 1);
assert!(matches!(merge_ops[0], MergeOperation::Merge { .. }));
assert_eq!(merge_ops[0].splits().len(), 2);
}
#[test]
fn test_merge_policy_last_merge_can_have_a_lower_merge_factor() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let mut splits = create_splits(vec![9_999_997, 9_999_998, 9_999_999]);
let merge_ops = merge_policy.operations(&mut splits);
assert_eq!(splits.len(), 1);
assert_eq!(splits[0].num_docs, 9_999_997);
assert_eq!(merge_ops.len(), 1);
assert!(matches!(merge_ops[0], MergeOperation::Merge { .. }));
assert_eq!(merge_ops[0].splits().len(), 2);
}
#[test]
fn test_merge_policy_no_merge_with_only_one_split() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
let mut splits = create_splits(vec![9_999_999]);
let merge_ops = merge_policy.operations(&mut splits);
assert_eq!(splits.len(), 1);
assert_eq!(splits[0].num_docs, 9_999_999);
assert!(merge_ops.is_empty());
}
#[test]
fn test_stable_multitenant_merge_policy_max_num_splits_worst_case() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
assert_eq!(merge_policy.max_num_splits_worst_case(99), 9);
assert_eq!(merge_policy.max_num_splits_worst_case(1_000_000), 27);
assert_eq!(merge_policy.max_num_splits_worst_case(2_000_000), 36);
assert_eq!(merge_policy.max_num_splits_worst_case(3_000_000), 36);
assert_eq!(merge_policy.max_num_splits_worst_case(4_000_000), 36);
assert_eq!(merge_policy.max_num_splits_worst_case(5_000_000), 45);
assert_eq!(merge_policy.max_num_splits_worst_case(7_000_000), 45);
assert_eq!(merge_policy.max_num_splits_worst_case(10_000_000), 45);
assert_eq!(merge_policy.max_num_splits_worst_case(20_000_000), 54);
assert_eq!(merge_policy.max_num_splits_worst_case(100_000_000), 63);
assert_eq!(merge_policy.max_num_splits_worst_case(1_000_000_000), 153);
}
#[test]
fn test_stable_multitenant_merge_policy_max_num_splits_ideal_case() {
let merge_policy = StableMultitenantWithTimestampMergePolicy::default();
assert_eq!(merge_policy.max_num_splits_ideal_case(1_000_000), 18);
assert_eq!(merge_policy.max_num_splits_ideal_case(99), 9);
assert_eq!(merge_policy.max_num_splits_ideal_case(2_000_000), 20);
assert_eq!(merge_policy.max_num_splits_ideal_case(3_000_000), 21);
assert_eq!(merge_policy.max_num_splits_ideal_case(4_000_000), 22);
assert_eq!(merge_policy.max_num_splits_ideal_case(5_000_000), 23);
assert_eq!(merge_policy.max_num_splits_ideal_case(7_000_000), 25);
assert_eq!(merge_policy.max_num_splits_ideal_case(10_000_000), 27);
assert_eq!(merge_policy.max_num_splits_ideal_case(100_000_000), 37);
assert_eq!(merge_policy.max_num_splits_ideal_case(1_000_000_000), 127);
}
#[test]
fn test_demux_one_operation_and_filter_out_irrelevant_splits() {
let demux_field_name = "demux_field_name";
let merge_policy = StableMultitenantWithTimestampMergePolicy {
demux_enabled: true,
demux_field_name: Some(demux_field_name.to_string()),
demux_factor: 6,
min_level_num_docs: 100_000,
merge_enabled: true,
merge_factor: 10,
max_merge_factor: 12,
split_num_docs_target: 10_000_000,
};
let mut demux_candidates = create_splits_with_tags(
vec![
10_000_000, 10_000_000, 12_000_000, 14_000_000, 10_000_000, 10_000_001, 10_000_002,
10_000_004, 10_000_005, 60_000_000,
],
demux_field_name,
&[0, 1, 2, 3, 3, 4, 5, 6, 10],
);
let mut other_splits =
create_splits_with_tags(vec![10_000_000], "other_demux_field_name", &[5]);
demux_candidates.append(&mut other_splits);
let merge_ops = merge_policy.demux_operations(&mut demux_candidates);
assert_eq!(demux_candidates.len(), 4);
assert_eq!(merge_ops.len(), 1);
assert!(matches!(merge_ops[0], MergeOperation::Demux { .. }));
assert_eq!(merge_ops[0].splits().len(), 6);
}
#[test]
fn test_demux_one_operation_with_1_normal_splits_and_1_huge_splits() {
let demux_field_name = "demux_field_name";
let merge_policy = StableMultitenantWithTimestampMergePolicy {
demux_enabled: true,
demux_field_name: Some(demux_field_name.to_string()),
..Default::default()
};
let mut demux_candidates = create_splits_with_tags(
vec![50_000_000, 10_000_000, 12_000_000],
demux_field_name,
&[2, 2, 2],
);
let merge_ops = merge_policy.demux_operations(&mut demux_candidates);
assert_eq!(demux_candidates.len(), 1);
assert_eq!(demux_candidates[0].split_id(), "split_02");
assert_eq!(merge_ops.len(), 1);
assert!(matches!(merge_ops[0], MergeOperation::Demux { .. }));
assert_eq!(merge_ops[0].splits().len(), 2);
}
#[test]
fn test_should_ignore_demux_operation_with_1_huge_split() {
let demux_field_name = "demux_field_name";
let merge_policy = StableMultitenantWithTimestampMergePolicy {
demux_field_name: Some(demux_field_name.to_string()),
..Default::default()
};
let mut demux_candidates =
create_splits_with_tags(vec![60_000_000], demux_field_name, &[2]);
let merge_ops = merge_policy.demux_operations(&mut demux_candidates);
assert_eq!(demux_candidates.len(), 1);
assert_eq!(merge_ops.len(), 0);
}
#[test]
fn test_demux_two_operations() {
let demux_field_name = "demux_field_name";
let merge_policy = StableMultitenantWithTimestampMergePolicy {
demux_enabled: true,
demux_field_name: Some(demux_field_name.to_string()),
merge_enabled: true,
..Default::default()
};
let mut splits = create_splits_with_tags(
vec![
19_999_999, 19_999_999, 10_000_000, 10_000_001, 10_000_002, 10_000_004, 10_000_005,
19_999_999, 19_999_999, 10_000_000, 10_000_001, 10_000_002, 10_000_004, 10_000_005,
],
demux_field_name,
&[5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5, 5],
);
let merge_ops = merge_policy.demux_operations(&mut splits);
assert_eq!(splits.len(), 5);
assert_eq!(merge_ops.len(), 2);
assert!(matches!(merge_ops[0], MergeOperation::Demux { .. }));
assert!(matches!(merge_ops[1], MergeOperation::Demux { .. }));
assert_eq!(merge_ops[0].splits().len(), 5);
assert_eq!(merge_ops[1].splits().len(), 4);
}
#[test]
fn test_stable_multitenant_merge_policy_merge_not_enabled() {
let merge_policy = StableMultitenantWithTimestampMergePolicy {
merge_enabled: false,
..Default::default()
};
let mut splits = create_splits(vec![100; 10]);
let merge_ops = merge_policy.operations(&mut splits);
assert_eq!(splits.len(), 10);
assert_eq!(merge_ops.len(), 0);
}
#[test]
fn test_stable_multitenant_merge_policy_demux_not_enabled() {
let demux_field_name = "demux_field_name";
let merge_policy = StableMultitenantWithTimestampMergePolicy {
demux_enabled: false,
demux_field_name: Some(demux_field_name.to_string()),
..Default::default()
};
let mut splits = create_splits_with_tags(vec![10_000_000; 10], demux_field_name, &[10; 10]);
let merge_ops = merge_policy.demux_operations(&mut splits);
assert_eq!(merge_ops.len(), 0);
}
}