From f5212d3631302390a458bfc23768a0c31c818bfa Mon Sep 17 00:00:00 2001 From: "Lei, HUANG" <6406592+v0y4g3r@users.noreply.github.com> Date: Tue, 1 Sep 2026 08:30:05 +0000 Subject: [PATCH] feat(mito2): add write cache upload hook (#8992) Add `WriteCacheUploadStoreWrapper` in `src/mito2/src/cache/write_cache.rs` and wire it through `src/mito2/src/worker.rs`. Signed-off-by: Lei, HUANG --- src/cmd/src/datanode/objbench.rs | 2 +- src/mito2/src/cache.rs | 1 + src/mito2/src/cache/write_cache.rs | 82 +++++++++++++++++++++++++++++- src/mito2/src/worker.rs | 9 +++- 4 files changed, 89 insertions(+), 5 deletions(-) diff --git a/src/cmd/src/datanode/objbench.rs b/src/cmd/src/datanode/objbench.rs index a01cad8468..ef2ff2af60 100644 --- a/src/cmd/src/datanode/objbench.rs +++ b/src/cmd/src/datanode/objbench.rs @@ -482,7 +482,7 @@ async fn build_cache_manager( puffin_manager: PuffinManagerFactory, intermediate_manager: IntermediateManager, ) -> error::Result { - let write_cache = write_cache_from_config(config, puffin_manager, intermediate_manager) + let write_cache = write_cache_from_config(config, puffin_manager, intermediate_manager, None) .await .map_err(|e| { error::IllegalConfigSnafu { diff --git a/src/mito2/src/cache.rs b/src/mito2/src/cache.rs index ac1f5b6691..642e7b1ab8 100644 --- a/src/mito2/src/cache.rs +++ b/src/mito2/src/cache.rs @@ -50,6 +50,7 @@ use smallvec::SmallVec; use snafu::{OptionExt, ResultExt}; use store_api::metadata::{RegionMetadata, RegionMetadataRef}; use store_api::storage::{ColumnId, ConcreteDataType, FileId, RegionId, TimeSeriesRowSelector}; +pub use write_cache::{WriteCacheUploadStoreWrapper, WriteCacheUploadStoreWrapperRef}; use crate::cache::cache_size::parquet_meta_size; use crate::cache::file_cache::{FileType, IndexKey}; diff --git a/src/mito2/src/cache/write_cache.rs b/src/mito2/src/cache/write_cache.rs index a85e3eb402..f105ed005e 100644 --- a/src/mito2/src/cache/write_cache.rs +++ b/src/mito2/src/cache/write_cache.rs @@ -42,6 +42,14 @@ use crate::sst::parquet::writer::ParquetWriter; use crate::sst::parquet::{SstInfo, WriteOptions}; use crate::sst::{DEFAULT_WRITE_BUFFER_SIZE, DEFAULT_WRITE_CONCURRENCY}; +/// Wraps the remote object store used by write cache uploads. +pub trait WriteCacheUploadStoreWrapper: Send + Sync { + /// Wraps an object store before uploading a cached file. + fn wrap(&self, store: ObjectStore) -> ObjectStore; +} + +pub type WriteCacheUploadStoreWrapperRef = Arc; + /// A cache for uploading files to remote object stores. /// /// It keeps files in local disk and then sends files to object stores. @@ -56,6 +64,8 @@ pub struct WriteCache { task_sender: UnboundedSender, /// Optional cache for manifest files. manifest_cache: Option, + /// Optional wrapper for remote stores used by uploads. + upload_store_wrapper: Option, } pub type WriteCacheRef = Arc; @@ -91,6 +101,7 @@ impl WriteCache { intermediate_manager, task_sender, manifest_cache, + upload_store_wrapper: None, }) } @@ -140,6 +151,15 @@ impl WriteCache { self.manifest_cache.clone() } + /// Sets the wrapper for remote stores used by uploads. + pub(crate) fn with_upload_store_wrapper( + mut self, + upload_store_wrapper: Option, + ) -> Self { + self.upload_store_wrapper = upload_store_wrapper; + self + } + /// Build the puffin manager pub(crate) fn build_puffin_manager(&self) -> SstPuffinManager { let store = self.file_cache.local_store(); @@ -379,7 +399,11 @@ impl WriteCache { .await .context(error::OpenDalSnafu)?; - let mut writer = remote_store + let upload_store = self.upload_store_wrapper.as_ref().map_or_else( + || remote_store.clone(), + |wrapper| wrapper.wrap(remote_store.clone()), + ); + let mut writer = upload_store .writer_with(upload_path) .chunk(DEFAULT_WRITE_BUFFER_SIZE.as_bytes() as usize) .concurrent(DEFAULT_WRITE_CONCURRENCY) @@ -501,7 +525,8 @@ impl UploadTracker { mod tests { use bytes::Bytes; use common_test_util::temp_dir::create_temp_dir; - use object_store::ATOMIC_WRITE_DIR; + use object_store::services::Memory; + use object_store::{ATOMIC_WRITE_DIR, ObjectStore}; use parquet::file::metadata::PageIndexPolicy; use store_api::region_request::PathType; use store_api::storage::FileId; @@ -521,6 +546,59 @@ mod tests { sst_file_handle_with_file_id, sst_region_metadata, }; + struct RedirectUploadStoreWrapper { + target_store: ObjectStore, + } + + impl WriteCacheUploadStoreWrapper for RedirectUploadStoreWrapper { + fn wrap(&self, _store: ObjectStore) -> ObjectStore { + self.target_store.clone() + } + } + + #[tokio::test] + async fn test_upload_uses_wrapped_remote_store() { + let env = TestEnv::new().await; + let local_store = ObjectStore::new(Memory::default()).unwrap().finish(); + let original_store = ObjectStore::new(Memory::default()).unwrap().finish(); + let target_store = ObjectStore::new(Memory::default()).unwrap().finish(); + let wrapper = Arc::new(RedirectUploadStoreWrapper { + target_store: target_store.clone(), + }); + let write_cache = WriteCache::new( + local_store.clone(), + ReadableSize::mb(10), + None, + None, + false, + env.get_puffin_manager(), + env.get_intermediate_manager(), + None, + ) + .await + .unwrap() + .with_upload_store_wrapper(Some(wrapper.clone())); + + let region_id = RegionId::new(1024, 1); + let file_id = FileId::random(); + let key = IndexKey::new(region_id, file_id, FileType::Parquet); + let cache_path = write_cache.file_cache.cache_file_path(key); + let upload_path = "wrapped-upload.parquet"; + let data = Bytes::from_static(b"wrapped upload data"); + local_store.write(&cache_path, data.clone()).await.unwrap(); + + write_cache + .upload(key, upload_path, &original_store) + .await + .unwrap(); + + assert!(original_store.stat(upload_path).await.is_err()); + assert_eq!( + target_store.read(upload_path).await.unwrap().to_vec(), + data.to_vec() + ); + } + #[tokio::test] async fn test_write_and_upload_sst() { // TODO(QuenKar): maybe find a way to create some object server for testing, diff --git a/src/mito2/src/worker.rs b/src/mito2/src/worker.rs index 63706cab4e..0c23affaf5 100644 --- a/src/mito2/src/worker.rs +++ b/src/mito2/src/worker.rs @@ -58,7 +58,7 @@ use tokio::sync::mpsc::{Receiver, Sender}; use tokio::sync::{Mutex, Semaphore, mpsc, oneshot, watch}; use crate::cache::write_cache::{WriteCache, WriteCacheRef}; -use crate::cache::{CacheManager, CacheManagerRef}; +use crate::cache::{CacheManager, CacheManagerRef, WriteCacheUploadStoreWrapperRef}; use crate::compaction::CompactionScheduler; use crate::compaction::memory_manager::{CompactionMemoryManager, new_compaction_memory_manager}; use crate::config::MitoConfig; @@ -195,10 +195,12 @@ impl WorkerGroup { let flush_semaphore = Arc::new(Semaphore::new(config.max_background_flushes)); // We use another scheduler to avoid purge jobs blocking other jobs. let purge_scheduler = Arc::new(LocalScheduler::new(config.max_background_purges)); + let upload_store_wrapper = plugins.get::(); let write_cache = write_cache_from_config( &config, puffin_manager_factory.clone(), intermediate_manager.clone(), + upload_store_wrapper, ) .await?; let cache_manager = Arc::new( @@ -415,6 +417,7 @@ impl WorkerGroup { &config, puffin_manager_factory.clone(), intermediate_manager.clone(), + None, ) .await?; let cache_manager = Arc::new( @@ -501,6 +504,7 @@ pub async fn write_cache_from_config( config: &MitoConfig, puffin_manager_factory: PuffinManagerFactory, intermediate_manager: IntermediateManager, + upload_store_wrapper: Option, ) -> Result> { if !config.enable_write_cache { return Ok(None); @@ -522,7 +526,8 @@ pub async fn write_cache_from_config( intermediate_manager, config.manifest_cache_size, ) - .await?; + .await? + .with_upload_store_wrapper(upload_store_wrapper); Ok(Some(Arc::new(cache))) }