From d34071000bf8e2defef9066a3138b7987b7517d2 Mon Sep 17 00:00:00 2001 From: Ning Sun Date: Fri, 5 Jun 2026 01:18:36 -0700 Subject: [PATCH] feat: load manifest version from compaction task --- src/mito2/src/compaction/compactor.rs | 9 +++++---- src/mito2/src/compaction/task.rs | 16 +++++----------- 2 files changed, 10 insertions(+), 15 deletions(-) diff --git a/src/mito2/src/compaction/compactor.rs b/src/mito2/src/compaction/compactor.rs index ef608ae0fc..d540b83b2c 100644 --- a/src/mito2/src/compaction/compactor.rs +++ b/src/mito2/src/compaction/compactor.rs @@ -27,6 +27,7 @@ use object_store::manager::ObjectStoreManagerRef; use partition::expr::PartitionExpr; use serde::{Deserialize, Serialize}; use snafu::{OptionExt, ResultExt}; +use store_api::ManifestVersion; use store_api::metadata::RegionMetadataRef; use store_api::region_request::PathType; use store_api::storage::RegionId; @@ -292,7 +293,7 @@ pub trait Compactor: Send + Sync + 'static { &self, compaction_region: &CompactionRegion, merge_output: MergeOutput, - ) -> Result; + ) -> Result<(RegionEdit, ManifestVersion)>; } /// Trait for merging a single compaction output into SST files. @@ -608,7 +609,7 @@ where &self, compaction_region: &CompactionRegion, merge_output: MergeOutput, - ) -> Result { + ) -> Result<(RegionEdit, ManifestVersion)> { // Write region edit to manifest. let edit = RegionEdit { files_to_add: merge_output.files_to_add, @@ -625,12 +626,12 @@ where let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit.clone())); // TODO: We might leak files if we fail to update manifest. We can add a cleanup task to remove them later. - compaction_region + let manifest_version = compaction_region .manifest_ctx .update_manifest_for_compaction(action_list) .await?; - Ok(edit) + Ok((edit, manifest_version)) } } diff --git a/src/mito2/src/compaction/task.rs b/src/mito2/src/compaction/task.rs index 0e3725b8e6..6cdb3bf1db 100644 --- a/src/mito2/src/compaction/task.rs +++ b/src/mito2/src/compaction/task.rs @@ -21,6 +21,7 @@ use common_memory_manager::OnExhaustedPolicy; use common_telemetry::{error, info, warn}; use itertools::Itertools; use snafu::ResultExt; +use store_api::ManifestVersion; use tokio::sync::mpsc; use crate::compaction::LocalCompactionState; @@ -250,7 +251,7 @@ impl CompactionTaskImpl { async fn update_manifest( &self, compaction_result: crate::compaction::compactor::MergeOutput, - ) -> error::Result { + ) -> error::Result<(RegionEdit, ManifestVersion)> { let _manifest_timer = COMPACTION_STAGE_ELAPSED .with_label_values(&["write_manifest"]) .start_timer(); @@ -305,16 +306,9 @@ impl CompactionTaskImpl { } } - async fn invoke_manifest_hook(&self, edit: &RegionEdit) { + async fn invoke_manifest_hook(&self, edit: &RegionEdit, manifest_version: ManifestVersion) { let hook: Option = self.compaction_region.plugins.get(); if let Some(hook) = hook { - let manifest_version = self - .compaction_region - .manifest_ctx - .manifest_manager - .read() - .await - .last_version(); hook.on_manifest_updated(self.compaction_region.region_id, edit, manifest_version) .await; } @@ -367,8 +361,8 @@ impl CompactionTask for CompactionTaskImpl { }) } else { match self.update_manifest(merge_output).await { - Ok(edit) => { - self.invoke_manifest_hook(&edit).await; + Ok((edit, manifest_version)) => { + self.invoke_manifest_hook(&edit, manifest_version).await; let senders = std::mem::take(&mut self.waiters); BackgroundNotify::CompactionFinished(CompactionFinished { region_id: self.compaction_region.region_id,