1use std::collections::{HashMap, HashSet};
16
17use snafu::ResultExt;
18use store_api::metadata::RegionMetadataRef;
19use store_api::storage::SequenceNumber;
20
21use crate::compaction::CompactionOutput;
22use crate::compaction::compactor::CompactionRegion;
23use crate::compaction::overlap::FileOverlapIndex;
24use crate::compaction::picker::{Picker, PickerOutput};
25use crate::compaction::twcs::TwcsPicker;
26use crate::error::{JoinSnafu, Result};
27use crate::sst::file::{FileHandle, RegionFileId};
28
29#[derive(Debug)]
45pub(super) struct LastNonNullPicker {
46 base_picker: TwcsPicker,
47}
48
49impl LastNonNullPicker {
50 pub(super) fn new(base_picker: TwcsPicker) -> Self {
51 Self { base_picker }
52 }
53}
54
55#[async_trait::async_trait]
56impl Picker for LastNonNullPicker {
57 async fn pick(&self, region: &CompactionRegion) -> Result<Option<PickerOutput>> {
58 let max_outputs = self.base_picker.max_background_tasks;
59
60 let Some(mut all_seeds) = self
62 .base_picker
63 .pick_with_output_limit(region, None)
64 .await?
65 else {
66 return Ok(None);
67 };
68 let version = region.current_version.clone();
69 common_runtime::spawn_blocking_compact(move || {
70 let expired: HashSet<_> = all_seeds
71 .expired_ssts
72 .iter()
73 .map(FileHandle::file_id)
74 .collect();
75 let all_files = version
78 .ssts
79 .levels()
80 .iter()
81 .flat_map(|level| level.files())
82 .filter(|file| !expired.contains(&file.file_id()))
83 .cloned()
84 .collect();
85 all_seeds.outputs =
86 build_closed_outputs(all_seeds.outputs, all_files, &version.metadata);
87 all_seeds.outputs.retain(|output| {
88 inputs_precede_memtables(&output.inputs, version.memtable_min_sequence)
89 });
90 if let Some(limit) = max_outputs {
91 let excess = all_seeds.outputs.len().saturating_sub(limit);
94 all_seeds.outputs.drain(..excess);
95 }
96 (!all_seeds.outputs.is_empty() || !all_seeds.expired_ssts.is_empty())
97 .then_some(all_seeds)
98 })
99 .await
100 .context(JoinSnafu)
101 }
102}
103
104pub(super) fn inputs_precede_memtables(
107 inputs: &[FileHandle],
108 memtable_min_sequence: Option<SequenceNumber>,
109) -> bool {
110 memtable_min_sequence.is_none_or(|min_sequence| {
111 inputs.iter().all(|file| {
112 file.meta_ref()
113 .sequence
114 .is_some_and(|max_sequence| max_sequence.get() < min_sequence)
115 })
116 })
117}
118
119fn build_closed_outputs(
123 seeds: Vec<CompactionOutput>,
124 files: Vec<FileHandle>,
125 metadata: &RegionMetadataRef,
126) -> Vec<CompactionOutput> {
127 let seed_by_file: HashMap<_, _> = seeds
130 .iter()
131 .enumerate()
132 .flat_map(|(i, seed)| seed.inputs.iter().map(move |file| (file.file_id(), i)))
133 .collect();
134 let mut seeds: Vec<_> = seeds.into_iter().map(Some).collect();
136 let mut remaining = FileOverlapIndex::new(files, metadata);
137 let mut outputs = Vec::new();
138 for seed_index in (0..seeds.len()).rev() {
140 let Some(output) = seeds[seed_index].take() else {
141 continue;
142 };
143 let output = expand_closure(output, &mut seeds, &seed_by_file, &mut remaining);
144 if !output.inputs.iter().any(FileHandle::compacting) {
146 outputs.push(output);
147 }
148 }
149 outputs.reverse();
150 outputs
151}
152
153fn expand_closure(
157 mut output: CompactionOutput,
158 seeds: &mut [Option<CompactionOutput>],
159 seed_by_file: &HashMap<RegionFileId, usize>,
160 remaining_files: &mut FileOverlapIndex<'_>,
161) -> CompactionOutput {
162 let mut selected = HashSet::new();
163 let inputs = std::mem::take(&mut output.inputs);
164 extend_inputs(&mut output, inputs, &mut selected, remaining_files);
165
166 let mut cursor = 0;
167 while cursor < output.inputs.len() {
168 let input = output.inputs[cursor].clone();
169 if let Some(&seed_index) = seed_by_file.get(&input.file_id())
170 && let Some(seed) = seeds[seed_index].take()
171 {
172 output.filter_deleted &= seed.filter_deleted;
174 extend_inputs(&mut output, seed.inputs, &mut selected, remaining_files);
175 }
176 let dependencies = remaining_files.drain_overlaps(&input);
177 extend_inputs(&mut output, dependencies, &mut selected, remaining_files);
178 cursor += 1;
179 }
180 output
181}
182
183fn extend_inputs(
186 output: &mut CompactionOutput,
187 inputs: Vec<FileHandle>,
188 selected: &mut HashSet<RegionFileId>,
189 remaining: &mut FileOverlapIndex<'_>,
190) {
191 for input in inputs {
192 if selected.insert(input.file_id()) {
193 remaining.remove(input.file_id());
194 output.inputs.push(input);
195 }
196 }
197}
198
199#[cfg(test)]
200mod tests {
201 use std::time::Duration;
202
203 use api::v1::region::compact_request;
204 use common_time::Timestamp;
205 use common_time::range::TimestampRange;
206 use store_api::storage::FileId;
207
208 use super::*;
209 use crate::compaction::picker::new_picker;
210 use crate::compaction::test_util::{
211 compaction_region_with_ssts, new_file_handle, new_file_handle_with_size_and_sequence,
212 };
213 use crate::region::options::{CompactionOptions, MergeMode};
214 use crate::sst::file::{FileMeta, RegionFileId};
215 use crate::test_util::memtable_util::metadata_for_test;
216
217 fn seed(inputs: Vec<FileHandle>) -> CompactionOutput {
218 CompactionOutput {
219 output_level: 1,
220 inputs,
221 filter_deleted: false,
222 output_time_range: None,
223 }
224 }
225
226 fn file_ids(files: &[FileHandle]) -> HashSet<RegionFileId> {
227 files.iter().map(FileHandle::file_id).collect()
228 }
229
230 #[rstest::rstest]
231 #[case::memtable_barrier(2, Some(20), false, Some(1), vec![vec![1]])]
232 #[case::busy_dependency(2, None, true, Some(1), vec![vec![1]])]
233 #[case::highest_priority(2, None, false, Some(1), vec![vec![2]])]
234 #[case::priority_order(2, None, false, Some(2), vec![vec![2], vec![1]])]
235 #[case::unlimited(2, None, false, None, vec![vec![2], vec![1], vec![0]])]
236 #[case::fewer_eligible_outputs(2, Some(20), false, Some(3), vec![vec![1], vec![0]])]
237 #[case::absorbed_seeds(1, None, false, Some(2), vec![vec![2, 1], vec![0]])]
238 #[case::busy_bridge_defers_both_seeds(1, None, true, Some(2), vec![vec![0]])]
239 #[case::memtable_barrier_defers_both_seeds(1, Some(20), false, Some(2), vec![vec![0]])]
240 #[tokio::test]
241 async fn test_output_limit_counts_eligible_closures(
242 #[case] dependency_start_window: usize,
243 #[case] memtable_min: Option<u64>,
244 #[case] busy: bool,
245 #[case] limit: Option<usize>,
246 #[case] expected_windows: Vec<Vec<usize>>,
247 ) {
248 let windows: Vec<Vec<_>> = (0..3)
253 .map(|window| {
254 (0..4)
255 .map(|i| FileMeta {
256 file_id: FileId::random(),
257 time_range: (
258 Timestamp::new_second(window * 3600),
259 Timestamp::new_second(window * 3600 + 10),
260 ),
261 level: 0,
262 file_size: 100,
263 sequence: std::num::NonZeroU64::new((window * 4 + i + 1) as u64),
264 ..Default::default()
265 })
266 .collect()
267 })
268 .collect();
269 let dependency = FileMeta {
270 file_id: FileId::random(),
271 time_range: (
272 Timestamp::new_second(dependency_start_window as i64 * 3600),
273 Timestamp::new_second(7210),
274 ),
275 level: 1,
276 file_size: 1_000_000,
277 sequence: std::num::NonZeroU64::new(20),
278 ..Default::default()
279 };
280 let dependency_id = dependency.file_id();
281 let files = windows.iter().flatten().cloned().chain([dependency]);
282 let mut region = compaction_region_with_ssts(files, Duration::from_secs(60)).await;
283 region.ttl = None;
284 region.region_options.merge_mode = Some(MergeMode::LastNonNull);
285 region.current_version.memtable_min_sequence = memtable_min;
286 let CompactionOptions::Twcs(opts) = &mut region.region_options.compaction;
287 opts.time_window = Some(Duration::from_secs(3600));
288 opts.active_window_trigger_file_num = 4;
289 opts.inactive_window_trigger_file_num = 4;
290 region.current_version.ssts.levels()[1]
291 .files()
292 .next()
293 .unwrap()
294 .set_compacting(busy);
295 let request = compact_request::Options::Regular(Default::default());
296 let picker = new_picker(&request, ®ion.region_options, limit, None);
297
298 let picked = picker.pick(®ion).await.unwrap().unwrap();
299 assert_eq!(expected_windows.len(), picked.outputs.len());
300 assert!(picked.expired_ssts.is_empty());
301 for (output, window_indices) in picked.outputs.iter().rev().zip(expected_windows) {
303 let mut expected: HashSet<_> = window_indices
304 .iter()
305 .flat_map(|&i| windows[i].iter().map(FileMeta::file_id))
306 .collect();
307 if window_indices.contains(&2) {
308 expected.insert(dependency_id);
309 }
310 assert_eq!(expected, file_ids(&output.inputs));
311 }
312
313 region.region_options.merge_mode = Some(MergeMode::LastRow);
315 let picker = new_picker(&request, ®ion.region_options, limit, None);
316 let picked = picker.pick(®ion).await.unwrap().unwrap();
317 assert_eq!(limit.unwrap_or(3).min(3), picked.outputs.len());
318 for (output, window) in picked.outputs.iter().rev().zip(windows.iter().rev()) {
319 let expected: HashSet<_> = window.iter().map(FileMeta::file_id).collect();
320 assert_eq!(expected, file_ids(&output.inputs));
321 }
322 }
323
324 #[rstest::rstest]
327 #[case(0, None, true)]
328 #[case(0, Some(2), false)]
329 #[case(1, Some(0), false)]
330 #[case(1, Some(2), true)]
331 #[case(2, Some(2), false)]
332 #[case(3, Some(2), false)]
333 #[case(u64::MAX - 1, Some(u64::MAX), true)]
334 #[case(u64::MAX, Some(u64::MAX), false)]
335 fn test_inputs_must_precede_pending_memtables(
336 #[case] file_sequence: u64,
337 #[case] memtable_min: Option<u64>,
338 #[case] expected: bool,
339 ) {
340 let file =
341 new_file_handle_with_size_and_sequence(FileId::random(), 0, 10, 0, file_sequence, 100);
342 assert_eq!(expected, inputs_precede_memtables(&[file], memtable_min));
343 }
344
345 #[test]
346 fn test_transitive_closure_exceeds_seed_limit_across_levels_and_windows() {
347 let metadata = metadata_for_test();
348 let files: Vec<_> = (0..32)
351 .map(|i| {
352 new_file_handle(
353 FileId::random(),
354 i * 3_600_000,
355 (i + 1) * 3_600_000,
356 (i % 2) as u8,
357 )
358 })
359 .collect();
360 let outputs =
361 build_closed_outputs(vec![seed(vec![files[0].clone()])], files.clone(), &metadata);
362 assert_eq!(1, outputs.len());
363 assert_eq!(files.len(), outputs[0].inputs.len());
364 assert_eq!(file_ids(&files), file_ids(&outputs[0].inputs));
365 }
366
367 #[rstest::rstest]
368 #[case(false)]
369 #[case(true)]
370 fn test_absorbed_seed_closes_disconnected_members_and_defers_busy_group(#[case] busy: bool) {
371 let metadata = metadata_for_test();
372 let files: Vec<_> = [(0, 10), (10, 20), (100, 110), (110, 120), (200, 210)]
373 .into_iter()
374 .map(|(start, end)| new_file_handle(FileId::random(), start, end, 0))
375 .collect();
376 files[3].set_compacting(busy);
377 let mut priority_seed = seed(vec![files[0].clone()]);
378 priority_seed.filter_deleted = true;
379 let seeds = vec![
380 seed(vec![files[4].clone()]),
381 seed(vec![files[1].clone(), files[2].clone()]),
382 priority_seed,
383 ];
384
385 let mut outputs = build_closed_outputs(seeds, files.clone(), &metadata);
386 assert_eq!(if busy { 1 } else { 2 }, outputs.len());
387 if !busy {
388 let first = outputs.pop().unwrap();
389 assert_eq!(4, first.inputs.len());
390 assert_eq!(file_ids(&files[..4]), file_ids(&first.inputs));
391 assert!(!first.filter_deleted);
393 }
394 assert_eq!(
395 vec![files[4].file_id()],
396 outputs[0]
397 .inputs
398 .iter()
399 .map(FileHandle::file_id)
400 .collect::<Vec<_>>()
401 );
402 }
403
404 #[rstest::rstest]
405 #[case("chain", 2048)]
406 #[case("dense_with_external", 1024)]
407 #[case("time_dense_pk_disjoint", 1)]
408 fn test_large_snapshot_closure(#[case] shape: &str, #[case] expected_count: usize) {
409 use crate::compaction::test_util::{
410 new_file_handle_with_size_sequence_and_primary_key_range, pk_range,
411 primary_key_metadata_for_test,
412 };
413
414 let metadata = primary_key_metadata_for_test();
415 let files = (0..2048)
416 .map(|i| {
417 let (start, end, pk) = match shape {
418 "chain" => (i, i + 1, None),
419 "dense_with_external" if i < 1024 => (0, 1, None),
420 "dense_with_external" => (i, i, None),
421 "time_dense_pk_disjoint" => {
422 let key = i.to_string();
423 (0, 1, pk_range(key.as_bytes(), key.as_bytes()))
424 }
425 _ => unreachable!(),
426 };
427 new_file_handle_with_size_sequence_and_primary_key_range(
428 FileId::random(),
429 start,
430 end,
431 0,
432 i as u64 + 1,
433 100,
434 pk,
435 )
436 })
437 .collect::<Vec<_>>();
438 let outputs =
439 build_closed_outputs(vec![seed(vec![files[0].clone()])], files.clone(), &metadata);
440 assert_eq!(1, outputs.len());
441 assert_eq!(expected_count, outputs[0].inputs.len());
442 assert_eq!(
443 file_ids(&files[..expected_count]),
444 file_ids(&outputs[0].inputs)
445 );
446 }
447
448 #[rstest::rstest]
449 #[case(MergeMode::LastRow, false, None, 4)]
450 #[case(MergeMode::LastRow, true, None, 4)]
451 #[case(MergeMode::LastNonNull, false, None, 5)]
452 #[case(MergeMode::LastNonNull, true, None, 0)]
453 #[case(MergeMode::LastRow, false, Some(5), 4)]
454 #[case(MergeMode::LastNonNull, false, Some(5), 0)]
455 #[case(MergeMode::LastNonNull, false, Some(6), 5)]
456 #[tokio::test]
457 async fn test_picker_dependencies_outside_request_window(
458 #[case] merge_mode: MergeMode,
459 #[case] busy: bool,
460 #[case] memtable_min: Option<u64>,
461 #[case] expected_count: usize,
462 ) {
463 let files = (0..5).map(|i| FileMeta {
464 file_id: FileId::random(),
465 time_range: (
466 Timestamp::new_second(0),
467 Timestamp::new_second(if i == 4 { 3601 } else { 10 }),
468 ),
469 level: if i == 4 { 1 } else { 0 },
470 file_size: if i == 4 { 1_000_000 } else { 100 },
471 sequence: std::num::NonZeroU64::new(i + 1),
472 ..Default::default()
473 });
474 let mut region = compaction_region_with_ssts(files, Duration::from_secs(60)).await;
475 region.ttl = None;
476 region.region_options.merge_mode = Some(merge_mode);
477 region.current_version.memtable_min_sequence = memtable_min;
478 let CompactionOptions::Twcs(opts) = &mut region.region_options.compaction;
479 opts.time_window = Some(Duration::from_secs(3600));
480 let dependency = region.current_version.ssts.levels()[1]
481 .files()
482 .next()
483 .unwrap();
484 dependency.set_compacting(busy);
485 let picker = new_picker(
486 &compact_request::Options::Regular(Default::default()),
487 ®ion.region_options,
488 Some(1),
489 TimestampRange::new(Timestamp::new_second(0), Timestamp::new_second(3600)),
490 );
491
492 let picked = picker.pick(®ion).await.unwrap();
493 if expected_count == 0 {
494 assert!(picked.is_none());
495 } else {
496 let picked = picked.unwrap();
497 assert_eq!(1, picked.outputs.len());
498 assert_eq!(expected_count, picked.outputs[0].inputs.len());
499 assert_eq!(
500 merge_mode == MergeMode::LastNonNull,
501 picked.outputs[0]
502 .inputs
503 .iter()
504 .any(|file| file.file_id() == dependency.file_id())
505 );
506 }
507 }
508
509 #[rstest::rstest]
510 #[case(MergeMode::LastNonNull, None, 3)]
511 #[case(MergeMode::LastNonNull, Some(2), 0)]
512 #[case(MergeMode::LastRow, Some(2), 3)]
513 #[tokio::test]
514 async fn test_strict_window_keeps_all_versions_in_disjoint_output_slices(
515 #[case] merge_mode: MergeMode,
516 #[case] memtable_min: Option<u64>,
517 #[case] expected_outputs: usize,
518 ) {
519 let files: Vec<_> = [(0, 10), (10, 10), (10, 20)]
520 .into_iter()
521 .enumerate()
522 .map(|(i, (start, end))| {
523 new_file_handle_with_size_and_sequence(
524 FileId::random(),
525 start * 1000,
526 end * 1000,
527 0,
528 i as u64 + 1,
529 100,
530 )
531 .meta_ref()
532 .clone()
533 })
534 .collect();
535 let mut region = compaction_region_with_ssts(files, Duration::from_secs(60)).await;
536 region.region_options.merge_mode = Some(merge_mode);
537 region.current_version.options.merge_mode = Some(merge_mode);
538 region.current_version.memtable_min_sequence = memtable_min;
539 let picker = new_picker(
540 &compact_request::Options::StrictWindow(api::v1::region::StrictWindow {
541 window_seconds: 10,
542 }),
543 ®ion.region_options,
544 Some(1),
545 TimestampRange::new(Timestamp::new_second(10), Timestamp::new_second(20)),
546 );
547
548 let picked = picker.pick(®ion).await.unwrap().unwrap();
549 assert_eq!(expected_outputs, picked.outputs.len());
554 for (i, output) in picked.outputs.iter().enumerate() {
555 assert_eq!(
556 TimestampRange::new(
557 Timestamp::new_second(i as i64 * 10),
558 Timestamp::new_second((i as i64 + 1) * 10),
559 ),
560 output.output_time_range
561 );
562 assert_eq!(if i == 1 { 3 } else { 1 }, output.inputs.len());
563 assert!(!output.filter_deleted);
564 }
565 }
566
567 #[tokio::test]
568 async fn test_busy_closure_keeps_expired_files_without_rewriting_them() {
569 let now = Timestamp::current_millis().value();
570 let files = (0..6).map(|i| FileMeta {
571 file_id: FileId::random(),
572 time_range: (
573 Timestamp::new_millisecond(if i == 5 { 0 } else { now - 10_000 }),
574 Timestamp::new_millisecond(if i == 5 { 10 } else { now }),
575 ),
576 level: if i == 4 { 1 } else { 0 },
577 file_size: if i == 4 { 1_000_000 } else { 100 },
578 sequence: std::num::NonZeroU64::new(i + 1),
579 ..Default::default()
580 });
581 let mut region = compaction_region_with_ssts(files, Duration::from_secs(3600)).await;
582 region.region_options.merge_mode = Some(MergeMode::LastNonNull);
583 let dependency = region.current_version.ssts.levels()[1]
584 .files()
585 .next()
586 .unwrap();
587 dependency.set_compacting(true);
588 let picker = new_picker(
589 &compact_request::Options::Regular(Default::default()),
590 ®ion.region_options,
591 Some(1),
592 None,
593 );
594 let picked = picker.pick(®ion).await.unwrap().unwrap();
595 assert!(picked.outputs.is_empty());
596 assert_eq!(1, picked.expired_ssts.len());
597 assert_eq!(
598 Timestamp::new_millisecond(10),
599 picked.expired_ssts[0].time_range().1
600 );
601 assert!(!picked.expired_ssts[0].compacting());
602
603 dependency.set_compacting(false);
604 let picked = picker.pick(®ion).await.unwrap().unwrap();
605 assert_eq!(1, picked.outputs.len());
606 assert_eq!(5, picked.outputs[0].inputs.len());
607 assert_eq!(1, picked.expired_ssts.len());
608 assert!(!file_ids(&picked.outputs[0].inputs).contains(&picked.expired_ssts[0].file_id()));
609 }
610}