1use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
16use std::fmt::Debug;
17
18use common_telemetry::info;
19use common_time::Timestamp;
20use common_time::range::TimestampRange;
21use common_time::timestamp::TimeUnit;
22use common_time::timestamp_millis::BucketAligned;
23use snafu::ResultExt;
24use store_api::storage::RegionId;
25
26use crate::compaction::CompactionOutput;
27use crate::compaction::buckets::infer_time_bucket;
28use crate::compaction::compactor::{CompactionRegion, CompactionVersion};
29use crate::compaction::last_non_null::inputs_precede_memtables;
30use crate::compaction::picker::{Picker, PickerOutput, get_expired_ssts};
31use crate::error::{JoinSnafu, Result};
32use crate::region::options::{CompactionOptions, MergeMode};
33use crate::sst::file::FileHandle;
34
35#[derive(Clone, Debug)]
39pub struct WindowedCompactionPicker {
40 compaction_time_window_seconds: Option<i64>,
41 time_range: Option<TimestampRange>,
42}
43
44impl WindowedCompactionPicker {
45 pub fn new(window_seconds: Option<i64>) -> Self {
46 Self {
47 compaction_time_window_seconds: window_seconds,
48 time_range: None,
49 }
50 }
51
52 pub(crate) fn with_time_range(mut self, time_range: Option<TimestampRange>) -> Self {
54 self.time_range = time_range;
55 self
56 }
57
58 fn calculate_time_window(
63 &self,
64 region_id: RegionId,
65 current_version: &CompactionVersion,
66 ) -> i64 {
67 self.compaction_time_window_seconds
68 .or(current_version
69 .compaction_time_window
70 .map(|t| t.as_secs() as i64))
71 .unwrap_or_else(|| {
72 let levels = current_version.ssts.levels();
73 let inferred = infer_time_bucket(levels[0].files());
74 info!(
75 "Compaction window for region {} is not present, inferring from files: {:?}",
76 region_id, inferred
77 );
78 inferred
79 })
80 }
81
82 fn pick_inner(
83 &self,
84 region_id: RegionId,
85 current_version: &CompactionVersion,
86 current_time: Timestamp,
87 ) -> (Vec<CompactionOutput>, Vec<FileHandle>, i64) {
88 let time_window = self.calculate_time_window(region_id, current_version);
89 info!(
90 "Compaction window for region: {} is {} seconds",
91 region_id, time_window
92 );
93
94 let expired_ssts = get_expired_ssts(
95 current_version.ssts.levels(),
96 current_version.options.ttl,
97 current_time,
98 );
99 if !expired_ssts.is_empty() {
100 info!("Expired SSTs in region {}: {:?}", region_id, expired_ssts);
101 }
102 let expired_file_ids = expired_ssts
103 .iter()
104 .map(|file| file.file_id())
105 .collect::<HashSet<_>>();
106
107 let last_non_null = current_version.options.merge_mode() == MergeMode::LastNonNull;
108 let windows = assign_files_to_time_windows(
111 time_window,
112 current_version
113 .ssts
114 .levels()
115 .iter()
116 .flat_map(|level| level.files.values())
117 .filter(|file| !expired_file_ids.contains(&file.file_id()))
118 .filter(|file| last_non_null || !file.compacting()),
119 );
120 let windows = filter_time_windows(windows, self.time_range);
121
122 if last_non_null
125 && windows.values().any(|(_, inputs)| {
126 inputs.iter().any(FileHandle::compacting)
127 || !inputs_precede_memtables(inputs, current_version.memtable_min_sequence)
128 })
129 {
130 return (vec![], expired_ssts, time_window);
131 }
132
133 (build_output(windows), expired_ssts, time_window)
134 }
135}
136
137#[async_trait::async_trait]
138impl Picker for WindowedCompactionPicker {
139 async fn pick(&self, compaction_region: &CompactionRegion) -> Result<Option<PickerOutput>> {
140 let picker = self.clone();
146 let region_id = compaction_region.current_version.metadata.region_id;
147 let current_version = compaction_region.current_version.clone();
148 let CompactionOptions::Twcs(options) = &compaction_region.region_options.compaction;
149 let max_file_size = options
151 .max_output_file_size
152 .map(|size| size.as_bytes())
153 .filter(|size| *size > 0)
154 .map(|size| size as usize);
155 let (outputs, expired_ssts, time_window) =
156 common_runtime::spawn_blocking_compact(move || {
157 picker.pick_inner(region_id, ¤t_version, Timestamp::current_millis())
158 })
159 .await
160 .context(JoinSnafu)?;
161
162 Ok(Some(PickerOutput {
163 outputs,
164 expired_ssts,
165 time_window_size: time_window,
166 max_file_size,
167 }))
168 }
169}
170
171fn filter_time_windows(
179 mut windows: BTreeMap<i64, (i64, Vec<FileHandle>)>,
180 time_range: Option<TimestampRange>,
181) -> BTreeMap<i64, (i64, Vec<FileHandle>)> {
182 let Some(time_range) = time_range else {
183 return windows;
184 };
185
186 let mut selected_windows = windows
187 .iter()
188 .filter_map(|(lower_bound, (upper_bound, _))| {
189 let window_start = Timestamp::new_second(*lower_bound);
190 let window_end = Timestamp::new_second(*upper_bound);
191 let starts_before_range_end = time_range
192 .end()
193 .is_none_or(|range_end| window_start < range_end);
194 let ends_after_range_start = time_range
195 .start()
196 .is_none_or(|range_start| range_start < window_end);
197 (starts_before_range_end && ends_after_range_start).then_some(*lower_bound)
198 })
199 .collect::<HashSet<_>>();
200
201 let mut file_windows = HashMap::new();
202 for (lower_bound, (_, files)) in &windows {
203 for file in files {
204 file_windows
205 .entry(file.file_id())
206 .or_insert_with(Vec::new)
207 .push(*lower_bound);
208 }
209 }
210
211 let mut pending_windows = selected_windows.iter().copied().collect::<VecDeque<_>>();
212 let mut visited_files = HashSet::new();
213 while let Some(lower_bound) = pending_windows.pop_front() {
214 let (_, files) = &windows[&lower_bound];
215 for file in files {
216 if !visited_files.insert(file.file_id()) {
217 continue;
218 }
219 for dependent_window in &file_windows[&file.file_id()] {
220 if selected_windows.insert(*dependent_window) {
221 pending_windows.push_back(*dependent_window);
222 }
223 }
224 }
225 }
226
227 windows.retain(|lower_bound, _| selected_windows.contains(lower_bound));
228 windows
229}
230
231fn build_output(windows: BTreeMap<i64, (i64, Vec<FileHandle>)>) -> Vec<CompactionOutput> {
232 let mut outputs = Vec::with_capacity(windows.len());
233 for (lower_bound, (upper_bound, files)) in windows {
234 let output_time_range = Some(
236 TimestampRange::new(
237 Timestamp::new_second(lower_bound),
238 Timestamp::new_second(upper_bound),
239 )
240 .unwrap(),
241 );
242
243 let output = CompactionOutput {
244 output_level: 1,
245 inputs: files,
246 filter_deleted: false,
247 output_time_range,
248 };
249 outputs.push(output);
250 }
251
252 outputs
253}
254
255fn assign_files_to_time_windows<'a>(
259 bucket_sec: i64,
260 files: impl Iterator<Item = &'a FileHandle>,
261) -> BTreeMap<i64, (i64, Vec<FileHandle>)> {
262 let mut buckets = BTreeMap::new();
263
264 for file in files {
265 let (start, end) = file.time_range();
266 let bounds = file_time_bucket_span(
267 start.convert_to(TimeUnit::Second).unwrap().value(),
269 end.convert_to(TimeUnit::Second).unwrap().value(),
270 bucket_sec,
271 );
272 for (lower_bound, upper_bound) in bounds {
273 let (_, files) = buckets
274 .entry(lower_bound)
275 .or_insert_with(|| (upper_bound, Vec::new()));
276 files.push(file.clone());
277 }
278 }
279 buckets
280}
281
282fn file_time_bucket_span(start_sec: i64, end_sec: i64, bucket_sec: i64) -> Vec<(i64, i64)> {
284 assert!(start_sec <= end_sec);
285
286 let mut start_aligned = start_sec.align_by_bucket(bucket_sec).unwrap_or(i64::MIN);
289 let end_aligned = end_sec
290 .align_by_bucket(bucket_sec)
291 .unwrap_or(start_aligned + (end_sec - start_sec));
292
293 let mut res = Vec::with_capacity(((end_aligned - start_aligned) / bucket_sec + 1) as usize);
294 while start_aligned <= end_aligned {
295 let window_size = if start_aligned % bucket_sec == 0 {
296 bucket_sec
297 } else {
298 (start_aligned % bucket_sec).abs()
299 };
300 let upper_bound = start_aligned.checked_add(window_size).unwrap_or(i64::MAX);
301 res.push((start_aligned, upper_bound));
302 start_aligned = upper_bound;
303 }
304 res
305}
306
307#[cfg(test)]
308mod tests {
309 use std::collections::HashSet;
310 use std::sync::Arc;
311 use std::time::Duration;
312
313 use common_base::readable_size::ReadableSize;
314 use common_time::Timestamp;
315 use common_time::range::TimestampRange;
316 use store_api::storage::{FileId, RegionId};
317
318 use crate::compaction::compactor::CompactionVersion;
319 use crate::compaction::picker::Picker;
320 use crate::compaction::test_util::compaction_region_with_ssts;
321 use crate::compaction::window::{WindowedCompactionPicker, file_time_bucket_span};
322 use crate::region::options::{CompactionOptions, MergeMode, RegionOptions};
323 use crate::sst::file::{FileMeta, Level};
324 use crate::sst::file_purger::NoopFilePurger;
325 use crate::sst::version::SstVersion;
326 use crate::test_util::memtable_util::metadata_for_test;
327
328 fn build_version(
329 files: &[(FileId, i64, i64, Level)],
330 ttl: Option<Duration>,
331 ) -> CompactionVersion {
332 let metadata = metadata_for_test();
333 let file_purger_ref = Arc::new(NoopFilePurger);
334
335 let mut ssts = SstVersion::new(metadata.clone());
336
337 ssts.add_files(
338 file_purger_ref,
339 files.iter().map(|(file_id, start, end, level)| FileMeta {
340 file_id: *file_id,
341 time_range: (
342 Timestamp::new_millisecond(*start),
343 Timestamp::new_millisecond(*end),
344 ),
345 level: *level,
346 ..Default::default()
347 }),
348 );
349
350 CompactionVersion {
351 metadata,
352 ssts: Arc::new(ssts),
353 memtable_min_sequence: None,
354 options: RegionOptions {
355 ttl: ttl.map(|t| t.into()),
356 auto_flush_interval: None,
357 compaction: Default::default(),
358 compaction_override: false,
359 storage: None,
360 append_mode: false,
361 skip_wal: false,
362 wal_options: Default::default(),
363 index_options: Default::default(),
364 memtable: None,
365 merge_mode: None,
366 sst_format: None,
367 max_row_group_row_count: None,
368 primary_key_encoding: None,
369 write_buffer_size: None,
370 preserve_row_sequence: false,
371 float_field_encoding: Default::default(),
372 },
373 compaction_time_window: None,
374 }
375 }
376
377 #[tokio::test]
378 async fn test_pick_output_file_size_threshold() {
379 let mut region = compaction_region_with_ssts(
380 [FileMeta {
381 file_id: FileId::random(),
382 time_range: (Timestamp::new_second(0), Timestamp::new_second(1)),
383 ..Default::default()
384 }],
385 Duration::from_secs(3600),
386 )
387 .await;
388 let picker = WindowedCompactionPicker::new(Some(3600));
389
390 for (size, expected) in [(None, None), (Some(0), None), (Some(1024), Some(1024))] {
391 let CompactionOptions::Twcs(options) = &mut region.region_options.compaction;
392 options.max_output_file_size = size.map(ReadableSize);
393
394 let output = picker.pick(®ion).await.unwrap().unwrap();
395 assert_eq!(1, output.outputs.len());
396 assert_eq!(expected, output.max_file_size);
397 }
398 }
399
400 #[test]
401 fn test_pick_expired_ssts_without_marking_compacting() {
402 let picker = WindowedCompactionPicker::new(None);
403 let files = vec![(FileId::random(), 0, 10, 0)];
404 let version = build_version(&files, Some(Duration::from_millis(1)));
405 let (outputs, expired_ssts, _) = picker.pick_inner(
406 RegionId::new(0, 0),
407 &version,
408 Timestamp::new_millisecond(12),
409 );
410
411 assert!(outputs.is_empty());
412 assert_eq!(1, expired_ssts.len());
413 assert!(expired_ssts.iter().all(|file| !file.compacting()));
414 }
415
416 const HOUR: i64 = 60 * 60 * 1000;
417
418 #[test]
419 fn test_infer_window() {
420 let picker = WindowedCompactionPicker::new(None);
421
422 let files = vec![
423 (FileId::random(), 0, HOUR, 0),
424 (FileId::random(), HOUR, HOUR * 2 - 1, 0),
425 ];
426
427 let version = build_version(&files, Some(Duration::from_millis(3 * HOUR as u64)));
428
429 let (outputs, expired_ssts, window_seconds) = picker.pick_inner(
430 RegionId::new(0, 0),
431 &version,
432 Timestamp::new_millisecond(HOUR * 2),
433 );
434 assert!(expired_ssts.is_empty());
435 assert_eq!(2 * HOUR / 1000, window_seconds);
436 assert_eq!(1, outputs.len());
437 assert_eq!(2, outputs[0].inputs.len());
438 }
439
440 #[test]
441 fn test_assign_files_to_windows() {
442 let picker = WindowedCompactionPicker::new(Some(HOUR / 1000));
443 let files = vec![
444 (FileId::random(), 0, 2 * HOUR - 1, 0),
445 (FileId::random(), HOUR, HOUR * 3 - 1, 0),
446 ];
447 let version = build_version(&files, Some(Duration::from_millis(3 * HOUR as u64)));
448 let (outputs, expired_ssts, window_seconds) = picker.pick_inner(
449 RegionId::new(0, 0),
450 &version,
451 Timestamp::new_millisecond(HOUR * 3),
452 );
453
454 assert!(expired_ssts.is_empty());
455 assert_eq!(HOUR / 1000, window_seconds);
456 assert_eq!(3, outputs.len());
457
458 assert_eq!(1, outputs[0].inputs.len());
459 assert_eq!(files[0].0, outputs[0].inputs[0].file_id().file_id());
460 assert_eq!(
461 TimestampRange::new(
462 Timestamp::new_millisecond(0),
463 Timestamp::new_millisecond(HOUR)
464 ),
465 outputs[0].output_time_range
466 );
467
468 assert_eq!(2, outputs[1].inputs.len());
469 assert_eq!(
470 TimestampRange::new(
471 Timestamp::new_millisecond(HOUR),
472 Timestamp::new_millisecond(2 * HOUR)
473 ),
474 outputs[1].output_time_range
475 );
476
477 assert_eq!(1, outputs[2].inputs.len());
478 assert_eq!(files[1].0, outputs[2].inputs[0].file_id().file_id());
479 assert_eq!(
480 TimestampRange::new(
481 Timestamp::new_millisecond(2 * HOUR),
482 Timestamp::new_millisecond(3 * HOUR)
483 ),
484 outputs[2].output_time_range
485 );
486 }
487
488 #[test]
489 fn test_pick_time_range_expands_for_cross_window_files() {
490 let time_range = TimestampRange::new(
491 Timestamp::new_millisecond(HOUR / 2),
492 Timestamp::new_millisecond(HOUR * 3 / 4),
493 )
494 .unwrap();
495 let picker =
496 WindowedCompactionPicker::new(Some(HOUR / 1000)).with_time_range(Some(time_range));
497 let files = vec![
498 (FileId::random(), 0, 2 * HOUR - 1, 0),
499 (FileId::random(), HOUR, HOUR * 3 - 1, 0),
500 (FileId::random(), 4 * HOUR, 5 * HOUR - 1, 0),
501 ];
502 let version = build_version(&files, None);
503
504 let (outputs, _, _) = picker.pick_inner(
505 RegionId::new(0, 0),
506 &version,
507 Timestamp::new_millisecond(6 * HOUR),
508 );
509
510 assert_eq!(3, outputs.len());
511 assert_eq!(
512 Some(TimestampRange::new(
513 Timestamp::new_millisecond(0),
514 Timestamp::new_millisecond(HOUR),
515 )),
516 outputs.first().map(|output| output.output_time_range)
517 );
518 assert_eq!(
519 Some(TimestampRange::new(
520 Timestamp::new_millisecond(2 * HOUR),
521 Timestamp::new_millisecond(3 * HOUR),
522 )),
523 outputs.last().map(|output| output.output_time_range)
524 );
525 }
526
527 #[test]
528 fn test_pick_time_range_expands_long_dependency_chain() {
529 const CHAIN_LEN: i64 = 128;
530
531 let time_range = TimestampRange::new(
532 Timestamp::new_millisecond(0),
533 Timestamp::new_millisecond(HOUR / 2),
534 )
535 .unwrap();
536 let picker =
537 WindowedCompactionPicker::new(Some(HOUR / 1000)).with_time_range(Some(time_range));
538 let files = (0..CHAIN_LEN)
539 .map(|window| (FileId::random(), window * HOUR, (window + 2) * HOUR - 1, 0))
540 .collect::<Vec<_>>();
541 let version = build_version(&files, None);
542
543 let (outputs, _, _) = picker.pick_inner(
544 RegionId::new(0, 0),
545 &version,
546 Timestamp::new_millisecond((CHAIN_LEN + 2) * HOUR),
547 );
548
549 assert_eq!(CHAIN_LEN as usize + 1, outputs.len());
550 }
551
552 #[test]
553 fn test_assign_compacting_files_to_windows() {
554 let picker = WindowedCompactionPicker::new(Some(HOUR / 1000));
555 let files = vec![
556 (FileId::random(), 0, 2 * HOUR - 1, 0),
557 (FileId::random(), HOUR, HOUR * 3 - 1, 0),
558 ];
559 let version = build_version(&files, Some(Duration::from_millis(3 * HOUR as u64)));
560 version.ssts.levels()[0]
561 .files()
562 .for_each(|f| f.set_compacting(true));
563 let (outputs, expired_ssts, window_seconds) = picker.pick_inner(
564 RegionId::new(0, 0),
565 &version,
566 Timestamp::new_millisecond(HOUR * 3),
567 );
568
569 assert!(expired_ssts.is_empty());
570 assert_eq!(HOUR / 1000, window_seconds);
571 assert!(outputs.is_empty());
572 }
573
574 #[rstest::rstest]
575 #[case::requested_window(0)]
576 #[case::transitive_window(1)]
577 #[case::unrelated_window(3)]
578 fn test_busy_dependencies_defer_last_non_null_plan_until_released(
579 #[case] busy_window: i64,
580 #[values(MergeMode::LastRow, MergeMode::LastNonNull)] merge_mode: MergeMode,
581 ) {
582 let [a, b, c, expired] = std::array::from_fn(|_| FileId::random());
586 let files = [
587 (a, 0, 2 * HOUR - 1, 0),
588 (b, busy_window * HOUR, (busy_window + 1) * HOUR - 1, 1),
589 (c, 0, 2 * HOUR - 1, 0),
590 (expired, -2 * HOUR, -2 * HOUR, 0),
591 ];
592 let mut version = build_version(&files, Some(Duration::from_millis(3 * HOUR as u64)));
593 version.options.merge_mode = Some(merge_mode);
594 assert_eq!(None, version.memtable_min_sequence);
596 let busy_file = version.ssts.levels()[1].files().next().unwrap();
597 let picker =
598 WindowedCompactionPicker::new(Some(HOUR / 1000)).with_time_range(TimestampRange::new(
599 Timestamp::new_millisecond(0),
600 Timestamp::new_millisecond(HOUR),
601 ));
602
603 for busy in [true, false] {
604 busy_file.set_compacting(busy);
605 let (outputs, expired_ssts, _) = picker.pick_inner(
606 version.metadata.region_id,
607 &version,
608 Timestamp::new_millisecond(2 * HOUR),
609 );
610 assert_eq!(1, expired_ssts.len());
611 assert_eq!(expired, expired_ssts[0].meta_ref().file_id);
612 assert!(!expired_ssts[0].compacting());
613
614 if busy && merge_mode == MergeMode::LastNonNull && busy_window < 2 {
615 assert!(outputs.is_empty(), "the entire dependent plan must wait");
616 continue;
617 }
618
619 assert_eq!(2, outputs.len());
620 for (window, output) in outputs.iter().enumerate() {
621 let window = window as i64;
622 let mut expected = HashSet::from([a, c]);
623 if !busy && busy_window == window {
624 expected.insert(b);
625 }
626 let actual: HashSet<_> = output
627 .inputs
628 .iter()
629 .map(|file| file.meta_ref().file_id)
630 .collect();
631 assert_eq!(expected, actual);
632 assert_eq!(
633 TimestampRange::new(
634 Timestamp::new_millisecond(window * HOUR),
635 Timestamp::new_millisecond((window + 1) * HOUR),
636 ),
637 output.output_time_range
638 );
639 }
640 }
641 }
642
643 #[test]
644 fn test_file_time_bucket_span() {
645 assert_eq!(
646 vec![(i64::MIN, i64::MIN + 8),],
647 file_time_bucket_span(i64::MIN, i64::MIN + 1, 10)
648 );
649
650 assert_eq!(
651 vec![(i64::MIN, i64::MIN + 8), (i64::MIN + 8, i64::MIN + 18)],
652 file_time_bucket_span(i64::MIN, i64::MIN + 8, 10)
653 );
654
655 assert_eq!(
656 vec![
657 (i64::MIN, i64::MIN + 8),
658 (i64::MIN + 8, i64::MIN + 18),
659 (i64::MIN + 18, i64::MIN + 28)
660 ],
661 file_time_bucket_span(i64::MIN, i64::MIN + 20, 10)
662 );
663
664 assert_eq!(
665 vec![(-10, 0), (0, 10), (10, 20)],
666 file_time_bucket_span(-1, 11, 10)
667 );
668
669 assert_eq!(
670 vec![(-3, 0), (0, 3), (3, 6)],
671 file_time_bucket_span(-1, 3, 3)
672 );
673
674 assert_eq!(vec![(0, 10)], file_time_bucket_span(0, 9, 10));
675
676 assert_eq!(
677 vec![(i64::MAX - (i64::MAX % 10), i64::MAX)],
678 file_time_bucket_span(i64::MAX - 1, i64::MAX, 10)
679 );
680 }
681}