diff --git a/src/indexer/index_writer.rs b/src/indexer/index_writer.rs index 40c2b678b..31293591c 100644 --- a/src/indexer/index_writer.rs +++ b/src/indexer/index_writer.rs @@ -897,7 +897,7 @@ mod tests { let index_writer = index.writer(3_000_000).unwrap(); assert_eq!( format!("{:?}", index_writer.get_merge_policy()), - "LogMergePolicy { min_merge_size: 8, min_layer_size: 10000, \ + "LogMergePolicy { min_merge_size: 8, max_merge_size: 10000000, min_layer_size: 10000, \ level_log_size: 0.75 }" ); let merge_policy = Box::new(NoMergePolicy::default()); diff --git a/src/indexer/log_merge_policy.rs b/src/indexer/log_merge_policy.rs index b51a2df54..591d5e840 100644 --- a/src/indexer/log_merge_policy.rs +++ b/src/indexer/log_merge_policy.rs @@ -6,12 +6,14 @@ use std::f64; const DEFAULT_LEVEL_LOG_SIZE: f64 = 0.75; const DEFAULT_MIN_LAYER_SIZE: u32 = 10_000; const DEFAULT_MIN_MERGE_SIZE: usize = 8; +const DEFAULT_MAX_MERGE_SIZE: usize = 10_000_000; /// `LogMergePolicy` tries tries to merge segments that have a similar number of /// documents. #[derive(Debug, Clone)] pub struct LogMergePolicy { min_merge_size: usize, + max_merge_size: usize, min_layer_size: u32, level_log_size: f64, } @@ -26,6 +28,12 @@ impl LogMergePolicy { self.min_merge_size = min_merge_size; } + /// Set the maximum number docs in a segment for it to be considered for + /// merging. + pub fn set_max_merge_size(&mut self, max_merge_size: usize) { + self.max_merge_size = max_merge_size; + } + /// Set the minimum segment size under which all segment belong /// to the same level. pub fn set_min_layer_size(&mut self, min_layer_size: u32) { @@ -53,6 +61,7 @@ impl MergePolicy for LogMergePolicy { let mut size_sorted_tuples = segments .iter() .map(SegmentMeta::num_docs) + .filter(|s| s <= &(self.max_merge_size as u32)) .enumerate() .collect::>(); @@ -86,6 +95,7 @@ impl Default for LogMergePolicy { fn default() -> LogMergePolicy { LogMergePolicy { min_merge_size: DEFAULT_MIN_MERGE_SIZE, + max_merge_size: DEFAULT_MAX_MERGE_SIZE, min_layer_size: DEFAULT_MIN_LAYER_SIZE, level_log_size: DEFAULT_LEVEL_LOG_SIZE, } @@ -104,6 +114,7 @@ mod tests { fn test_merge_policy() -> LogMergePolicy { let mut log_merge_policy = LogMergePolicy::default(); log_merge_policy.set_min_merge_size(3); + log_merge_policy.set_max_merge_size(100_000); log_merge_policy.set_min_layer_size(2); log_merge_policy } @@ -141,11 +152,11 @@ mod tests { create_random_segment_meta(10), create_random_segment_meta(10), create_random_segment_meta(10), - create_random_segment_meta(1000), - create_random_segment_meta(1000), - create_random_segment_meta(1000), - create_random_segment_meta(10000), - create_random_segment_meta(10000), + create_random_segment_meta(1_000), + create_random_segment_meta(1_000), + create_random_segment_meta(1_000), + create_random_segment_meta(10_000), + create_random_segment_meta(10_000), create_random_segment_meta(10), create_random_segment_meta(10), create_random_segment_meta(10), @@ -182,4 +193,19 @@ mod tests { let result_list = test_merge_policy().compute_merge_candidates(&test_input); assert_eq!(result_list.len(), 1); } + + #[test] + fn test_large_merge_segments() { + let test_input = vec![ + create_random_segment_meta(1_000_000), + create_random_segment_meta(100_001), + create_random_segment_meta(100_000), + create_random_segment_meta(100_000), + create_random_segment_meta(100_000), + ]; + let result_list = test_merge_policy().compute_merge_candidates(&test_input); + // Do not include large segments + assert_eq!(result_list.len(), 1); + assert_eq!(result_list[0].0.len(), 3) + } }