1use std::num::NonZero;
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;
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::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::{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) compaction_time_window: Option<Duration>,
78}
79
80impl From<VersionRef> for CompactionVersion {
81 fn from(value: VersionRef) -> Self {
82 Self {
83 metadata: value.metadata.clone(),
84 options: value.options.clone(),
85 ssts: value.ssts.clone(),
86 compaction_time_window: value.compaction_time_window,
87 }
88 }
89}
90
91#[derive(Clone)]
94pub struct CompactionRegion {
95 pub region_id: RegionId,
96 pub region_options: RegionOptions,
97
98 pub(crate) engine_config: Arc<MitoConfig>,
99 pub(crate) region_metadata: RegionMetadataRef,
100 pub(crate) cache_manager: CacheManagerRef,
101 pub access_layer: AccessLayerRef,
103 pub(crate) manifest_ctx: Arc<ManifestContext>,
104 pub(crate) current_version: CompactionVersion,
105 pub(crate) file_purger: Option<Arc<LocalFilePurger>>,
106 pub(crate) ttl: Option<TimeToLive>,
107
108 pub max_parallelism: usize,
113
114 pub(crate) plugins: Plugins,
115}
116
117#[derive(Clone)]
119pub struct OpenCompactionRegionRequest {
120 pub region_id: RegionId,
121 pub table_dir: String,
122 pub path_type: PathType,
123 pub region_options: RegionOptions,
124 pub max_parallelism: usize,
125 pub plugins: Plugins,
127}
128
129pub async fn open_compaction_region(
132 req: &OpenCompactionRegionRequest,
133 mito_config: &MitoConfig,
134 object_store_manager: ObjectStoreManagerRef,
135 ttl_provider: Either<TimeToLive, SchemaMetadataManagerRef>,
136) -> Result<CompactionRegion> {
137 let object_store = {
138 let name = &req.region_options.storage;
139 if let Some(name) = name {
140 object_store_manager
141 .find(name)
142 .with_context(|| ObjectStoreNotFoundSnafu {
143 object_store: name.clone(),
144 })?
145 } else {
146 object_store_manager.default_object_store()
147 }
148 };
149
150 let access_layer = {
151 let puffin_manager_factory = PuffinManagerFactory::new(
152 &mito_config.index.aux_path,
153 mito_config.index.staging_size.as_bytes(),
154 Some(mito_config.index.write_buffer_size.as_bytes() as _),
155 mito_config.index.staging_ttl,
156 )
157 .await?;
158 let intermediate_manager =
159 IntermediateManager::init_fs(mito_config.index.aux_path.clone()).await?;
160
161 Arc::new(AccessLayer::new(
162 &req.table_dir,
163 req.path_type,
164 object_store.clone(),
165 puffin_manager_factory,
166 intermediate_manager,
167 ))
168 };
169
170 let manifest_manager = {
171 let region_dir = region_dir_from_table_dir(&req.table_dir, req.region_id, req.path_type);
172 let region_manifest_options =
173 RegionManifestOptions::new(mito_config, ®ion_dir, object_store);
174
175 RegionManifestManager::open(region_manifest_options, &Default::default())
176 .await?
177 .with_context(|| EmptyRegionDirSnafu {
178 region_id: req.region_id,
179 region_dir: region_dir_from_table_dir(&req.table_dir, req.region_id, req.path_type),
180 })?
181 };
182
183 let manifest = manifest_manager.manifest();
184 let region_metadata = manifest.metadata.clone();
185 let hook: Option<RegionHookRef> = req.plugins.get();
186 let manifest_ctx = Arc::new(ManifestContext::new(
187 manifest_manager,
188 RegionRoleState::Leader(RegionLeaderState::Writable),
189 hook,
190 ));
191
192 let file_purger = {
193 let purge_scheduler = Arc::new(LocalScheduler::new(mito_config.max_background_purges));
194 Arc::new(LocalFilePurger::new(
195 purge_scheduler.clone(),
196 access_layer.clone(),
197 None,
198 ))
199 };
200
201 let current_version = {
202 let mut ssts = SstVersion::new();
203 ssts.add_files(file_purger.clone(), manifest.files.values().cloned());
204 CompactionVersion {
205 metadata: region_metadata.clone(),
206 options: req.region_options.clone(),
207 ssts: Arc::new(ssts),
208 compaction_time_window: manifest.compaction_time_window,
209 }
210 };
211
212 let ttl = match ttl_provider {
213 Either::Left(ttl) => ttl,
215 Either::Right(schema_metadata_manager) => {
217 let (_, ttl) =
218 find_dynamic_options(req.region_id, &req.region_options, &schema_metadata_manager)
219 .await
220 .unwrap_or_else(|e| {
221 warn!(e; "Failed to get ttl for region: {}", region_metadata.region_id);
222 (
223 crate::region::options::CompactionOptions::default(),
224 TimeToLive::default(),
225 )
226 });
227 ttl
228 }
229 };
230
231 Ok(CompactionRegion {
232 region_id: req.region_id,
233 region_options: req.region_options.clone(),
234 engine_config: Arc::new(mito_config.clone()),
235 region_metadata: region_metadata.clone(),
236 cache_manager: Arc::new(CacheManager::default()),
237 access_layer,
238 manifest_ctx,
239 current_version,
240 file_purger: Some(file_purger),
241 ttl: Some(ttl),
242 max_parallelism: req.max_parallelism,
243 plugins: req.plugins.clone(),
244 })
245}
246
247impl CompactionRegion {
248 pub fn file_purger(&self) -> Option<Arc<LocalFilePurger>> {
250 self.file_purger.clone()
251 }
252
253 pub async fn stop_purger_scheduler(&self) -> Result<()> {
255 if let Some(file_purger) = &self.file_purger {
256 file_purger.stop_scheduler().await
257 } else {
258 Ok(())
259 }
260 }
261
262 pub async fn invoke_sst_hook(&self, merge_output: &MergeOutput) {
267 let Some(hook) = self.plugins.get::<RegionHookRef>() else {
268 return;
269 };
270
271 let synthesized: Vec<SstInfo>;
276 let infos: &[SstInfo] = if merge_output.sst_infos.is_empty() {
277 synthesized = merge_output
278 .files_to_add
279 .iter()
280 .map(sst_info_from_file_meta)
281 .collect();
282 &synthesized
283 } else {
284 if merge_output.sst_infos.len() != merge_output.files_to_add.len() {
288 warn!(
289 "sst_infos length ({}) does not match files_to_add length ({}) for region {}",
290 merge_output.sst_infos.len(),
291 merge_output.files_to_add.len(),
292 self.region_id
293 );
294 }
295 &merge_output.sst_infos
296 };
297
298 let files: Vec<SstFileInfo<'_>> = merge_output
299 .files_to_add
300 .iter()
301 .zip(infos)
302 .map(|(meta, info)| SstFileInfo {
303 sst_info_ref: info,
304 file_meta: meta,
305 })
306 .collect();
307 hook.on_sst_files_written(self.region_id, &self.region_metadata, &files)
308 .await;
309 }
310}
311
312fn sst_info_from_file_meta(meta: &FileMeta) -> SstInfo {
335 SstInfo {
336 file_id: meta.file_id,
337 time_range: meta.time_range,
338 file_size: meta.file_size,
339 max_row_group_uncompressed_size: meta.max_row_group_uncompressed_size,
340 num_rows: meta.num_rows as usize,
341 num_row_groups: meta.num_row_groups,
342 num_series: meta.num_series,
343 ..Default::default()
344 }
345}
346
347#[derive(Default, Debug, Serialize, Deserialize)]
349pub struct MergeOutput {
350 pub files_to_add: Vec<FileMeta>,
351 pub files_to_remove: Vec<FileMeta>,
352 pub compaction_time_window: Option<i64>,
353 #[serde(skip)]
354 pub sst_infos: Vec<SstInfo>,
355}
356
357impl MergeOutput {
358 pub fn is_empty(&self) -> bool {
359 self.files_to_add.is_empty() && self.files_to_remove.is_empty()
360 }
361
362 pub fn input_file_size(&self) -> u64 {
363 self.files_to_remove.iter().map(|f| f.file_size).sum()
364 }
365
366 pub fn output_file_size(&self) -> u64 {
367 self.files_to_add.iter().map(|f| f.file_size).sum()
368 }
369}
370
371#[async_trait::async_trait]
373pub trait Compactor: Send + Sync + 'static {
374 async fn merge_ssts(
376 &self,
377 compaction_region: &CompactionRegion,
378 picker_output: PickerOutput,
379 ) -> Result<MergeOutput>;
380
381 async fn update_manifest(
383 &self,
384 compaction_region: &CompactionRegion,
385 merge_output: MergeOutput,
386 ) -> Result<(RegionEdit, ManifestVersion)>;
387}
388
389#[async_trait::async_trait]
394pub trait SstMerger: Send + Sync + 'static {
395 async fn merge_single_output(
396 &self,
397 compaction_region: CompactionRegion,
398 output: CompactionOutput,
399 write_opts: WriteOptions,
400 ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)>;
401}
402
403#[derive(Clone)]
405pub struct DefaultSstMerger;
406
407#[async_trait::async_trait]
408impl SstMerger for DefaultSstMerger {
409 async fn merge_single_output(
410 &self,
411 compaction_region: CompactionRegion,
412 output: CompactionOutput,
413 write_opts: WriteOptions,
414 ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
415 let region_id = compaction_region.region_id;
416 let storage = compaction_region.region_options.storage.clone();
417 let index_options = compaction_region
418 .current_version
419 .options
420 .index_options
421 .clone();
422 let append_mode = compaction_region.current_version.options.append_mode;
423 let merge_mode = compaction_region.current_version.options.merge_mode();
424 let flat_format = compaction_region
425 .region_options
426 .sst_format
427 .map(|format| format == FormatType::Flat)
428 .unwrap_or(compaction_region.engine_config.default_flat_format);
429
430 let index_config = compaction_region.engine_config.index.clone();
431 let inverted_index_config = compaction_region.engine_config.inverted_index.clone();
432 let fulltext_index_config = compaction_region.engine_config.fulltext_index.clone();
433 let bloom_filter_index_config = compaction_region.engine_config.bloom_filter_index.clone();
434 #[cfg(feature = "vector_index")]
435 let vector_index_config = compaction_region.engine_config.vector_index.clone();
436
437 let input_file_names = output
438 .inputs
439 .iter()
440 .map(|f| f.file_id().to_string())
441 .join(",");
442 let max_sequence = output
443 .inputs
444 .iter()
445 .map(|f| f.meta_ref().sequence)
446 .max()
447 .flatten();
448 let builder = CompactionSstReaderBuilder {
449 metadata: compaction_region.region_metadata.clone(),
450 sst_layer: compaction_region.access_layer.clone(),
451 cache: compaction_region.cache_manager.clone(),
452 inputs: &output.inputs,
453 append_mode,
454 filter_deleted: output.filter_deleted,
455 time_range: output.output_time_range,
456 merge_mode,
457 };
458 let source = builder.build_flat_sst_reader().await?;
459
460 let mut metrics = Metrics::new(WriteType::Compaction);
461 let region_metadata = compaction_region.region_metadata.clone();
462 let sst_infos = compaction_region
463 .access_layer
464 .write_sst(
465 SstWriteRequest {
466 op_type: OperationType::Compact,
467 metadata: region_metadata.clone(),
468 source,
469 cache_manager: compaction_region.cache_manager.clone(),
470 storage,
471 max_sequence: max_sequence.map(NonZero::get),
472 sst_write_format: if flat_format {
473 FormatType::Flat
474 } else {
475 FormatType::PrimaryKey
476 },
477 index_options,
478 index_config,
479 inverted_index_config,
480 fulltext_index_config,
481 bloom_filter_index_config,
482 #[cfg(feature = "vector_index")]
483 vector_index_config,
484 },
485 &write_opts,
486 &mut metrics,
487 )
488 .await?;
489 let partition_expr = match ®ion_metadata.partition_expr {
491 None => None,
492 Some(json_str) if json_str.is_empty() => None,
493 Some(json_str) => PartitionExpr::from_json_str(json_str).with_context(|_| {
494 InvalidPartitionExprSnafu {
495 expr: json_str.clone(),
496 }
497 })?,
498 };
499
500 let output_files = sst_infos
501 .iter()
502 .map(|sst_info| {
503 let pk_range = sst_info
504 .file_metadata
505 .as_ref()
506 .and_then(|meta| extract_primary_key_range(meta, ®ion_metadata));
507 let (primary_key_min, primary_key_max) = match pk_range {
508 Some((min, max)) => (Some(min), Some(max)),
509 None => (None, None),
510 };
511
512 FileMeta {
513 region_id,
514 file_id: sst_info.file_id,
515 time_range: sst_info.time_range,
516 level: output.output_level,
517 file_size: sst_info.file_size,
518 max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
519 available_indexes: sst_info.index_metadata.build_available_indexes(),
520 indexes: sst_info.index_metadata.build_indexes(),
521 index_file_size: sst_info.index_metadata.file_size,
522 index_version: 0,
523 num_rows: sst_info.num_rows as u64,
524 num_row_groups: sst_info.num_row_groups,
525 sequence: max_sequence,
526 partition_expr: partition_expr.clone(),
527 num_series: sst_info.num_series,
528 primary_key_min,
529 primary_key_max,
530 }
531 })
532 .collect::<Vec<_>>();
533 let output_file_names = output_files.iter().map(|f| f.file_id.to_string()).join(",");
534 info!(
535 "Region {} compaction inputs: [{}], outputs: [{}], flat_format: {}, metrics: {:?}",
536 region_id, input_file_names, output_file_names, flat_format, metrics
537 );
538 metrics.observe();
539 Ok((output_files, sst_infos.into_iter().collect()))
540 }
541}
542
543pub struct DefaultCompactor<M = DefaultSstMerger> {
548 merger: M,
549 cancel_handle: Arc<CancellationHandle>,
550 uncommitted: Option<UncommittedSsts>,
551}
552
553#[cfg(test)]
554impl<M: SstMerger> DefaultCompactor<M> {
555 pub fn with_merger(merger: M) -> Self {
556 Self {
557 merger,
558 cancel_handle: Arc::new(CancellationHandle::default()),
559 uncommitted: None,
560 }
561 }
562}
563
564impl DefaultCompactor {
565 pub fn new_with_cancel_handle(cancel_handle: Arc<CancellationHandle>) -> Self {
571 Self {
572 merger: DefaultSstMerger,
573 cancel_handle,
574 uncommitted: None,
575 }
576 }
577
578 pub(crate) fn with_cancel_handle(
579 cancel_handle: Arc<CancellationHandle>,
580 uncommitted: UncommittedSsts,
581 ) -> Self {
582 Self {
583 merger: DefaultSstMerger,
584 cancel_handle,
585 uncommitted: Some(uncommitted),
586 }
587 }
588}
589
590#[async_trait::async_trait]
591impl<M: SstMerger> Compactor for DefaultCompactor<M>
592where
593 M: Clone,
594{
595 async fn merge_ssts(
596 &self,
597 compaction_region: &CompactionRegion,
598 mut picker_output: PickerOutput,
599 ) -> Result<MergeOutput> {
600 let internal_parallelism = compaction_region.max_parallelism.max(1);
601 let compaction_time_window = picker_output.time_window_size;
602 let region_id = compaction_region.region_id;
603
604 let mut tasks: Vec<(Vec<FileMeta>, _)> = Vec::with_capacity(picker_output.outputs.len());
607
608 for output in picker_output.outputs.drain(..) {
609 let inputs_to_remove: Vec<_> =
610 output.inputs.iter().map(|f| f.meta_ref().clone()).collect();
611 let write_opts = WriteOptions {
612 write_buffer_size: compaction_region.engine_config.sst_write_buffer_size,
613 max_file_size: picker_output.max_file_size,
614 row_group_size: compaction_region.region_options.row_group_size(),
615 };
616 let merger = self.merger.clone();
617 let compaction_region = compaction_region.clone();
618 let uncommitted = self.uncommitted.clone();
619 let fut = async move {
620 let result = merger
621 .merge_single_output(compaction_region, output, write_opts)
622 .await;
623 if let (Some(uncommitted), Ok((_, infos))) = (&uncommitted, &result) {
624 uncommitted.track(infos);
625 }
626 result
627 };
628 tasks.push((inputs_to_remove, fut));
629 }
630
631 let hook: Option<RegionHookRef> = compaction_region.plugins.get();
632 let mut output_files = Vec::with_capacity(tasks.len());
633 let mut all_sst_infos: Vec<SstInfo> = Vec::new();
634 let mut compacted_inputs = Vec::with_capacity(
635 tasks.iter().map(|(inputs, _)| inputs.len()).sum::<usize>()
636 + picker_output.expired_ssts.len(),
637 );
638
639 while !tasks.is_empty() {
640 let mut chunk: Vec<(Vec<FileMeta>, _)> = Vec::with_capacity(internal_parallelism);
641 for _ in 0..internal_parallelism {
642 if let Some(task) = tasks.pop() {
643 chunk.push(task);
644 }
645 }
646 let mut spawned: Vec<_> = chunk
647 .into_iter()
648 .map(|(inputs, fut)| {
649 let handle = common_runtime::spawn_compact(fut);
650 (inputs, handle)
651 })
652 .collect();
653
654 while let Some((inputs, mut handle)) = spawned.pop() {
655 let abort_handle = handle.abort_handle();
656 match CancellableFuture::new(&mut handle, self.cancel_handle.clone()).await {
657 Ok(Ok(Ok((files, infos)))) => {
658 output_files.extend(files);
659 if hook.is_some() {
660 all_sst_infos.extend(infos);
661 }
662 compacted_inputs.extend(inputs);
663 }
664 Ok(Ok(Err(e))) => {
665 warn!(
666 e; "Failed to merge compaction output for region: {}, inputs: [{}]",
667 region_id,
668 inputs.iter().map(|f| f.file_id.to_string()).join(",")
669 );
670 }
671 Ok(Err(e)) => {
672 warn!(
673 "Region {} compaction task join error for inputs: [{}], skipping: {}",
674 region_id,
675 inputs.iter().map(|f| f.file_id.to_string()).join(","),
676 e
677 );
678 if self.cancel_handle.is_cancelled() {
681 abort_handle.abort();
682 for (_, handle) in &spawned {
683 handle.abort();
684 }
685 for (_, handle) in spawned {
686 let _ = handle.await;
687 }
688 }
689 return Err(e).context(error::JoinSnafu);
690 }
691 Err(_) => {
692 debug!(
693 "Compaction merge cancelled for region: {}, aborting remaining {} spawned tasks",
694 region_id,
695 spawned.len(),
696 );
697 abort_handle.abort();
698 for (_, handle) in &spawned {
699 handle.abort();
700 }
701 let _ = handle.await;
702 for (_, handle) in spawned {
703 let _ = handle.await;
704 }
705 break;
706 }
707 }
708 }
709
710 if self.cancel_handle.is_cancelled() {
711 info!("Compaction merge cancelled for region: {}", region_id);
712 break;
713 }
714 }
715
716 compacted_inputs.extend(
718 picker_output
719 .expired_ssts
720 .iter()
721 .map(|f| f.meta_ref().clone()),
722 );
723
724 Ok(MergeOutput {
725 files_to_add: output_files,
726 files_to_remove: compacted_inputs,
727 compaction_time_window: Some(compaction_time_window),
728 sst_infos: all_sst_infos,
729 })
730 }
731
732 async fn update_manifest(
733 &self,
734 compaction_region: &CompactionRegion,
735 merge_output: MergeOutput,
736 ) -> Result<(RegionEdit, ManifestVersion)> {
737 let edit = RegionEdit {
739 files_to_add: merge_output.files_to_add,
740 files_to_remove: merge_output.files_to_remove,
741 timestamp_ms: Some(chrono::Utc::now().timestamp_millis()),
743 compaction_time_window: merge_output
744 .compaction_time_window
745 .map(|seconds| Duration::from_secs(seconds as u64)),
746 flushed_entry_id: None,
747 flushed_sequence: None,
748 committed_sequence: None,
749 };
750
751 let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit.clone()));
752 let manifest_version = compaction_region
754 .manifest_ctx
755 .update_manifest_for_compaction(action_list)
756 .await?;
757
758 Ok((edit, manifest_version))
759 }
760}
761
762#[cfg(test)]
763mod tests {
764 use std::sync::atomic::{AtomicUsize, Ordering};
765 use std::sync::{Arc, Mutex};
766 use std::time::Duration;
767
768 use store_api::storage::{FileId, RegionId};
769 use tokio::time::sleep;
770
771 use super::{DefaultCompactor, *};
772 use crate::cache::CacheManager;
773 use crate::compaction::picker::PickerOutput;
774 use crate::error::Result;
775 use crate::sst::file::FileHandle;
776 use crate::sst::file_purger::NoopFilePurger;
777 use crate::sst::version::SstVersion;
778 use crate::test_util::memtable_util::metadata_for_test;
779 use crate::test_util::scheduler_util::SchedulerEnv;
780
781 fn dummy_file_meta() -> FileMeta {
782 FileMeta {
783 region_id: RegionId::new(1, 1),
784 file_id: FileId::random(),
785 file_size: 100,
786 ..Default::default()
787 }
788 }
789
790 fn new_file_handle(meta: FileMeta) -> FileHandle {
791 FileHandle::new(meta, Arc::new(NoopFilePurger))
792 }
793
794 async fn new_test_compaction_region() -> CompactionRegion {
797 let env = SchedulerEnv::new().await;
798 let metadata = metadata_for_test();
799 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
800 CompactionRegion {
801 region_id: RegionId::new(1, 1),
802 region_options: RegionOptions::default(),
803 engine_config: Arc::new(MitoConfig::default()),
804 region_metadata: metadata.clone(),
805 cache_manager: Arc::new(CacheManager::default()),
806 access_layer: env.access_layer.clone(),
807 manifest_ctx,
808 current_version: CompactionVersion {
809 metadata,
810 options: RegionOptions::default(),
811 ssts: Arc::new(SstVersion::new()),
812 compaction_time_window: None,
813 },
814 file_purger: None,
815 ttl: None,
816 max_parallelism: 1,
817 plugins: Plugins::new(),
818 }
819 }
820
821 #[derive(Clone)]
825 struct MockMerger {
826 results: Arc<Mutex<Vec<Result<Vec<FileMeta>>>>>,
827 call_idx: Arc<AtomicUsize>,
828 }
829
830 impl MockMerger {
831 fn new(results: Vec<Result<Vec<FileMeta>>>) -> Self {
832 Self {
833 results: Arc::new(Mutex::new(results)),
834 call_idx: Arc::new(AtomicUsize::new(0)),
835 }
836 }
837 }
838
839 #[async_trait::async_trait]
840 impl SstMerger for MockMerger {
841 async fn merge_single_output(
842 &self,
843 _compaction_region: CompactionRegion,
844 _output: CompactionOutput,
845 _write_opts: WriteOptions,
846 ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
847 let idx = self.call_idx.fetch_add(1, Ordering::SeqCst);
848 match self.results.lock().unwrap().get(idx) {
849 Some(Ok(files)) => Ok((files.clone(), Vec::new())),
850 Some(Err(_)) => error::InvalidMetaSnafu {
851 reason: format!("simulated failure at index {idx}"),
852 }
853 .fail(),
854 None => panic!("MockMerger: no result configured for call index {idx}"),
855 }
856 }
857 }
858
859 #[tokio::test]
860 async fn test_partial_merge_failure_collects_only_successful_outputs() {
861 common_telemetry::init_default_ut_logging();
862
863 let compaction_region = new_test_compaction_region().await;
864
865 let input_meta_0 = dummy_file_meta();
867 let input_meta_1 = dummy_file_meta();
868 let input_meta_2 = dummy_file_meta();
869
870 let output_meta_0 = vec![dummy_file_meta()];
871 let output_meta_2 = vec![dummy_file_meta(), dummy_file_meta()];
872
873 let merger = MockMerger::new(vec![
874 Ok(output_meta_0.clone()),
875 Err(error::InvalidMetaSnafu {
876 reason: "boom".to_string(),
877 }
878 .build()),
879 Ok(output_meta_2.clone()),
880 ]);
881 let compactor = DefaultCompactor::with_merger(merger);
882
883 let picker_output = PickerOutput {
884 outputs: vec![
885 CompactionOutput {
886 output_level: 1,
887 inputs: vec![new_file_handle(input_meta_0.clone())],
888 filter_deleted: false,
889 output_time_range: None,
890 },
891 CompactionOutput {
892 output_level: 1,
893 inputs: vec![new_file_handle(input_meta_1.clone())],
894 filter_deleted: false,
895 output_time_range: None,
896 },
897 CompactionOutput {
898 output_level: 1,
899 inputs: vec![new_file_handle(input_meta_2.clone())],
900 filter_deleted: false,
901 output_time_range: None,
902 },
903 ],
904 expired_ssts: vec![],
905 time_window_size: 3600,
906 max_file_size: None,
907 };
908
909 let merge_output = compactor
910 .merge_ssts(&compaction_region, picker_output)
911 .await
912 .unwrap();
913
914 assert_eq!(merge_output.files_to_add.len(), 3);
916 assert_eq!(merge_output.files_to_remove.len(), 2);
918
919 let removed_ids: Vec<_> = merge_output
920 .files_to_remove
921 .iter()
922 .map(|f| f.file_id)
923 .collect();
924 assert!(removed_ids.contains(&input_meta_0.file_id));
925 assert!(removed_ids.contains(&input_meta_2.file_id));
926 assert!(!removed_ids.contains(&input_meta_1.file_id));
928 }
929
930 #[tokio::test]
931 async fn test_all_outputs_succeed() {
932 common_telemetry::init_default_ut_logging();
933
934 let compaction_region = new_test_compaction_region().await;
935 let input_meta = dummy_file_meta();
936 let output_meta = vec![dummy_file_meta()];
937
938 let merger = MockMerger::new(vec![Ok(output_meta.clone())]);
939 let compactor = DefaultCompactor::with_merger(merger);
940
941 let picker_output = PickerOutput {
942 outputs: vec![CompactionOutput {
943 output_level: 1,
944 inputs: vec![new_file_handle(input_meta.clone())],
945 filter_deleted: false,
946 output_time_range: None,
947 }],
948 expired_ssts: vec![],
949 time_window_size: 3600,
950 max_file_size: None,
951 };
952
953 let merge_output = compactor
954 .merge_ssts(&compaction_region, picker_output)
955 .await
956 .unwrap();
957
958 assert_eq!(merge_output.files_to_add.len(), 1);
959 assert_eq!(merge_output.files_to_add[0].file_id, output_meta[0].file_id);
960 assert_eq!(merge_output.files_to_remove.len(), 1);
961 assert_eq!(merge_output.files_to_remove[0].file_id, input_meta.file_id);
962 }
963
964 #[tokio::test]
965 async fn test_expired_ssts_always_removed() {
966 common_telemetry::init_default_ut_logging();
967
968 let compaction_region = new_test_compaction_region().await;
969 let input_meta = dummy_file_meta();
970 let expired_meta = dummy_file_meta();
971
972 let merger = MockMerger::new(vec![Err(error::InvalidMetaSnafu {
974 reason: "fail".to_string(),
975 }
976 .build())]);
977 let compactor = DefaultCompactor::with_merger(merger);
978
979 let picker_output = PickerOutput {
980 outputs: vec![CompactionOutput {
981 output_level: 1,
982 inputs: vec![new_file_handle(input_meta.clone())],
983 filter_deleted: false,
984 output_time_range: None,
985 }],
986 expired_ssts: vec![new_file_handle(expired_meta.clone())],
987 time_window_size: 3600,
988 max_file_size: None,
989 };
990
991 let merge_output = compactor
992 .merge_ssts(&compaction_region, picker_output)
993 .await
994 .unwrap();
995
996 assert!(merge_output.files_to_add.is_empty());
998 assert_eq!(merge_output.files_to_remove.len(), 1);
1000 assert_eq!(
1001 merge_output.files_to_remove[0].file_id,
1002 expired_meta.file_id
1003 );
1004 }
1005
1006 #[derive(Clone)]
1007 struct BlockingMerger {
1008 call_idx: Arc<AtomicUsize>,
1009 }
1010
1011 #[async_trait::async_trait]
1012 impl SstMerger for BlockingMerger {
1013 async fn merge_single_output(
1014 &self,
1015 _compaction_region: CompactionRegion,
1016 _output: CompactionOutput,
1017 _write_opts: WriteOptions,
1018 ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
1019 self.call_idx.fetch_add(1, Ordering::SeqCst);
1020 std::future::pending().await
1021 }
1022 }
1023
1024 #[tokio::test(flavor = "multi_thread")]
1025 async fn test_merge_ssts_cancels_spawned_tasks() {
1026 common_telemetry::init_default_ut_logging();
1027
1028 let mut compaction_region = new_test_compaction_region().await;
1029 compaction_region.max_parallelism = 2;
1030
1031 let cancel_handle = Arc::new(CancellationHandle::default());
1032 let call_idx = Arc::new(AtomicUsize::new(0));
1033 let compactor = DefaultCompactor {
1034 merger: BlockingMerger {
1035 call_idx: call_idx.clone(),
1036 },
1037 cancel_handle: cancel_handle.clone(),
1038 uncommitted: None,
1039 };
1040
1041 let picker_output = PickerOutput {
1042 outputs: vec![
1043 CompactionOutput {
1044 output_level: 1,
1045 inputs: vec![new_file_handle(dummy_file_meta())],
1046 filter_deleted: false,
1047 output_time_range: None,
1048 },
1049 CompactionOutput {
1050 output_level: 1,
1051 inputs: vec![new_file_handle(dummy_file_meta())],
1052 filter_deleted: false,
1053 output_time_range: None,
1054 },
1055 CompactionOutput {
1056 output_level: 1,
1057 inputs: vec![new_file_handle(dummy_file_meta())],
1058 filter_deleted: false,
1059 output_time_range: None,
1060 },
1061 ],
1062 expired_ssts: vec![],
1063 time_window_size: 3600,
1064 max_file_size: None,
1065 };
1066
1067 let task = tokio::spawn(async move {
1068 compactor
1069 .merge_ssts(&compaction_region, picker_output)
1070 .await
1071 });
1072
1073 sleep(Duration::from_millis(100)).await;
1074 cancel_handle.cancel();
1075
1076 let merge_output = task
1077 .await
1078 .expect("merge_ssts should stop after cancellation")
1079 .unwrap();
1080
1081 let started = call_idx.load(Ordering::SeqCst);
1082
1083 assert!(merge_output.files_to_add.is_empty());
1084 assert!(merge_output.files_to_remove.is_empty());
1085 assert_eq!(started, 2);
1086 }
1087}