feat(mito2): pass operation type to write cache upload hook (#9012)

Extend `WriteCacheUploadStoreWrapper::wrap` with the `OperationType` of the upload so implementations can apply per-operation policies (e.g. throttling compaction uploads but not flush uploads). Flush and compaction paths forward their existing `SstWriteRequest::op_type`; `put_and_upload_sst` is flush-only and index rebuild uploads are reported as compaction uploads.

Files: `src/mito2/src/cache/write_cache.rs`, `src/mito2/src/sst/index.rs`.

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>
This commit is contained in:
Lei, HUANG
2026-09-03 08:10:05 +00:00
committed by GitHub
parent 05c65f54a8
commit d351b7d471
2 changed files with 45 additions and 9 deletions
+37 -8
View File
@@ -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<dyn WriteCacheUploadStoreWrapper>;
@@ -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<Option<OperationType>>,
}
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(),
+8 -1
View File
@@ -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);