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 <ratuthomm@gmail.com>
This commit is contained in:
Lei, HUANG
2026-09-01 08:30:05 +00:00
committed by GitHub
parent 00d43b29ad
commit f5212d3631
4 changed files with 89 additions and 5 deletions
+1 -1
View File
@@ -482,7 +482,7 @@ async fn build_cache_manager(
puffin_manager: PuffinManagerFactory,
intermediate_manager: IntermediateManager,
) -> error::Result<CacheManagerRef> {
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 {
+1
View File
@@ -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};
+80 -2
View File
@@ -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<dyn WriteCacheUploadStoreWrapper>;
/// 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<RegionLoadCacheTask>,
/// Optional cache for manifest files.
manifest_cache: Option<ManifestCache>,
/// Optional wrapper for remote stores used by uploads.
upload_store_wrapper: Option<WriteCacheUploadStoreWrapperRef>,
}
pub type WriteCacheRef = Arc<WriteCache>;
@@ -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<WriteCacheUploadStoreWrapperRef>,
) -> 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,
+7 -2
View File
@@ -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::<WriteCacheUploadStoreWrapperRef>();
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<WriteCacheUploadStoreWrapperRef>,
) -> Result<Option<WriteCacheRef>> {
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)))
}