fix: keep compaction pruning, metadata, and index work on compact runtime (#9304)

* fix: run compaction pruner tasks on compact runtime

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix: keep compaction metadata and index work on compact runtime

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

---------

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>
This commit is contained in:
Lei, HUANG
2026-09-22 13:12:39 +00:00
committed by GitHub
parent 045441e3cc
commit a9e2a89b7b
7 changed files with 193 additions and 81 deletions
+7 -1
View File
@@ -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<SstInfoArray> {
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),
+142 -68
View File
@@ -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<Arc<CachedSstMeta>> {
async fn decode(&self, runtime: &Runtime) -> Result<Arc<CachedSstMeta>> {
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<RegionMetadataRef>,
page_index_policy: PageIndexPolicy,
runtime: &Runtime,
) -> Result<Arc<CachedSstMeta>> {
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<RegionMetadataRef>,
page_index_policy: PageIndexPolicy,
runtime: &Runtime,
) -> Result<SstMetaPreparation> {
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<Arc<CachedSstMeta>> {
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(&region_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()
);
+5 -2
View File
@@ -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<SstMetaPreparation> {
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<Arc<CachedSstMeta>> {
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
+7 -5
View File
@@ -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 {
+7 -1
View File
@@ -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
{
+17 -3
View File
@@ -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!(
+8 -1
View File
@@ -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) => {