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::picker::{Picker, PickerOutput, get_expired_ssts};
30use crate::error::{JoinSnafu, Result};
31use crate::sst::file::FileHandle;
32
33#[derive(Clone, Debug)]
37pub struct WindowedCompactionPicker {
38 compaction_time_window_seconds: Option<i64>,
39 time_range: Option<TimestampRange>,
40}
41
42impl WindowedCompactionPicker {
43 pub fn new(window_seconds: Option<i64>) -> Self {
44 Self {
45 compaction_time_window_seconds: window_seconds,
46 time_range: None,
47 }
48 }
49
50 pub(crate) fn with_time_range(mut self, time_range: Option<TimestampRange>) -> Self {
52 self.time_range = time_range;
53 self
54 }
55
56 fn calculate_time_window(
61 &self,
62 region_id: RegionId,
63 current_version: &CompactionVersion,
64 ) -> i64 {
65 self.compaction_time_window_seconds
66 .or(current_version
67 .compaction_time_window
68 .map(|t| t.as_secs() as i64))
69 .unwrap_or_else(|| {
70 let levels = current_version.ssts.levels();
71 let inferred = infer_time_bucket(levels[0].files());
72 info!(
73 "Compaction window for region {} is not present, inferring from files: {:?}",
74 region_id, inferred
75 );
76 inferred
77 })
78 }
79
80 fn pick_inner(
81 &self,
82 region_id: RegionId,
83 current_version: &CompactionVersion,
84 current_time: Timestamp,
85 ) -> (Vec<CompactionOutput>, Vec<FileHandle>, i64) {
86 let time_window = self.calculate_time_window(region_id, current_version);
87 info!(
88 "Compaction window for region: {} is {} seconds",
89 region_id, time_window
90 );
91
92 let expired_ssts = get_expired_ssts(
93 current_version.ssts.levels(),
94 current_version.options.ttl,
95 current_time,
96 );
97 if !expired_ssts.is_empty() {
98 info!("Expired SSTs in region {}: {:?}", region_id, expired_ssts);
99 }
100 let expired_file_ids = expired_ssts
101 .iter()
102 .map(|file| file.file_id())
103 .collect::<HashSet<_>>();
104
105 let windows = assign_files_to_time_windows(
106 time_window,
107 current_version
108 .ssts
109 .levels()
110 .iter()
111 .flat_map(|level| level.files.values())
112 .filter(|file| !expired_file_ids.contains(&file.file_id())),
113 );
114 let windows = filter_time_windows(windows, self.time_range);
115
116 (build_output(windows), expired_ssts, time_window)
117 }
118}
119
120#[async_trait::async_trait]
121impl Picker for WindowedCompactionPicker {
122 async fn pick(&self, compaction_region: &CompactionRegion) -> Result<Option<PickerOutput>> {
123 let picker = self.clone();
124 let region_id = compaction_region.current_version.metadata.region_id;
125 let current_version = compaction_region.current_version.clone();
126 let (outputs, expired_ssts, time_window) =
127 common_runtime::spawn_blocking_compact(move || {
128 picker.pick_inner(region_id, ¤t_version, Timestamp::current_millis())
129 })
130 .await
131 .context(JoinSnafu)?;
132
133 Ok(Some(PickerOutput {
134 outputs,
135 expired_ssts,
136 time_window_size: time_window,
137 max_file_size: None, }))
139 }
140}
141
142fn filter_time_windows(
150 mut windows: BTreeMap<i64, (i64, Vec<FileHandle>)>,
151 time_range: Option<TimestampRange>,
152) -> BTreeMap<i64, (i64, Vec<FileHandle>)> {
153 let Some(time_range) = time_range else {
154 return windows;
155 };
156
157 let mut selected_windows = windows
158 .iter()
159 .filter_map(|(lower_bound, (upper_bound, _))| {
160 let window_start = Timestamp::new_second(*lower_bound);
161 let window_end = Timestamp::new_second(*upper_bound);
162 let starts_before_range_end = time_range
163 .end()
164 .is_none_or(|range_end| window_start < range_end);
165 let ends_after_range_start = time_range
166 .start()
167 .is_none_or(|range_start| range_start < window_end);
168 (starts_before_range_end && ends_after_range_start).then_some(*lower_bound)
169 })
170 .collect::<HashSet<_>>();
171
172 let mut file_windows = HashMap::new();
173 for (lower_bound, (_, files)) in &windows {
174 for file in files {
175 file_windows
176 .entry(file.file_id())
177 .or_insert_with(Vec::new)
178 .push(*lower_bound);
179 }
180 }
181
182 let mut pending_windows = selected_windows.iter().copied().collect::<VecDeque<_>>();
183 let mut visited_files = HashSet::new();
184 while let Some(lower_bound) = pending_windows.pop_front() {
185 let (_, files) = &windows[&lower_bound];
186 for file in files {
187 if !visited_files.insert(file.file_id()) {
188 continue;
189 }
190 for dependent_window in &file_windows[&file.file_id()] {
191 if selected_windows.insert(*dependent_window) {
192 pending_windows.push_back(*dependent_window);
193 }
194 }
195 }
196 }
197
198 windows.retain(|lower_bound, _| selected_windows.contains(lower_bound));
199 windows
200}
201
202fn build_output(windows: BTreeMap<i64, (i64, Vec<FileHandle>)>) -> Vec<CompactionOutput> {
203 let mut outputs = Vec::with_capacity(windows.len());
204 for (lower_bound, (upper_bound, files)) in windows {
205 let output_time_range = Some(
207 TimestampRange::new(
208 Timestamp::new_second(lower_bound),
209 Timestamp::new_second(upper_bound),
210 )
211 .unwrap(),
212 );
213
214 let output = CompactionOutput {
215 output_level: 1,
216 inputs: files,
217 filter_deleted: false,
218 output_time_range,
219 };
220 outputs.push(output);
221 }
222
223 outputs
224}
225
226fn assign_files_to_time_windows<'a>(
230 bucket_sec: i64,
231 files: impl Iterator<Item = &'a FileHandle>,
232) -> BTreeMap<i64, (i64, Vec<FileHandle>)> {
233 let mut buckets = BTreeMap::new();
234
235 for file in files {
236 if file.compacting() {
237 continue;
238 }
239 let (start, end) = file.time_range();
240 let bounds = file_time_bucket_span(
241 start.convert_to(TimeUnit::Second).unwrap().value(),
243 end.convert_to(TimeUnit::Second).unwrap().value(),
244 bucket_sec,
245 );
246 for (lower_bound, upper_bound) in bounds {
247 let (_, files) = buckets
248 .entry(lower_bound)
249 .or_insert_with(|| (upper_bound, Vec::new()));
250 files.push(file.clone());
251 }
252 }
253 buckets
254}
255
256fn file_time_bucket_span(start_sec: i64, end_sec: i64, bucket_sec: i64) -> Vec<(i64, i64)> {
258 assert!(start_sec <= end_sec);
259
260 let mut start_aligned = start_sec.align_by_bucket(bucket_sec).unwrap_or(i64::MIN);
263 let end_aligned = end_sec
264 .align_by_bucket(bucket_sec)
265 .unwrap_or(start_aligned + (end_sec - start_sec));
266
267 let mut res = Vec::with_capacity(((end_aligned - start_aligned) / bucket_sec + 1) as usize);
268 while start_aligned <= end_aligned {
269 let window_size = if start_aligned % bucket_sec == 0 {
270 bucket_sec
271 } else {
272 (start_aligned % bucket_sec).abs()
273 };
274 let upper_bound = start_aligned.checked_add(window_size).unwrap_or(i64::MAX);
275 res.push((start_aligned, upper_bound));
276 start_aligned = upper_bound;
277 }
278 res
279}
280
281#[cfg(test)]
282mod tests {
283 use std::sync::Arc;
284 use std::time::Duration;
285
286 use common_time::Timestamp;
287 use common_time::range::TimestampRange;
288 use store_api::storage::{FileId, RegionId};
289
290 use crate::compaction::compactor::CompactionVersion;
291 use crate::compaction::window::{WindowedCompactionPicker, file_time_bucket_span};
292 use crate::region::options::RegionOptions;
293 use crate::sst::file::{FileMeta, Level};
294 use crate::sst::file_purger::NoopFilePurger;
295 use crate::sst::version::SstVersion;
296 use crate::test_util::memtable_util::metadata_for_test;
297
298 fn build_version(
299 files: &[(FileId, i64, i64, Level)],
300 ttl: Option<Duration>,
301 ) -> CompactionVersion {
302 let metadata = metadata_for_test();
303 let file_purger_ref = Arc::new(NoopFilePurger);
304
305 let mut ssts = SstVersion::new();
306
307 ssts.add_files(
308 file_purger_ref,
309 files.iter().map(|(file_id, start, end, level)| FileMeta {
310 file_id: *file_id,
311 time_range: (
312 Timestamp::new_millisecond(*start),
313 Timestamp::new_millisecond(*end),
314 ),
315 level: *level,
316 ..Default::default()
317 }),
318 );
319
320 CompactionVersion {
321 metadata,
322 ssts: Arc::new(ssts),
323 options: RegionOptions {
324 ttl: ttl.map(|t| t.into()),
325 auto_flush_interval: None,
326 compaction: Default::default(),
327 compaction_override: false,
328 storage: None,
329 append_mode: false,
330 skip_wal: false,
331 wal_options: Default::default(),
332 index_options: Default::default(),
333 memtable: None,
334 merge_mode: None,
335 sst_format: None,
336 max_row_group_row_count: None,
337 primary_key_encoding: None,
338 write_buffer_size: None,
339 },
340 compaction_time_window: None,
341 }
342 }
343
344 #[test]
345 fn test_pick_expired_ssts_without_marking_compacting() {
346 let picker = WindowedCompactionPicker::new(None);
347 let files = vec![(FileId::random(), 0, 10, 0)];
348 let version = build_version(&files, Some(Duration::from_millis(1)));
349 let (outputs, expired_ssts, _) = picker.pick_inner(
350 RegionId::new(0, 0),
351 &version,
352 Timestamp::new_millisecond(12),
353 );
354
355 assert!(outputs.is_empty());
356 assert_eq!(1, expired_ssts.len());
357 assert!(expired_ssts.iter().all(|file| !file.compacting()));
358 }
359
360 const HOUR: i64 = 60 * 60 * 1000;
361
362 #[test]
363 fn test_infer_window() {
364 let picker = WindowedCompactionPicker::new(None);
365
366 let files = vec![
367 (FileId::random(), 0, HOUR, 0),
368 (FileId::random(), HOUR, HOUR * 2 - 1, 0),
369 ];
370
371 let version = build_version(&files, Some(Duration::from_millis(3 * HOUR as u64)));
372
373 let (outputs, expired_ssts, window_seconds) = picker.pick_inner(
374 RegionId::new(0, 0),
375 &version,
376 Timestamp::new_millisecond(HOUR * 2),
377 );
378 assert!(expired_ssts.is_empty());
379 assert_eq!(2 * HOUR / 1000, window_seconds);
380 assert_eq!(1, outputs.len());
381 assert_eq!(2, outputs[0].inputs.len());
382 }
383
384 #[test]
385 fn test_assign_files_to_windows() {
386 let picker = WindowedCompactionPicker::new(Some(HOUR / 1000));
387 let files = vec![
388 (FileId::random(), 0, 2 * HOUR - 1, 0),
389 (FileId::random(), HOUR, HOUR * 3 - 1, 0),
390 ];
391 let version = build_version(&files, Some(Duration::from_millis(3 * HOUR as u64)));
392 let (outputs, expired_ssts, window_seconds) = picker.pick_inner(
393 RegionId::new(0, 0),
394 &version,
395 Timestamp::new_millisecond(HOUR * 3),
396 );
397
398 assert!(expired_ssts.is_empty());
399 assert_eq!(HOUR / 1000, window_seconds);
400 assert_eq!(3, outputs.len());
401
402 assert_eq!(1, outputs[0].inputs.len());
403 assert_eq!(files[0].0, outputs[0].inputs[0].file_id().file_id());
404 assert_eq!(
405 TimestampRange::new(
406 Timestamp::new_millisecond(0),
407 Timestamp::new_millisecond(HOUR)
408 ),
409 outputs[0].output_time_range
410 );
411
412 assert_eq!(2, outputs[1].inputs.len());
413 assert_eq!(
414 TimestampRange::new(
415 Timestamp::new_millisecond(HOUR),
416 Timestamp::new_millisecond(2 * HOUR)
417 ),
418 outputs[1].output_time_range
419 );
420
421 assert_eq!(1, outputs[2].inputs.len());
422 assert_eq!(files[1].0, outputs[2].inputs[0].file_id().file_id());
423 assert_eq!(
424 TimestampRange::new(
425 Timestamp::new_millisecond(2 * HOUR),
426 Timestamp::new_millisecond(3 * HOUR)
427 ),
428 outputs[2].output_time_range
429 );
430 }
431
432 #[test]
433 fn test_pick_time_range_expands_for_cross_window_files() {
434 let time_range = TimestampRange::new(
435 Timestamp::new_millisecond(HOUR / 2),
436 Timestamp::new_millisecond(HOUR * 3 / 4),
437 )
438 .unwrap();
439 let picker =
440 WindowedCompactionPicker::new(Some(HOUR / 1000)).with_time_range(Some(time_range));
441 let files = vec![
442 (FileId::random(), 0, 2 * HOUR - 1, 0),
443 (FileId::random(), HOUR, HOUR * 3 - 1, 0),
444 (FileId::random(), 4 * HOUR, 5 * HOUR - 1, 0),
445 ];
446 let version = build_version(&files, None);
447
448 let (outputs, _, _) = picker.pick_inner(
449 RegionId::new(0, 0),
450 &version,
451 Timestamp::new_millisecond(6 * HOUR),
452 );
453
454 assert_eq!(3, outputs.len());
455 assert_eq!(
456 Some(TimestampRange::new(
457 Timestamp::new_millisecond(0),
458 Timestamp::new_millisecond(HOUR),
459 )),
460 outputs.first().map(|output| output.output_time_range)
461 );
462 assert_eq!(
463 Some(TimestampRange::new(
464 Timestamp::new_millisecond(2 * HOUR),
465 Timestamp::new_millisecond(3 * HOUR),
466 )),
467 outputs.last().map(|output| output.output_time_range)
468 );
469 }
470
471 #[test]
472 fn test_pick_time_range_expands_long_dependency_chain() {
473 const CHAIN_LEN: i64 = 128;
474
475 let time_range = TimestampRange::new(
476 Timestamp::new_millisecond(0),
477 Timestamp::new_millisecond(HOUR / 2),
478 )
479 .unwrap();
480 let picker =
481 WindowedCompactionPicker::new(Some(HOUR / 1000)).with_time_range(Some(time_range));
482 let files = (0..CHAIN_LEN)
483 .map(|window| (FileId::random(), window * HOUR, (window + 2) * HOUR - 1, 0))
484 .collect::<Vec<_>>();
485 let version = build_version(&files, None);
486
487 let (outputs, _, _) = picker.pick_inner(
488 RegionId::new(0, 0),
489 &version,
490 Timestamp::new_millisecond((CHAIN_LEN + 2) * HOUR),
491 );
492
493 assert_eq!(CHAIN_LEN as usize + 1, outputs.len());
494 }
495
496 #[test]
497 fn test_assign_compacting_files_to_windows() {
498 let picker = WindowedCompactionPicker::new(Some(HOUR / 1000));
499 let files = vec![
500 (FileId::random(), 0, 2 * HOUR - 1, 0),
501 (FileId::random(), HOUR, HOUR * 3 - 1, 0),
502 ];
503 let version = build_version(&files, Some(Duration::from_millis(3 * HOUR as u64)));
504 version.ssts.levels()[0]
505 .files()
506 .for_each(|f| f.set_compacting(true));
507 let (outputs, expired_ssts, window_seconds) = picker.pick_inner(
508 RegionId::new(0, 0),
509 &version,
510 Timestamp::new_millisecond(HOUR * 3),
511 );
512
513 assert!(expired_ssts.is_empty());
514 assert_eq!(HOUR / 1000, window_seconds);
515 assert!(outputs.is_empty());
516 }
517
518 #[test]
519 fn test_file_time_bucket_span() {
520 assert_eq!(
521 vec![(i64::MIN, i64::MIN + 8),],
522 file_time_bucket_span(i64::MIN, i64::MIN + 1, 10)
523 );
524
525 assert_eq!(
526 vec![(i64::MIN, i64::MIN + 8), (i64::MIN + 8, i64::MIN + 18)],
527 file_time_bucket_span(i64::MIN, i64::MIN + 8, 10)
528 );
529
530 assert_eq!(
531 vec![
532 (i64::MIN, i64::MIN + 8),
533 (i64::MIN + 8, i64::MIN + 18),
534 (i64::MIN + 18, i64::MIN + 28)
535 ],
536 file_time_bucket_span(i64::MIN, i64::MIN + 20, 10)
537 );
538
539 assert_eq!(
540 vec![(-10, 0), (0, 10), (10, 20)],
541 file_time_bucket_span(-1, 11, 10)
542 );
543
544 assert_eq!(
545 vec![(-3, 0), (0, 3), (3, 6)],
546 file_time_bucket_span(-1, 3, 3)
547 );
548
549 assert_eq!(vec![(0, 10)], file_time_bucket_span(0, 9, 10));
550
551 assert_eq!(
552 vec![(i64::MAX - (i64::MAX % 10), i64::MAX)],
553 file_time_bucket_span(i64::MAX - 1, i64::MAX, 10)
554 );
555 }
556}