diff --git a/src/mito2/src/access_layer.rs b/src/mito2/src/access_layer.rs index 337386aa2ff..685cc0d2fc7 100644 --- a/src/mito2/src/access_layer.rs +++ b/src/mito2/src/access_layer.rs @@ -18,6 +18,7 @@ use std::time::{Duration, Instant}; use async_stream::try_stream; use common_base::readable_size::ReadableSize; +use common_runtime::runtime::RuntimeTrait; use common_telemetry::warn; use common_time::Timestamp; use futures::{Stream, TryStreamExt}; @@ -348,6 +349,7 @@ impl AccessLayer { write_opts: &WriteOptions, metrics: &mut Metrics, ) -> Result { + let op_type = request.op_type; let region_id = request.metadata.region_id; let region_metadata = request.metadata.clone(); let cache_manager = request.cache_manager.clone(); @@ -426,6 +428,10 @@ impl AccessLayer { // Put parquet metadata to cache manager. if !sst_info.is_empty() && cache_manager.sst_meta_cache_enabled() { + let runtime = match op_type { + OperationType::Compact => common_runtime::compact_runtime(), + OperationType::Flush => common_runtime::global_runtime(), + }; for sst in &sst_info { if let Some(parquet_metadata) = &sst.file_metadata { let file_id = RegionFileId::new(region_id, sst.file_id); @@ -445,7 +451,7 @@ impl AccessLayer { // Compact cache preparation is best-effort. Run the entire operation in one // detached blocking task so it neither blocks an async worker nor delays the // SST write. - common_runtime::spawn_blocking_global(move || { + runtime.spawn_blocking(move || { match prepare_sst_meta_sync( &file_path, Arc::unwrap_or_clone(parquet_metadata), diff --git a/src/mito2/src/cache.rs b/src/mito2/src/cache.rs index 0d43d31f434..ef150dfb070 100644 --- a/src/mito2/src/cache.rs +++ b/src/mito2/src/cache.rs @@ -31,6 +31,8 @@ use std::sync::{Arc, RwLock, Weak}; use bytes::Bytes; use common_base::readable_size::ReadableSize; use common_datasource::compression::CompressionType; +use common_runtime::Runtime; +use common_runtime::runtime::RuntimeTrait; use common_telemetry::warn; use datatypes::arrow::buffer::BooleanBuffer; use datatypes::arrow::record_batch::RecordBatch; @@ -238,7 +240,7 @@ impl SstMetaPreparation { } impl CompactSstMeta { - async fn decode(&self) -> Result> { + async fn decode(&self, runtime: &Runtime) -> Result> { let mut decoded_guard = self.decoded.lock().await; if let Some(decoded) = decoded_guard.upgrade() { return Ok(decoded); @@ -248,33 +250,34 @@ impl CompactSstMeta { let decoded_size = self.decoded_size; let region_metadata = self.region_metadata.upgrade(); let page_index_policy = self.page_index_policy; - let decoded = common_runtime::spawn_blocking_global(move || { - let bytes = zstd::bulk::decompress(&encoded_metadata, decoded_size).context( - DecompressObjectSnafu { - compress_type: CompressionType::Zstd, + let decoded = runtime + .spawn_blocking(move || { + let bytes = zstd::bulk::decompress(&encoded_metadata, decoded_size).context( + DecompressObjectSnafu { + compress_type: CompressionType::Zstd, + path: "cached SST metadata", + }, + )?; + let bytes = Bytes::from(bytes); + let mut reader = ParquetMetaDataReader::new() + .with_column_index_policy(PageIndexPolicy::Skip) + .with_offset_index_policy(page_index_policy); + reader.try_parse(&bytes).context(ReadParquetSnafu { path: "cached SST metadata", - }, - )?; - let bytes = Bytes::from(bytes); - let mut reader = ParquetMetaDataReader::new() - .with_column_index_policy(PageIndexPolicy::Skip) - .with_offset_index_policy(page_index_policy); - reader.try_parse(&bytes).context(ReadParquetSnafu { - path: "cached SST metadata", - })?; - let metadata = reader.finish().context(ReadParquetSnafu { - path: "cached SST metadata", - })?; - CachedSstMeta::try_new_with_page_index_policy( - "cached SST metadata", - metadata, - region_metadata, - page_index_policy, - ) - .map(Arc::new) - }) - .await - .context(JoinSnafu)??; + })?; + let metadata = reader.finish().context(ReadParquetSnafu { + path: "cached SST metadata", + })?; + CachedSstMeta::try_new_with_page_index_policy( + "cached SST metadata", + metadata, + region_metadata, + page_index_policy, + ) + .map(Arc::new) + }) + .await + .context(JoinSnafu)??; *decoded_guard = Arc::downgrade(&decoded); Ok(decoded) @@ -285,46 +288,50 @@ impl CompactSstMeta { } } -/// Decodes SST metadata without preparing a compact cache entry. +/// Decodes SST metadata on the given blocking runtime without preparing a compact cache entry. pub(crate) async fn decode_sst_meta( file_path: &str, parquet_metadata: ParquetMetaData, region_metadata: Option, page_index_policy: PageIndexPolicy, + runtime: &Runtime, ) -> Result> { let file_path = file_path.to_string(); - common_runtime::spawn_blocking_global(move || { - let parquet_metadata = strip_column_indexes(parquet_metadata); - CachedSstMeta::try_new_with_page_index_policy( - &file_path, - parquet_metadata, - region_metadata, - page_index_policy, - ) - .map(Arc::new) - }) - .await - .context(JoinSnafu)? + runtime + .spawn_blocking(move || { + let parquet_metadata = strip_column_indexes(parquet_metadata); + CachedSstMeta::try_new_with_page_index_policy( + &file_path, + parquet_metadata, + region_metadata, + page_index_policy, + ) + .map(Arc::new) + }) + .await + .context(JoinSnafu)? } -/// Decodes SST metadata and attempts to encode both cache representations on the blocking runtime. +/// Decodes SST metadata and attempts to encode both cache representations on the given blocking runtime. pub(crate) async fn prepare_sst_meta( file_path: &str, parquet_metadata: ParquetMetaData, region_metadata: Option, page_index_policy: PageIndexPolicy, + runtime: &Runtime, ) -> Result { let file_path = file_path.to_string(); - common_runtime::spawn_blocking_global(move || { - prepare_sst_meta_sync( - &file_path, - parquet_metadata, - region_metadata, - page_index_policy, - ) - }) - .await - .context(JoinSnafu)? + runtime + .spawn_blocking(move || { + prepare_sst_meta_sync( + &file_path, + parquet_metadata, + region_metadata, + page_index_policy, + ) + }) + .await + .context(JoinSnafu)? } /// Synchronously prepares SST metadata. Callers must run this on a blocking runtime. @@ -671,6 +678,16 @@ pub enum CacheStrategy { } impl CacheStrategy { + /// Returns the runtime for CPU-bound SST metadata work for this request. + pub(crate) fn sst_meta_runtime(&self) -> Runtime { + match self { + CacheStrategy::Compaction(_) => common_runtime::compact_runtime(), + CacheStrategy::EnableAll(_) | CacheStrategy::Disabled => { + common_runtime::global_runtime() + } + } + } + /// Returns whether the SST metadata cache is enabled for this strategy. pub(crate) fn sst_meta_cache_enabled(&self) -> bool { match self { @@ -691,7 +708,12 @@ impl CacheStrategy { match self { CacheStrategy::EnableAll(cache_manager) | CacheStrategy::Compaction(cache_manager) => { cache_manager - .get_sst_meta_data(file_id, metrics, page_index_policy) + .get_sst_meta_data( + file_id, + metrics, + page_index_policy, + &self.sst_meta_runtime(), + ) .await } CacheStrategy::Disabled => { @@ -1057,6 +1079,7 @@ impl CacheManager { file_id: RegionFileId, metrics: &mut MetadataCacheMetrics, page_index_policy: PageIndexPolicy, + runtime: &Runtime, ) -> Option> { let cache_key = SstMetaKey(file_id.region_id(), file_id.file_id()); let compact = self @@ -1075,7 +1098,7 @@ impl CacheManager { } CACHE_MISS.with_label_values(&[SST_META_DECODED_TYPE]).inc(); - match compact.decode().await { + match compact.decode(runtime).await { Ok(decoded) => { self.put_sst_meta_data(file_id, decoded.clone()); return Some(decoded); @@ -1101,7 +1124,7 @@ impl CacheManager { let file_cache = write_cache.file_cache(); if self.sst_meta_cache_enabled() { if let Some(metadata) = file_cache - .get_sst_meta_data(key, metrics, page_index_policy) + .get_sst_meta_data(key, metrics, page_index_policy, runtime) .await { metrics.file_cache_hit += 1; @@ -1122,7 +1145,7 @@ impl CacheManager { return Some(decoded); } } else if let Some(decoded) = file_cache - .get_decoded_sst_meta_data(key, metrics, page_index_policy) + .get_decoded_sst_meta_data(key, metrics, page_index_policy, runtime) .await { metrics.file_cache_hit += 1; @@ -1209,7 +1232,7 @@ impl CacheManager { let compact = self .get_compact_sst_meta(&key) .filter(|metadata| metadata.satisfies_page_index_policy(page_index_policy))?; - match compact.decode().await { + match compact.decode(&common_runtime::global_runtime()).await { Ok(metadata) => Some(metadata), Err(err) => { warn!(err; "Failed to decode compact SST metadata, region_id: {}, file_id: {}", file_id.region_id(), file_id.file_id()); @@ -2144,7 +2167,12 @@ mod tests { cache.put_parquet_meta_data(file_id, metadata, None); assert!( cache - .get_sst_meta_data(file_id, &mut metrics, Default::default()) + .get_sst_meta_data( + file_id, + &mut metrics, + Default::default(), + &common_runtime::global_runtime() + ) .await .is_none() ); @@ -2205,14 +2233,24 @@ mod tests { let file_id = RegionFileId::new(region_id, FileId::random()); assert!( cache - .get_sst_meta_data(file_id, &mut metrics, Default::default()) + .get_sst_meta_data( + file_id, + &mut metrics, + Default::default(), + &common_runtime::global_runtime() + ) .await .is_none() ); let (metadata, region_metadata) = sst_parquet_meta(); cache.put_parquet_meta_data(file_id, metadata, None); let cached = cache - .get_sst_meta_data(file_id, &mut metrics, Default::default()) + .get_sst_meta_data( + file_id, + &mut metrics, + Default::default(), + &common_runtime::global_runtime(), + ) .await .unwrap(); assert_eq!(region_metadata, cached.region_metadata()); @@ -2230,7 +2268,12 @@ mod tests { cache.remove_parquet_meta_data(file_id); assert!( cache - .get_sst_meta_data(file_id, &mut metrics, Default::default()) + .get_sst_meta_data( + file_id, + &mut metrics, + Default::default(), + &common_runtime::global_runtime() + ) .await .is_none() ); @@ -2247,7 +2290,12 @@ mod tests { cache.put_parquet_meta_data(file_id, metadata, Some(region_metadata.clone())); let cached = cache - .get_sst_meta_data(file_id, &mut metrics, Default::default()) + .get_sst_meta_data( + file_id, + &mut metrics, + Default::default(), + &common_runtime::global_runtime(), + ) .await .unwrap(); assert!(Arc::ptr_eq(®ion_metadata, &cached.region_metadata())); @@ -2285,9 +2333,15 @@ mod tests { .set_offset_index(Some(offset_indexes)) .build(); let expected_rows = metadata.file_metadata().num_rows(); - let prepared = prepare_sst_meta("test.parquet", metadata, None, PageIndexPolicy::Required) - .await - .unwrap(); + let prepared = prepare_sst_meta( + "test.parquet", + metadata, + None, + PageIndexPolicy::Required, + &common_runtime::global_runtime(), + ) + .await + .unwrap(); let SstMetaPreparation::Prepared(prepared) = prepared else { panic!("valid metadata should produce a compact cache entry"); }; @@ -2296,7 +2350,12 @@ mod tests { cache.put_prepared_sst_meta(file_id, prepared, false); let mut metrics = MetadataCacheMetrics::default(); let cached = cache - .get_sst_meta_data(file_id, &mut metrics, PageIndexPolicy::Required) + .get_sst_meta_data( + file_id, + &mut metrics, + PageIndexPolicy::Required, + &common_runtime::global_runtime(), + ) .await .unwrap(); @@ -2377,7 +2436,12 @@ mod tests { let mut metrics = MetadataCacheMetrics::default(); assert!( cache - .get_sst_meta_data(file_id, &mut metrics, PageIndexPolicy::Optional) + .get_sst_meta_data( + file_id, + &mut metrics, + PageIndexPolicy::Optional, + &common_runtime::global_runtime() + ) .await .is_none() ); @@ -2397,7 +2461,12 @@ mod tests { let mut metrics = MetadataCacheMetrics::default(); assert!( cache - .get_sst_meta_data(file_id, &mut metrics, PageIndexPolicy::Optional) + .get_sst_meta_data( + file_id, + &mut metrics, + PageIndexPolicy::Optional, + &common_runtime::global_runtime() + ) .await .is_some() ); @@ -2406,7 +2475,12 @@ mod tests { let mut metrics = MetadataCacheMetrics::default(); assert!( cache - .get_sst_meta_data(file_id, &mut metrics, PageIndexPolicy::Skip) + .get_sst_meta_data( + file_id, + &mut metrics, + PageIndexPolicy::Skip, + &common_runtime::global_runtime() + ) .await .is_some() ); diff --git a/src/mito2/src/cache/file_cache.rs b/src/mito2/src/cache/file_cache.rs index 0594a1ac10d..4cb60ed8d48 100644 --- a/src/mito2/src/cache/file_cache.rs +++ b/src/mito2/src/cache/file_cache.rs @@ -21,6 +21,7 @@ use std::time::{Duration, Instant}; use bytes::Bytes; use common_base::readable_size::ReadableSize; +use common_runtime::Runtime; use common_telemetry::{debug, error, info, warn}; use futures::{AsyncWriteExt, FutureExt, TryStreamExt}; use moka::future::Cache; @@ -627,12 +628,13 @@ impl FileCache { key: IndexKey, cache_metrics: &mut MetadataCacheMetrics, page_index_policy: PageIndexPolicy, + runtime: &Runtime, ) -> Option { let file_path = self.inner.cache_file_path(key); let metadata = self .get_parquet_meta_data(key, cache_metrics, page_index_policy) .await?; - match prepare_sst_meta(&file_path, metadata, None, page_index_policy).await { + match prepare_sst_meta(&file_path, metadata, None, page_index_policy, runtime).await { Ok(metadata) => Some(metadata), Err(err) => { CACHE_MISS @@ -653,12 +655,13 @@ impl FileCache { key: IndexKey, cache_metrics: &mut MetadataCacheMetrics, page_index_policy: PageIndexPolicy, + runtime: &Runtime, ) -> Option> { let file_path = self.inner.cache_file_path(key); let metadata = self .get_parquet_meta_data(key, cache_metrics, page_index_policy) .await?; - match decode_sst_meta(&file_path, metadata, None, page_index_policy).await { + match decode_sst_meta(&file_path, metadata, None, page_index_policy, runtime).await { Ok(metadata) => Some(metadata), Err(err) => { CACHE_MISS diff --git a/src/mito2/src/read/pruner.rs b/src/mito2/src/read/pruner.rs index 19c487ed0b8..ef09f5f48ab 100644 --- a/src/mito2/src/read/pruner.rs +++ b/src/mito2/src/read/pruner.rs @@ -350,12 +350,14 @@ impl Pruner { enable_predicate_prefilter, }); - // Spawn worker tasks with their receivers + // Keep pruning and prefetching on the runtime of the originating workload. for (worker_id, rx) in receivers.into_iter().enumerate() { - let inner_clone = inner.clone(); - common_runtime::spawn_query(async move { - Self::worker_loop(worker_id, rx, inner_clone).await; - }); + let worker = Self::worker_loop(worker_id, rx, inner.clone()); + if inner.stream_ctx.input.compaction { + common_runtime::spawn_compact(worker); + } else { + common_runtime::spawn_query(worker); + } } Self { diff --git a/src/mito2/src/region/opener.rs b/src/mito2/src/region/opener.rs index fb5d260fe6f..84a0e13833f 100644 --- a/src/mito2/src/region/opener.rs +++ b/src/mito2/src/region/opener.rs @@ -1247,7 +1247,12 @@ async fn preload_parquet_meta_cache_for_files( let key = IndexKey::new(file_id.region_id(), file_id.file_id(), FileType::Parquet); if let Some(metadata) = write_cache .file_cache() - .get_sst_meta_data(key, &mut cache_metrics, PageIndexPolicy::Optional) + .get_sst_meta_data( + key, + &mut cache_metrics, + PageIndexPolicy::Optional, + &common_runtime::global_runtime(), + ) .await { let decoded = metadata.decoded(); @@ -1297,6 +1302,7 @@ async fn preload_parquet_meta_cache_for_files( // instead of substituting the region's current schema. None, PageIndexPolicy::Optional, + &common_runtime::global_runtime(), ) .await { diff --git a/src/mito2/src/sst/index.rs b/src/mito2/src/sst/index.rs index 0c3ede4fb06..ed3ec36a473 100644 --- a/src/mito2/src/sst/index.rs +++ b/src/mito2/src/sst/index.rs @@ -57,8 +57,8 @@ use crate::cache::{CacheManagerRef, CacheStrategy}; use crate::config::VectorIndexConfig; use crate::config::{BloomFilterConfig, FulltextIndexConfig, InvertedIndexConfig}; use crate::error::{ - BuildIndexAsyncSnafu, Error, RegionClosedSnafu, RegionDroppedSnafu, RegionTruncatedSnafu, - Result, + BuildIndexAsyncSnafu, Error, JoinSnafu, RegionClosedSnafu, RegionDroppedSnafu, + RegionTruncatedSnafu, Result, }; use crate::metrics::{ INDEX_ARTIFACT_CLEANUP_FAILURE_TOTAL, INDEX_CREATE_MEMORY_USAGE, INDEX_PUBLICATION_STALE_TOTAL, @@ -774,7 +774,21 @@ impl IndexBuildTask { self.source.file_meta.file_id, )) .await; - match self.index_build(version_control).await { + let result = if self.reason == IndexBuildType::Compact { + let mut task = self.clone(); + // Keep the scheduler slot occupied until the compact runtime finishes the build. + match common_runtime::spawn_compact( + async move { task.index_build(version_control).await }, + ) + .await + { + Ok(result) => result, + Err(err) => Err(err).context(JoinSnafu), + } + } else { + self.index_build(version_control).await + }; + match result { Ok(outcome) => self.on_success(outcome).await, Err(e) => { warn!( diff --git a/src/mito2/src/sst/parquet/reader.rs b/src/mito2/src/sst/parquet/reader.rs index f54812f97d4..9584ff0a6af 100644 --- a/src/mito2/src/sst/parquet/reader.rs +++ b/src/mito2/src/sst/parquet/reader.rs @@ -851,7 +851,14 @@ impl ParquetReaderBuilder { let metadata = metadata_loader.load(cache_metrics).await?; let decoded = if self.cache_strategy.sst_meta_cache_enabled() { - let metadata = prepare_sst_meta(file_path, metadata, None, page_index_policy).await?; + let metadata = prepare_sst_meta( + file_path, + metadata, + None, + page_index_policy, + &self.cache_strategy.sst_meta_runtime(), + ) + .await?; let decoded = metadata.decoded(); match metadata { SstMetaPreparation::Prepared(metadata) => {