1use std::num::NonZeroU64;
16use std::sync::Arc;
17use std::time::Duration;
18
19use common_base::Plugins;
20use common_base::cancellation::{CancellableFuture, CancellationHandle};
21use common_meta::key::SchemaMetadataManagerRef;
22use common_telemetry::{debug, info, warn};
23use common_time::TimeToLive;
24use either::Either;
25use itertools::Itertools;
26use object_store::manager::ObjectStoreManagerRef;
27use partition::expr::PartitionExpr;
28use serde::{Deserialize, Serialize};
29use snafu::{OptionExt, ResultExt};
30use store_api::ManifestVersion;
31use store_api::metadata::RegionMetadataRef;
32use store_api::region_request::PathType;
33use store_api::storage::{RegionId, SequenceNumber};
34
35use crate::access_layer::{
36 AccessLayer, AccessLayerRef, Metrics, OperationType, SstWriteRequest, WriteType,
37};
38use crate::cache::{CacheManager, CacheManagerRef};
39use crate::compaction::picker::PickerOutput;
40use crate::compaction::reader::CompactionSstReaderBuilder;
41use crate::compaction::{CompactionOutput, find_dynamic_options};
42use crate::config::MitoConfig;
43use crate::engine::region_hook::{RegionHookRef, SstFileInfo};
44use crate::error;
45use crate::error::{
46 EmptyRegionDirSnafu, InvalidPartitionExprSnafu, ObjectStoreNotFoundSnafu, Result,
47};
48use crate::manifest::action::{RegionEdit, RegionMetaAction, RegionMetaActionList};
49use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions};
50use crate::region::options::{MergeMode, RegionOptions};
51use crate::region::version::VersionRef;
52use crate::region::{ManifestContext, RegionLeaderState, RegionRoleState};
53use crate::schedule::scheduler::LocalScheduler;
54use crate::sst::FormatType;
55use crate::sst::file::{FileHandle, FileMeta, UncommittedSsts};
56use crate::sst::file_purger::LocalFilePurger;
57use crate::sst::index::intermediate::IntermediateManager;
58use crate::sst::index::puffin_manager::PuffinManagerFactory;
59use crate::sst::location::region_dir_from_table_dir;
60use crate::sst::parquet::metadata::extract_primary_key_range;
61use crate::sst::parquet::{SstInfo, WriteOptions};
62use crate::sst::version::{SstVersion, SstVersionRef};
63
64#[derive(Clone)]
66pub struct CompactionVersion {
67 pub(crate) metadata: RegionMetadataRef,
72 pub(crate) options: RegionOptions,
74 pub(crate) ssts: SstVersionRef,
76 pub(crate) memtable_min_sequence: Option<SequenceNumber>,
80 pub(crate) compaction_time_window: Option<Duration>,
82}
83
84impl From<VersionRef> for CompactionVersion {
85 fn from(value: VersionRef) -> Self {
86 Self {
87 metadata: value.metadata.clone(),
88 options: value.options.clone(),
89 ssts: value.ssts.clone(),
90 memtable_min_sequence: if value.options.merge_mode() == MergeMode::LastNonNull {
91 value.memtables.min_sequence()
92 } else {
93 None
94 },
95 compaction_time_window: value.compaction_time_window,
96 }
97 }
98}
99
100#[derive(Clone)]
103pub struct CompactionRegion {
104 pub region_id: RegionId,
105 pub region_options: RegionOptions,
106
107 pub(crate) engine_config: Arc<MitoConfig>,
108 pub(crate) region_metadata: RegionMetadataRef,
109 pub(crate) cache_manager: CacheManagerRef,
110 pub access_layer: AccessLayerRef,
112 pub(crate) manifest_ctx: Arc<ManifestContext>,
113 pub(crate) current_version: CompactionVersion,
114 pub(crate) file_purger: Option<Arc<LocalFilePurger>>,
115 pub(crate) ttl: Option<TimeToLive>,
116
117 pub max_parallelism: usize,
122
123 pub(crate) plugins: Plugins,
124}
125
126#[derive(Clone)]
128pub struct OpenCompactionRegionRequest {
129 pub region_id: RegionId,
130 pub table_dir: String,
131 pub path_type: PathType,
132 pub region_options: RegionOptions,
133 pub max_parallelism: usize,
134 pub plugins: Plugins,
136}
137
138pub async fn open_compaction_region(
141 req: &OpenCompactionRegionRequest,
142 mito_config: &MitoConfig,
143 object_store_manager: ObjectStoreManagerRef,
144 ttl_provider: Either<TimeToLive, SchemaMetadataManagerRef>,
145) -> Result<CompactionRegion> {
146 let object_store = {
147 let name = &req.region_options.storage;
148 if let Some(name) = name {
149 object_store_manager
150 .find(name)
151 .with_context(|| ObjectStoreNotFoundSnafu {
152 object_store: name.clone(),
153 })?
154 } else {
155 object_store_manager.default_object_store()
156 }
157 };
158
159 let access_layer = {
160 let puffin_manager_factory = PuffinManagerFactory::new(
161 &mito_config.index.aux_path,
162 mito_config.index.staging_size.as_bytes(),
163 Some(mito_config.index.write_buffer_size.as_bytes() as _),
164 mito_config.index.staging_ttl,
165 )
166 .await?;
167 let intermediate_manager =
168 IntermediateManager::init_fs(mito_config.index.aux_path.clone()).await?;
169
170 Arc::new(AccessLayer::new(
171 &req.table_dir,
172 req.path_type,
173 object_store.clone(),
174 puffin_manager_factory,
175 intermediate_manager,
176 ))
177 };
178
179 let manifest_manager = {
180 let region_dir = region_dir_from_table_dir(&req.table_dir, req.region_id, req.path_type);
181 let region_manifest_options =
182 RegionManifestOptions::new(mito_config, ®ion_dir, object_store);
183
184 RegionManifestManager::open(region_manifest_options, &Default::default())
185 .await?
186 .with_context(|| EmptyRegionDirSnafu {
187 region_id: req.region_id,
188 region_dir: region_dir_from_table_dir(&req.table_dir, req.region_id, req.path_type),
189 })?
190 };
191
192 let manifest = manifest_manager.manifest();
193 let region_metadata = manifest.metadata.clone();
194 let hook: Option<RegionHookRef> = req.plugins.get();
195 let manifest_ctx = Arc::new(ManifestContext::new(
196 manifest_manager,
197 RegionRoleState::Leader(RegionLeaderState::Writable),
198 hook,
199 ));
200
201 let file_purger = {
202 let purge_scheduler = Arc::new(LocalScheduler::new(mito_config.max_background_purges));
203 Arc::new(LocalFilePurger::new(
204 purge_scheduler.clone(),
205 access_layer.clone(),
206 None,
207 ))
208 };
209
210 let current_version = {
211 let mut ssts = SstVersion::new(region_metadata.clone());
212 ssts.add_files(file_purger.clone(), manifest.files.values().cloned());
213 CompactionVersion {
214 metadata: region_metadata.clone(),
215 options: req.region_options.clone(),
216 ssts: Arc::new(ssts),
217 memtable_min_sequence: Some(0),
220 compaction_time_window: manifest.compaction_time_window,
221 }
222 };
223
224 let ttl = match ttl_provider {
225 Either::Left(ttl) => ttl,
227 Either::Right(schema_metadata_manager) => {
229 let (_, ttl) =
230 find_dynamic_options(req.region_id, &req.region_options, &schema_metadata_manager)
231 .await
232 .unwrap_or_else(|e| {
233 warn!(e; "Failed to get ttl for region: {}", region_metadata.region_id);
234 (
235 crate::region::options::CompactionOptions::default(),
236 TimeToLive::default(),
237 )
238 });
239 ttl
240 }
241 };
242
243 Ok(CompactionRegion {
244 region_id: req.region_id,
245 region_options: req.region_options.clone(),
246 engine_config: Arc::new(mito_config.clone()),
247 region_metadata: region_metadata.clone(),
248 cache_manager: Arc::new(CacheManager::default()),
249 access_layer,
250 manifest_ctx,
251 current_version,
252 file_purger: Some(file_purger),
253 ttl: Some(ttl),
254 max_parallelism: req.max_parallelism,
255 plugins: req.plugins.clone(),
256 })
257}
258
259impl CompactionRegion {
260 pub fn file_purger(&self) -> Option<Arc<LocalFilePurger>> {
262 self.file_purger.clone()
263 }
264
265 pub async fn stop_purger_scheduler(&self) -> Result<()> {
267 if let Some(file_purger) = &self.file_purger {
268 file_purger.stop_scheduler().await
269 } else {
270 Ok(())
271 }
272 }
273
274 pub async fn invoke_sst_hook(&self, merge_output: &MergeOutput) {
279 let Some(hook) = self.plugins.get::<RegionHookRef>() else {
280 return;
281 };
282
283 let synthesized: Vec<SstInfo>;
288 let infos: &[SstInfo] = if merge_output.sst_infos.is_empty() {
289 synthesized = merge_output
290 .files_to_add
291 .iter()
292 .map(sst_info_from_file_meta)
293 .collect();
294 &synthesized
295 } else {
296 if merge_output.sst_infos.len() != merge_output.files_to_add.len() {
300 warn!(
301 "sst_infos length ({}) does not match files_to_add length ({}) for region {}",
302 merge_output.sst_infos.len(),
303 merge_output.files_to_add.len(),
304 self.region_id
305 );
306 }
307 &merge_output.sst_infos
308 };
309
310 let files: Vec<SstFileInfo<'_>> = merge_output
311 .files_to_add
312 .iter()
313 .zip(infos)
314 .map(|(meta, info)| SstFileInfo {
315 sst_info_ref: info,
316 file_meta: meta,
317 })
318 .collect();
319 hook.on_sst_files_written(self.region_id, &self.region_metadata, &files)
320 .await;
321 }
322}
323
324fn sst_info_from_file_meta(meta: &FileMeta) -> SstInfo {
347 SstInfo {
348 file_id: meta.file_id,
349 time_range: meta.time_range,
350 file_size: meta.file_size,
351 max_row_group_uncompressed_size: meta.max_row_group_uncompressed_size,
352 num_rows: meta.num_rows as usize,
353 num_row_groups: meta.num_row_groups,
354 num_series: meta.num_series,
355 ..Default::default()
356 }
357}
358
359#[derive(Default, Debug, Serialize, Deserialize)]
361pub struct MergeOutput {
362 pub files_to_add: Vec<FileMeta>,
363 pub files_to_remove: Vec<FileMeta>,
364 pub compaction_time_window: Option<i64>,
365 #[serde(skip)]
366 pub sst_infos: Vec<SstInfo>,
367}
368
369impl MergeOutput {
370 pub fn is_empty(&self) -> bool {
371 self.files_to_add.is_empty() && self.files_to_remove.is_empty()
372 }
373
374 pub fn input_file_size(&self) -> u64 {
375 self.files_to_remove.iter().map(|f| f.file_size).sum()
376 }
377
378 pub fn output_file_size(&self) -> u64 {
379 self.files_to_add.iter().map(|f| f.file_size).sum()
380 }
381}
382
383#[async_trait::async_trait]
385pub trait Compactor: Send + Sync + 'static {
386 async fn merge_ssts(
388 &self,
389 compaction_region: &CompactionRegion,
390 picker_output: PickerOutput,
391 ) -> Result<MergeOutput>;
392
393 async fn update_manifest(
395 &self,
396 compaction_region: &CompactionRegion,
397 merge_output: MergeOutput,
398 ) -> Result<(RegionEdit, ManifestVersion)>;
399}
400
401#[async_trait::async_trait]
406pub trait SstMerger: Send + Sync + 'static {
407 async fn merge_single_output(
408 &self,
409 compaction_region: CompactionRegion,
410 output: CompactionOutput,
411 write_opts: WriteOptions,
412 ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)>;
413}
414
415#[derive(Clone)]
417pub struct DefaultSstMerger;
418
419struct OutputSequenceMetadata {
421 file_sequence: Option<NonZeroU64>,
422 exact_sequence_trusted: bool,
423}
424
425fn output_sequence_metadata(
429 region: &CompactionRegion,
430 inputs: &[FileHandle],
431) -> OutputSequenceMetadata {
432 let exact_sequence_trusted = region.region_options.preserve_row_sequence
433 && inputs
434 .iter()
435 .all(|f| f.is_effective_target_sequence_trusted(region.region_id));
436 OutputSequenceMetadata {
437 file_sequence: known_max_input_sequence(inputs),
438 exact_sequence_trusted,
439 }
440}
441
442fn known_max_input_sequence(inputs: &[FileHandle]) -> Option<NonZeroU64> {
447 let mut max: Option<NonZeroU64> = None;
448 for input in inputs {
449 let sequence = input.meta_ref().sequence?;
450 max = Some(max.map_or(sequence, |current| current.max(sequence)));
451 }
452 max
453}
454
455#[async_trait::async_trait]
456impl SstMerger for DefaultSstMerger {
457 async fn merge_single_output(
458 &self,
459 compaction_region: CompactionRegion,
460 output: CompactionOutput,
461 write_opts: WriteOptions,
462 ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
463 let region_id = compaction_region.region_id;
464 let storage = compaction_region.region_options.storage.clone();
465 let index_options = compaction_region
466 .current_version
467 .options
468 .index_options
469 .clone();
470 let append_mode = compaction_region.current_version.options.append_mode;
471 let merge_mode = compaction_region.current_version.options.merge_mode();
472 let flat_format = compaction_region
473 .region_options
474 .sst_format
475 .map(|format| format == FormatType::Flat)
476 .unwrap_or(compaction_region.engine_config.default_flat_format);
477
478 let index_config = compaction_region.engine_config.index.clone();
479 let inverted_index_config = compaction_region.engine_config.inverted_index.clone();
480 let fulltext_index_config = compaction_region.engine_config.fulltext_index.clone();
481 let bloom_filter_index_config = compaction_region.engine_config.bloom_filter_index.clone();
482 #[cfg(feature = "vector_index")]
483 let vector_index_config = compaction_region.engine_config.vector_index.clone();
484
485 let input_file_names = output
486 .inputs
487 .iter()
488 .map(|f| f.file_id().to_string())
489 .join(",");
490 let sequence_metadata = output_sequence_metadata(&compaction_region, &output.inputs);
491 let builder = CompactionSstReaderBuilder {
492 metadata: compaction_region.region_metadata.clone(),
493 sst_layer: compaction_region.access_layer.clone(),
494 cache: compaction_region.cache_manager.clone(),
495 inputs: &output.inputs,
496 append_mode,
497 filter_deleted: output.filter_deleted,
498 time_range: output.output_time_range,
499 merge_mode,
500 };
501 let source = builder.build_flat_sst_reader().await?;
502
503 let mut metrics = Metrics::new(WriteType::Compaction);
504 let region_metadata = compaction_region.region_metadata.clone();
505 let sst_infos = compaction_region
506 .access_layer
507 .write_sst(
508 SstWriteRequest {
509 op_type: OperationType::Compact,
510 metadata: region_metadata.clone(),
511 source,
512 cache_manager: compaction_region.cache_manager.clone(),
513 storage,
514 max_sequence: None,
518 sst_write_format: if flat_format {
519 FormatType::Flat
520 } else {
521 FormatType::PrimaryKey
522 },
523 preserve_row_sequence: sequence_metadata.exact_sequence_trusted,
524 index_options,
525 index_config,
526 inverted_index_config,
527 fulltext_index_config,
528 bloom_filter_index_config,
529 #[cfg(feature = "vector_index")]
530 vector_index_config,
531 },
532 &write_opts,
533 &mut metrics,
534 )
535 .await?;
536 let partition_expr = match ®ion_metadata.partition_expr {
538 None => None,
539 Some(json_str) if json_str.is_empty() => None,
540 Some(json_str) => PartitionExpr::from_json_str(json_str).with_context(|_| {
541 InvalidPartitionExprSnafu {
542 expr: json_str.clone(),
543 }
544 })?,
545 };
546
547 let output_files = sst_infos
548 .iter()
549 .map(|sst_info| {
550 let pk_range = sst_info
551 .file_metadata
552 .as_ref()
553 .and_then(|meta| extract_primary_key_range(meta, ®ion_metadata));
554 let (primary_key_min, primary_key_max) = match pk_range {
555 Some((min, max)) => (Some(min), Some(max)),
556 None => (None, None),
557 };
558
559 FileMeta {
560 region_id,
561 file_id: sst_info.file_id,
562 time_range: sst_info.time_range,
563 level: output.output_level,
564 file_size: sst_info.file_size,
565 max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
566 available_indexes: sst_info.index_metadata.build_available_indexes(),
567 indexes: sst_info.index_metadata.build_indexes(),
568 index_file_size: sst_info.index_metadata.file_size,
569 index_version: 0,
570 num_rows: sst_info.num_rows as u64,
571 num_row_groups: sst_info.num_row_groups,
572 sequence: sequence_metadata.file_sequence,
573 partition_expr: partition_expr.clone(),
574 num_series: sst_info.num_series,
575 primary_key_min,
576 primary_key_max,
577 preserve_row_sequence: sequence_metadata.exact_sequence_trusted,
578 }
579 })
580 .collect::<Vec<_>>();
581 let output_file_names = output_files.iter().map(|f| f.file_id.to_string()).join(",");
582 info!(
583 "Region {} compaction inputs: [{}], outputs: [{}], flat_format: {}, metrics: {:?}",
584 region_id, input_file_names, output_file_names, flat_format, metrics
585 );
586 metrics.observe();
587 Ok((output_files, sst_infos.into_iter().collect()))
588 }
589}
590
591pub struct DefaultCompactor<M = DefaultSstMerger> {
596 merger: M,
597 cancel_handle: Arc<CancellationHandle>,
598 uncommitted: Option<UncommittedSsts>,
599}
600
601#[cfg(test)]
602impl<M: SstMerger> DefaultCompactor<M> {
603 pub fn with_merger(merger: M) -> Self {
604 Self {
605 merger,
606 cancel_handle: Arc::new(CancellationHandle::default()),
607 uncommitted: None,
608 }
609 }
610}
611
612impl DefaultCompactor {
613 pub fn new_with_cancel_handle(cancel_handle: Arc<CancellationHandle>) -> Self {
619 Self {
620 merger: DefaultSstMerger,
621 cancel_handle,
622 uncommitted: None,
623 }
624 }
625
626 pub(crate) fn with_cancel_handle(
627 cancel_handle: Arc<CancellationHandle>,
628 uncommitted: UncommittedSsts,
629 ) -> Self {
630 Self {
631 merger: DefaultSstMerger,
632 cancel_handle,
633 uncommitted: Some(uncommitted),
634 }
635 }
636}
637
638#[async_trait::async_trait]
639impl<M: SstMerger> Compactor for DefaultCompactor<M>
640where
641 M: Clone,
642{
643 async fn merge_ssts(
644 &self,
645 compaction_region: &CompactionRegion,
646 mut picker_output: PickerOutput,
647 ) -> Result<MergeOutput> {
648 let internal_parallelism = compaction_region.max_parallelism.max(1);
649 let compaction_time_window = picker_output.time_window_size;
650 let region_id = compaction_region.region_id;
651
652 let mut tasks: Vec<(Vec<FileMeta>, _)> = Vec::with_capacity(picker_output.outputs.len());
655
656 for output in picker_output.outputs.drain(..) {
657 let inputs_to_remove: Vec<_> =
658 output.inputs.iter().map(|f| f.meta_ref().clone()).collect();
659 let write_opts = WriteOptions {
660 write_buffer_size: compaction_region.engine_config.sst_write_buffer_size,
661 max_file_size: picker_output.max_file_size,
662 row_group_size: compaction_region.region_options.row_group_size(),
663 float_field_encoding: compaction_region.region_options.float_field_encoding,
664 };
665 let merger = self.merger.clone();
666 let compaction_region = compaction_region.clone();
667 let uncommitted = self.uncommitted.clone();
668 let fut = async move {
669 let result = merger
670 .merge_single_output(compaction_region, output, write_opts)
671 .await;
672 if let (Some(uncommitted), Ok((_, infos))) = (&uncommitted, &result) {
673 uncommitted.track(infos);
674 }
675 result
676 };
677 tasks.push((inputs_to_remove, fut));
678 }
679
680 let hook: Option<RegionHookRef> = compaction_region.plugins.get();
681 let mut output_files = Vec::with_capacity(tasks.len());
682 let mut all_sst_infos: Vec<SstInfo> = Vec::new();
683 let mut compacted_inputs = Vec::with_capacity(
684 tasks.iter().map(|(inputs, _)| inputs.len()).sum::<usize>()
685 + picker_output.expired_ssts.len(),
686 );
687
688 while !tasks.is_empty() {
689 let mut chunk: Vec<(Vec<FileMeta>, _)> = Vec::with_capacity(internal_parallelism);
690 for _ in 0..internal_parallelism {
691 if let Some(task) = tasks.pop() {
692 chunk.push(task);
693 }
694 }
695 let mut spawned: Vec<_> = chunk
696 .into_iter()
697 .map(|(inputs, fut)| {
698 let handle = common_runtime::spawn_compact(fut);
699 (inputs, handle)
700 })
701 .collect();
702
703 while let Some((inputs, mut handle)) = spawned.pop() {
704 let abort_handle = handle.abort_handle();
705 match CancellableFuture::new(&mut handle, self.cancel_handle.clone()).await {
706 Ok(Ok(Ok((files, infos)))) => {
707 output_files.extend(files);
708 if hook.is_some() {
709 all_sst_infos.extend(infos);
710 }
711 compacted_inputs.extend(inputs);
712 }
713 Ok(Ok(Err(e))) => {
714 warn!(
715 e; "Failed to merge compaction output for region: {}, inputs: [{}]",
716 region_id,
717 inputs.iter().map(|f| f.file_id.to_string()).join(",")
718 );
719 }
720 Ok(Err(e)) => {
721 warn!(
722 "Region {} compaction task join error for inputs: [{}], skipping: {}",
723 region_id,
724 inputs.iter().map(|f| f.file_id.to_string()).join(","),
725 e
726 );
727 if self.cancel_handle.is_cancelled() {
730 abort_handle.abort();
731 for (_, handle) in &spawned {
732 handle.abort();
733 }
734 for (_, handle) in spawned {
735 let _ = handle.await;
736 }
737 }
738 return Err(e).context(error::JoinSnafu);
739 }
740 Err(_) => {
741 debug!(
742 "Compaction merge cancelled for region: {}, aborting remaining {} spawned tasks",
743 region_id,
744 spawned.len(),
745 );
746 abort_handle.abort();
747 for (_, handle) in &spawned {
748 handle.abort();
749 }
750 let _ = handle.await;
751 for (_, handle) in spawned {
752 let _ = handle.await;
753 }
754 break;
755 }
756 }
757 }
758
759 if self.cancel_handle.is_cancelled() {
760 info!("Compaction merge cancelled for region: {}", region_id);
761 break;
762 }
763 }
764
765 compacted_inputs.extend(
767 picker_output
768 .expired_ssts
769 .iter()
770 .map(|f| f.meta_ref().clone()),
771 );
772
773 Ok(MergeOutput {
774 files_to_add: output_files,
775 files_to_remove: compacted_inputs,
776 compaction_time_window: Some(compaction_time_window),
777 sst_infos: all_sst_infos,
778 })
779 }
780
781 async fn update_manifest(
782 &self,
783 compaction_region: &CompactionRegion,
784 merge_output: MergeOutput,
785 ) -> Result<(RegionEdit, ManifestVersion)> {
786 let edit = RegionEdit {
788 files_to_add: merge_output.files_to_add,
789 files_to_remove: merge_output.files_to_remove,
790 timestamp_ms: Some(chrono::Utc::now().timestamp_millis()),
792 compaction_time_window: merge_output
793 .compaction_time_window
794 .map(|seconds| Duration::from_secs(seconds as u64)),
795 flushed_entry_id: None,
796 flushed_sequence: None,
797 committed_sequence: None,
798 };
799
800 let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit.clone()));
801 let manifest_version = compaction_region
803 .manifest_ctx
804 .update_manifest_for_compaction(action_list)
805 .await?;
806
807 Ok((edit, manifest_version))
808 }
809}
810
811#[cfg(test)]
812mod tests {
813 use std::sync::atomic::{AtomicUsize, Ordering};
814 use std::sync::{Arc, Mutex};
815 use std::time::Duration;
816
817 use store_api::storage::{FileId, RegionId};
818 use tokio::time::sleep;
819
820 use super::{DefaultCompactor, *};
821 use crate::cache::CacheManager;
822 use crate::compaction::picker::PickerOutput;
823 use crate::error::Result;
824 use crate::sst::file::FileHandle;
825 use crate::sst::file_purger::NoopFilePurger;
826 use crate::sst::version::SstVersion;
827 use crate::test_util::memtable_util::metadata_for_test;
828 use crate::test_util::scheduler_util::SchedulerEnv;
829
830 fn dummy_file_meta() -> FileMeta {
831 FileMeta {
832 region_id: RegionId::new(1, 1),
833 file_id: FileId::random(),
834 file_size: 100,
835 ..Default::default()
836 }
837 }
838
839 fn new_file_handle(meta: FileMeta) -> FileHandle {
840 FileHandle::new(meta, Arc::new(NoopFilePurger))
841 }
842
843 #[test]
844 fn test_effective_target_trust_and_max_sequence() {
845 let target = RegionId::new(1, 1);
846 let local_trusted = new_file_handle(FileMeta {
847 region_id: target,
848 sequence: NonZeroU64::new(3),
849 preserve_row_sequence: true,
850 ..dummy_file_meta()
851 });
852 let foreign_untrusted_marker = new_file_handle(FileMeta {
853 region_id: RegionId::new(1, 2),
854 sequence: NonZeroU64::new(9),
855 preserve_row_sequence: false,
856 ..dummy_file_meta()
857 });
858 let local_untrusted = new_file_handle(FileMeta {
859 region_id: target,
860 sequence: NonZeroU64::new(11),
861 preserve_row_sequence: false,
862 ..dummy_file_meta()
863 });
864 assert!(local_trusted.is_effective_target_sequence_trusted(target));
865 assert!(foreign_untrusted_marker.is_effective_target_sequence_trusted(target));
866 assert!(!local_untrusted.is_effective_target_sequence_trusted(target));
867 assert_eq!(
868 NonZeroU64::new(9),
869 known_max_input_sequence(&[local_trusted, foreign_untrusted_marker])
870 );
871 }
872
873 #[test]
874 fn test_known_max_input_sequence() {
875 let meta_with = |sequence| FileMeta {
876 sequence,
877 ..dummy_file_meta()
878 };
879
880 let inputs = vec![
881 new_file_handle(meta_with(NonZeroU64::new(3))),
882 new_file_handle(meta_with(NonZeroU64::new(9))),
883 new_file_handle(meta_with(NonZeroU64::new(5))),
884 ];
885 assert_eq!(NonZeroU64::new(9), known_max_input_sequence(&inputs));
886
887 let inputs = vec![
888 new_file_handle(meta_with(NonZeroU64::new(9))),
889 new_file_handle(meta_with(None)),
890 ];
891 assert_eq!(None, known_max_input_sequence(&inputs));
892
893 assert_eq!(None, known_max_input_sequence(&[]));
894 }
895
896 #[tokio::test]
897 async fn test_output_sequence_metadata_preserves_trust_and_unknown_bounds() {
898 let mut region = new_test_compaction_region().await;
899 let local = region.region_id;
900 let foreign = RegionId::new(1, 2);
901 let file = |region_id, sequence: Option<u64>, preserve_row_sequence| {
902 new_file_handle(FileMeta {
903 region_id,
904 sequence: sequence.and_then(NonZeroU64::new),
905 preserve_row_sequence,
906 ..dummy_file_meta()
907 })
908 };
909 let cases = [
912 (
913 vec![file(local, Some(3), true), file(foreign, Some(9), false)],
914 Some(9),
915 true,
916 ),
917 (
918 vec![file(local, Some(3), true), file(local, Some(6), false)],
919 Some(6),
920 false,
921 ),
922 (
923 vec![file(foreign, Some(9), false), file(foreign, Some(11), true)],
924 Some(11),
925 true,
926 ),
927 (
928 vec![file(local, Some(3), true), file(local, None, true)],
929 None,
930 true,
931 ),
932 (
933 vec![file(local, None, false), file(local, Some(6), true)],
934 None,
935 false,
936 ),
937 (
938 vec![file(local, Some(3), true), file(foreign, None, true)],
939 None,
940 false,
941 ),
942 (vec![], None, true),
943 (
944 vec![file(local, Some(u64::MAX), false)],
945 Some(u64::MAX),
946 false,
947 ),
948 ];
949 for committed_sequence in [100, 101] {
950 region
951 .manifest_ctx
952 .manifest_manager
953 .write()
954 .await
955 .update(
956 RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
957 committed_sequence: Some(committed_sequence),
958 files_to_add: vec![],
959 files_to_remove: vec![],
960 timestamp_ms: None,
961 compaction_time_window: None,
962 flushed_entry_id: None,
963 flushed_sequence: None,
964 })),
965 false,
966 )
967 .await
968 .unwrap();
969 for preserve in [false, true] {
970 region.region_options.preserve_row_sequence = preserve;
971 for (inputs, known_bound, inputs_trusted) in &cases {
972 let metadata = output_sequence_metadata(®ion, inputs);
973 let trusted = preserve && *inputs_trusted;
974 assert_eq!(trusted, metadata.exact_sequence_trusted);
975 assert_eq!(
976 known_bound.and_then(NonZeroU64::new),
977 metadata.file_sequence
978 );
979 }
980 }
981 }
982 }
983
984 #[rstest::rstest]
988 #[case(true)]
989 #[case(false)]
990 #[tokio::test]
991 async fn test_independent_outputs_preserve_delete_in_every_visibility_state(
992 #[case] flat_format: bool,
993 #[values(false, true)] unknown_put_bound: bool,
994 ) {
995 use datatypes::arrow::array::AsArray;
996 use datatypes::arrow::datatypes::TimestampMillisecondType;
997 use store_api::region_engine::RegionEngine;
998 use store_api::region_request::RegionRequest;
999
1000 use crate::compaction::CompactionOutput;
1001 use crate::compaction::compactor::{CompactionRegion, Compactor, DefaultCompactor};
1002 use crate::compaction::picker::PickerOutput;
1003 use crate::compaction::reader::CompactionSstReaderBuilder;
1004 use crate::engine::compaction_test::{delete_and_flush, put_and_flush};
1005 use crate::region::options::MergeMode;
1006 use crate::sst::file::FileHandle;
1007 use crate::sst::file_purger::NoopFilePurger;
1008 use crate::test_util::{CreateRequestBuilder, TestEnv, rows_schema};
1009
1010 let mut env = TestEnv::new().await;
1011 let region_id = RegionId::new(1, 1);
1012 let config = MitoConfig {
1013 default_flat_format: flat_format,
1014 min_compaction_interval: Duration::from_secs(3600),
1015 ..Default::default()
1016 };
1017 let engine = env.create_engine(config).await;
1018 let request = CreateRequestBuilder::new().build();
1019 let columns = rows_schema(&request);
1020 engine
1021 .handle_request(region_id, RegionRequest::Create(request))
1022 .await
1023 .unwrap();
1024 put_and_flush(&engine, region_id, &columns, 0..1).await;
1025 delete_and_flush(&engine, region_id, &columns, 0..1).await;
1026 put_and_flush(&engine, region_id, &columns, 10..11).await;
1027 let region = engine.get_region(region_id).unwrap();
1028 let version = region.version();
1029 let mut files: Vec<_> = version
1030 .ssts
1031 .levels()
1032 .iter()
1033 .flat_map(|level| level.files())
1034 .cloned()
1035 .collect();
1036 assert_eq!(3, files.len());
1037 files.sort_unstable_by_key(|file| file.meta_ref().sequence);
1038 if unknown_put_bound {
1039 let mut meta = files[0].meta_ref().clone();
1042 meta.sequence = None;
1043 files[0] = FileHandle::new(meta, Arc::new(NoopFilePurger));
1044 }
1045 let compaction_region = CompactionRegion {
1046 region_id,
1047 region_options: version.options.clone(),
1048 engine_config: Arc::new(MitoConfig {
1049 default_flat_format: flat_format,
1050 ..Default::default()
1051 }),
1052 region_metadata: version.metadata.clone(),
1053 cache_manager: engine.cache_manager(),
1054 access_layer: region.access_layer.clone(),
1055 manifest_ctx: region.manifest_ctx.clone(),
1056 current_version: version.into(),
1057 file_purger: None,
1058 ttl: None,
1059 max_parallelism: 2,
1060 plugins: common_base::Plugins::new(),
1061 };
1062 let puts = vec![files[0].clone(), files[2].clone()];
1063 let deletes = vec![files[1].clone()];
1064 let picker_output = PickerOutput {
1065 outputs: [puts.clone(), deletes.clone()]
1066 .into_iter()
1067 .map(|inputs| CompactionOutput {
1068 output_level: 1,
1069 inputs,
1070 filter_deleted: false,
1071 output_time_range: None,
1072 })
1073 .collect(),
1074 expired_ssts: vec![],
1075 time_window_size: 3600,
1076 max_file_size: None,
1077 };
1078 let merged = DefaultCompactor::with_merger(DefaultSstMerger)
1079 .merge_ssts(&compaction_region, picker_output)
1080 .await
1081 .unwrap();
1082 assert_eq!(2, merged.files_to_add.len());
1083 let outputs: Vec<_> = merged
1084 .files_to_add
1085 .into_iter()
1086 .map(|meta| FileHandle::new(meta, Arc::new(NoopFilePurger)))
1087 .collect();
1088 let put_output = outputs
1089 .iter()
1090 .find(|file| file.meta_ref().num_rows == 2)
1091 .unwrap();
1092 let delete_output = outputs
1093 .iter()
1094 .find(|file| file.meta_ref().num_rows == 1)
1095 .unwrap();
1096 assert_eq!(
1097 if unknown_put_bound {
1098 None
1099 } else {
1100 NonZeroU64::new(3)
1101 },
1102 put_output.meta_ref().sequence,
1103 );
1104 assert_eq!(NonZeroU64::new(2), delete_output.meta_ref().sequence);
1105 assert!(
1106 outputs
1107 .iter()
1108 .all(|file| !file.meta_ref().preserve_row_sequence)
1109 );
1110
1111 for (put_done, delete_done) in [(false, false), (true, false), (false, true), (true, true)]
1114 {
1115 let mut visible = if put_done {
1116 vec![put_output.clone()]
1117 } else {
1118 puts.clone()
1119 };
1120 visible.extend(if delete_done {
1121 vec![delete_output.clone()]
1122 } else {
1123 deletes.clone()
1124 });
1125 let mut reader = CompactionSstReaderBuilder {
1126 metadata: compaction_region.region_metadata.clone(),
1127 sst_layer: region.access_layer.clone(),
1128 cache: engine.cache_manager(),
1129 inputs: &visible,
1130 append_mode: false,
1131 filter_deleted: true,
1132 time_range: None,
1133 merge_mode: MergeMode::LastRow,
1134 }
1135 .build_flat_sst_reader()
1136 .await
1137 .unwrap();
1138 let mut timestamps = Vec::new();
1139 while let Some(batch) = reader.next_batch().await.unwrap() {
1140 timestamps.extend(
1141 batch
1142 .column_by_name("ts")
1143 .unwrap()
1144 .as_primitive::<TimestampMillisecondType>()
1145 .values()
1146 .iter()
1147 .copied(),
1148 );
1149 }
1150 assert_eq!(
1151 vec![10_000],
1152 timestamps,
1153 "put_done={put_done}, delete_done={delete_done}"
1154 );
1155 }
1156 }
1157
1158 async fn new_test_compaction_region() -> CompactionRegion {
1161 let env = SchedulerEnv::new().await;
1162 let metadata = metadata_for_test();
1163 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
1164 CompactionRegion {
1165 region_id: RegionId::new(1, 1),
1166 region_options: RegionOptions::default(),
1167 engine_config: Arc::new(MitoConfig::default()),
1168 region_metadata: metadata.clone(),
1169 cache_manager: Arc::new(CacheManager::default()),
1170 access_layer: env.access_layer.clone(),
1171 manifest_ctx,
1172 current_version: CompactionVersion {
1173 metadata: metadata.clone(),
1174 options: RegionOptions::default(),
1175 ssts: Arc::new(SstVersion::new(metadata)),
1176 memtable_min_sequence: None,
1177 compaction_time_window: None,
1178 },
1179 file_purger: None,
1180 ttl: None,
1181 max_parallelism: 1,
1182 plugins: Plugins::new(),
1183 }
1184 }
1185
1186 #[derive(Clone)]
1190 struct MockMerger {
1191 results: Arc<Mutex<Vec<Result<Vec<FileMeta>>>>>,
1192 call_idx: Arc<AtomicUsize>,
1193 }
1194
1195 impl MockMerger {
1196 fn new(results: Vec<Result<Vec<FileMeta>>>) -> Self {
1197 Self {
1198 results: Arc::new(Mutex::new(results)),
1199 call_idx: Arc::new(AtomicUsize::new(0)),
1200 }
1201 }
1202 }
1203
1204 #[async_trait::async_trait]
1205 impl SstMerger for MockMerger {
1206 async fn merge_single_output(
1207 &self,
1208 _compaction_region: CompactionRegion,
1209 _output: CompactionOutput,
1210 _write_opts: WriteOptions,
1211 ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
1212 let idx = self.call_idx.fetch_add(1, Ordering::SeqCst);
1213 match self.results.lock().unwrap().get(idx) {
1214 Some(Ok(files)) => Ok((files.clone(), Vec::new())),
1215 Some(Err(_)) => error::InvalidMetaSnafu {
1216 reason: format!("simulated failure at index {idx}"),
1217 }
1218 .fail(),
1219 None => panic!("MockMerger: no result configured for call index {idx}"),
1220 }
1221 }
1222 }
1223
1224 #[tokio::test]
1225 async fn test_partial_merge_failure_collects_only_successful_outputs() {
1226 common_telemetry::init_default_ut_logging();
1227
1228 let compaction_region = new_test_compaction_region().await;
1229
1230 let input_meta_0 = dummy_file_meta();
1232 let input_meta_1 = dummy_file_meta();
1233 let input_meta_2 = dummy_file_meta();
1234
1235 let output_meta_0 = vec![dummy_file_meta()];
1236 let output_meta_2 = vec![dummy_file_meta(), dummy_file_meta()];
1237
1238 let merger = MockMerger::new(vec![
1239 Ok(output_meta_0.clone()),
1240 Err(error::InvalidMetaSnafu {
1241 reason: "boom".to_string(),
1242 }
1243 .build()),
1244 Ok(output_meta_2.clone()),
1245 ]);
1246 let compactor = DefaultCompactor::with_merger(merger);
1247
1248 let picker_output = PickerOutput {
1249 outputs: vec![
1250 CompactionOutput {
1251 output_level: 1,
1252 inputs: vec![new_file_handle(input_meta_0.clone())],
1253 filter_deleted: false,
1254 output_time_range: None,
1255 },
1256 CompactionOutput {
1257 output_level: 1,
1258 inputs: vec![new_file_handle(input_meta_1.clone())],
1259 filter_deleted: false,
1260 output_time_range: None,
1261 },
1262 CompactionOutput {
1263 output_level: 1,
1264 inputs: vec![new_file_handle(input_meta_2.clone())],
1265 filter_deleted: false,
1266 output_time_range: None,
1267 },
1268 ],
1269 expired_ssts: vec![],
1270 time_window_size: 3600,
1271 max_file_size: None,
1272 };
1273
1274 let merge_output = compactor
1275 .merge_ssts(&compaction_region, picker_output)
1276 .await
1277 .unwrap();
1278
1279 assert_eq!(merge_output.files_to_add.len(), 3);
1281 assert_eq!(merge_output.files_to_remove.len(), 2);
1283
1284 let removed_ids: Vec<_> = merge_output
1285 .files_to_remove
1286 .iter()
1287 .map(|f| f.file_id)
1288 .collect();
1289 assert!(removed_ids.contains(&input_meta_0.file_id));
1290 assert!(removed_ids.contains(&input_meta_2.file_id));
1291 assert!(!removed_ids.contains(&input_meta_1.file_id));
1293 }
1294
1295 #[tokio::test]
1296 async fn test_all_outputs_succeed() {
1297 common_telemetry::init_default_ut_logging();
1298
1299 let compaction_region = new_test_compaction_region().await;
1300 let input_meta = dummy_file_meta();
1301 let output_meta = vec![dummy_file_meta()];
1302
1303 let merger = MockMerger::new(vec![Ok(output_meta.clone())]);
1304 let compactor = DefaultCompactor::with_merger(merger);
1305
1306 let picker_output = PickerOutput {
1307 outputs: vec![CompactionOutput {
1308 output_level: 1,
1309 inputs: vec![new_file_handle(input_meta.clone())],
1310 filter_deleted: false,
1311 output_time_range: None,
1312 }],
1313 expired_ssts: vec![],
1314 time_window_size: 3600,
1315 max_file_size: None,
1316 };
1317
1318 let merge_output = compactor
1319 .merge_ssts(&compaction_region, picker_output)
1320 .await
1321 .unwrap();
1322
1323 assert_eq!(merge_output.files_to_add.len(), 1);
1324 assert_eq!(merge_output.files_to_add[0].file_id, output_meta[0].file_id);
1325 assert_eq!(merge_output.files_to_remove.len(), 1);
1326 assert_eq!(merge_output.files_to_remove[0].file_id, input_meta.file_id);
1327 }
1328
1329 #[tokio::test]
1330 async fn test_expired_ssts_always_removed() {
1331 common_telemetry::init_default_ut_logging();
1332
1333 let compaction_region = new_test_compaction_region().await;
1334 let input_meta = dummy_file_meta();
1335 let expired_meta = dummy_file_meta();
1336
1337 let merger = MockMerger::new(vec![Err(error::InvalidMetaSnafu {
1339 reason: "fail".to_string(),
1340 }
1341 .build())]);
1342 let compactor = DefaultCompactor::with_merger(merger);
1343
1344 let picker_output = PickerOutput {
1345 outputs: vec![CompactionOutput {
1346 output_level: 1,
1347 inputs: vec![new_file_handle(input_meta.clone())],
1348 filter_deleted: false,
1349 output_time_range: None,
1350 }],
1351 expired_ssts: vec![new_file_handle(expired_meta.clone())],
1352 time_window_size: 3600,
1353 max_file_size: None,
1354 };
1355
1356 let merge_output = compactor
1357 .merge_ssts(&compaction_region, picker_output)
1358 .await
1359 .unwrap();
1360
1361 assert!(merge_output.files_to_add.is_empty());
1363 assert_eq!(merge_output.files_to_remove.len(), 1);
1365 assert_eq!(
1366 merge_output.files_to_remove[0].file_id,
1367 expired_meta.file_id
1368 );
1369 }
1370
1371 #[derive(Clone)]
1372 struct BlockingMerger {
1373 call_idx: Arc<AtomicUsize>,
1374 }
1375
1376 #[async_trait::async_trait]
1377 impl SstMerger for BlockingMerger {
1378 async fn merge_single_output(
1379 &self,
1380 _compaction_region: CompactionRegion,
1381 _output: CompactionOutput,
1382 _write_opts: WriteOptions,
1383 ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
1384 self.call_idx.fetch_add(1, Ordering::SeqCst);
1385 std::future::pending().await
1386 }
1387 }
1388
1389 #[tokio::test(flavor = "multi_thread")]
1390 async fn test_merge_ssts_cancels_spawned_tasks() {
1391 common_telemetry::init_default_ut_logging();
1392
1393 let mut compaction_region = new_test_compaction_region().await;
1394 compaction_region.max_parallelism = 2;
1395
1396 let cancel_handle = Arc::new(CancellationHandle::default());
1397 let call_idx = Arc::new(AtomicUsize::new(0));
1398 let compactor = DefaultCompactor {
1399 merger: BlockingMerger {
1400 call_idx: call_idx.clone(),
1401 },
1402 cancel_handle: cancel_handle.clone(),
1403 uncommitted: None,
1404 };
1405
1406 let picker_output = PickerOutput {
1407 outputs: vec![
1408 CompactionOutput {
1409 output_level: 1,
1410 inputs: vec![new_file_handle(dummy_file_meta())],
1411 filter_deleted: false,
1412 output_time_range: None,
1413 },
1414 CompactionOutput {
1415 output_level: 1,
1416 inputs: vec![new_file_handle(dummy_file_meta())],
1417 filter_deleted: false,
1418 output_time_range: None,
1419 },
1420 CompactionOutput {
1421 output_level: 1,
1422 inputs: vec![new_file_handle(dummy_file_meta())],
1423 filter_deleted: false,
1424 output_time_range: None,
1425 },
1426 ],
1427 expired_ssts: vec![],
1428 time_window_size: 3600,
1429 max_file_size: None,
1430 };
1431
1432 let task = tokio::spawn(async move {
1433 compactor
1434 .merge_ssts(&compaction_region, picker_output)
1435 .await
1436 });
1437
1438 sleep(Duration::from_millis(100)).await;
1439 cancel_handle.cancel();
1440
1441 let merge_output = task
1442 .await
1443 .expect("merge_ssts should stop after cancellation")
1444 .unwrap();
1445
1446 let started = call_idx.load(Ordering::SeqCst);
1447
1448 assert!(merge_output.files_to_add.is_empty());
1449 assert!(merge_output.files_to_remove.is_empty());
1450 assert_eq!(started, 2);
1451 }
1452}