diff --git a/src/mito2/src/cache/write_cache.rs b/src/mito2/src/cache/write_cache.rs index d16bc186fe..b464e366cb 100644 --- a/src/mito2/src/cache/write_cache.rs +++ b/src/mito2/src/cache/write_cache.rs @@ -26,7 +26,7 @@ use store_api::storage::RegionId; use tokio::sync::mpsc::{UnboundedSender, unbounded_channel}; use crate::access_layer::{ - FilePathProvider, Metrics, RegionFilePathFactory, SstInfoArray, SstWriteRequest, + FilePathProvider, Metrics, OperationType, RegionFilePathFactory, SstInfoArray, SstWriteRequest, TempFileCleaner, WriteCachePathProvider, WriteType, new_fs_cache_store, }; use crate::cache::file_cache::{FileCache, FileCacheRef, FileType, IndexKey, IndexValue}; @@ -45,7 +45,12 @@ 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; + /// + /// `op_type` identifies the origin of the upload so implementations can + /// apply different policies per operation, e.g. throttling compaction + /// uploads but not flush uploads. Index rebuild uploads are reported as + /// [`OperationType::Compact`]. + fn wrap(&self, store: ObjectStore, op_type: OperationType) -> ObjectStore; } pub type WriteCacheUploadStoreWrapperRef = Arc; @@ -205,7 +210,13 @@ impl WriteCache { .build_sst_file_path(region_file_id); if let Err(e) = self - .upload(parquet_key, &remote_path, &upload_request.remote_store) + .upload( + parquet_key, + &remote_path, + &upload_request.remote_store, + // `put_and_upload_sst` is only used by the flush path. + OperationType::Flush, + ) .await { // Clean up cache on failure @@ -292,6 +303,7 @@ impl WriteCache { let mut upload_tracker = UploadTracker::new(region_id); let mut err = None; + let op_type = write_request.op_type; let remote_store = &upload_request.remote_store; for sst in &sst_info { let parquet_key = IndexKey::new(region_id, sst.file_id, FileType::Parquet); @@ -299,7 +311,10 @@ impl WriteCache { .dest_path_provider .build_sst_file_path(RegionFileId::new(region_id, sst.file_id)); let start = Instant::now(); - if let Err(e) = self.upload(parquet_key, &parquet_path, remote_store).await { + if let Err(e) = self + .upload(parquet_key, &parquet_path, remote_store, op_type) + .await + { err = Some(e); break; } @@ -312,7 +327,10 @@ impl WriteCache { .dest_path_provider .build_index_file_path(RegionFileId::new(region_id, sst.file_id)); let start = Instant::now(); - if let Err(e) = self.upload(puffin_key, &puffin_path, remote_store).await { + if let Err(e) = self + .upload(puffin_key, &puffin_path, remote_store, op_type) + .await + { err = Some(e); break; } @@ -381,6 +399,7 @@ impl WriteCache { index_key: IndexKey, upload_path: &str, remote_store: &ObjectStore, + op_type: OperationType, ) -> Result<()> { let region_id = index_key.region_id; let file_id = index_key.file_id; @@ -406,7 +425,7 @@ impl WriteCache { let upload_store = self.upload_store_wrapper.as_ref().map_or_else( || remote_store.clone(), - |wrapper| wrapper.wrap(remote_store.clone()), + |wrapper| wrapper.wrap(remote_store.clone(), op_type), ); let mut writer = upload_store .writer_with(upload_path) @@ -528,6 +547,8 @@ impl UploadTracker { #[cfg(test)] mod tests { + use std::sync::Mutex; + use bytes::Bytes; use common_test_util::temp_dir::create_temp_dir; use object_store::services::Memory; @@ -553,10 +574,12 @@ mod tests { struct RedirectUploadStoreWrapper { target_store: ObjectStore, + last_op_type: Mutex>, } impl WriteCacheUploadStoreWrapper for RedirectUploadStoreWrapper { - fn wrap(&self, _store: ObjectStore) -> ObjectStore { + fn wrap(&self, _store: ObjectStore, op_type: OperationType) -> ObjectStore { + *self.last_op_type.lock().unwrap() = Some(op_type); self.target_store.clone() } } @@ -569,6 +592,7 @@ mod tests { let target_store = ObjectStore::new(Memory::default()).unwrap().finish(); let wrapper = Arc::new(RedirectUploadStoreWrapper { target_store: target_store.clone(), + last_op_type: Mutex::new(None), }); let write_cache = WriteCache::new( local_store.clone(), @@ -593,10 +617,15 @@ mod tests { local_store.write(&cache_path, data.clone()).await.unwrap(); write_cache - .upload(key, upload_path, &original_store) + .upload(key, upload_path, &original_store, OperationType::Compact) .await .unwrap(); + assert_eq!( + *wrapper.last_op_type.lock().unwrap(), + Some(OperationType::Compact) + ); + assert!(original_store.stat(upload_path).await.is_err()); assert_eq!( target_store.read(upload_path).await.unwrap().to_vec(), diff --git a/src/mito2/src/sst/index.rs b/src/mito2/src/sst/index.rs index db7089c201..6b26cf5255 100644 --- a/src/mito2/src/sst/index.rs +++ b/src/mito2/src/sst/index.rs @@ -1015,7 +1015,14 @@ impl IndexBuildTask { ) .build_index_file_path_with_version(index_id); if let Err(e) = write_cache - .upload(puffin_key, &puffin_path, remote_store) + // Index rebuild is background maintenance work, so its uploads + // are reported as compaction uploads. + .upload( + puffin_key, + &puffin_path, + remote_store, + OperationType::Compact, + ) .await { err = Some(e);