diff --git a/crates/paimon/src/common/options.rs b/crates/paimon/src/common/options.rs index c23cd9318..d482b6446 100644 --- a/crates/paimon/src/common/options.rs +++ b/crates/paimon/src/common/options.rs @@ -87,6 +87,9 @@ impl CatalogOptions { /// Comma-separated file types eligible for local caching. pub const LOCAL_CACHE_WHITELIST: &'static str = "local-cache.whitelist"; + + /// Comma-separated file extensions that bypass the local cache. + pub const LOCAL_CACHE_EXCLUDE_EXTENSIONS: &'static str = "local-cache.exclude-extensions"; } /// Configuration options container. diff --git a/crates/paimon/src/io/cache/mod.rs b/crates/paimon/src/io/cache/mod.rs index 3f661a47c..3561686a1 100644 --- a/crates/paimon/src/io/cache/mod.rs +++ b/crates/paimon/src/io/cache/mod.rs @@ -44,6 +44,7 @@ pub(crate) struct LocalCache { namespace: String, block_size: u64, whitelist: HashSet, + excluded_extensions: HashSet, file_size_capacity: usize, } @@ -71,6 +72,7 @@ impl LocalCache { namespace: String::new(), block_size, whitelist: FileType::parse_whitelist(whitelist), + excluded_extensions: HashSet::new(), file_size_capacity: DEFAULT_FILE_SIZE_CAPACITY, }) } @@ -98,6 +100,7 @@ impl LocalCache { namespace: config.namespace, block_size: config.block_size, whitelist: config.whitelist, + excluded_extensions: config.excluded_extensions, file_size_capacity, }) } @@ -111,7 +114,16 @@ impl LocalCache { } pub(super) fn is_cacheable(&self, path: &str) -> bool { - !FileType::is_mutable(path) && self.whitelist.contains(&FileType::classify(path)) + let extension = path + .rsplit('/') + .next() + .and_then(|name| name.rsplit_once('.')) + .filter(|(stem, _)| !stem.is_empty()); + !FileType::is_mutable(path) + && !extension.is_some_and(|(_, ext)| { + self.excluded_extensions.contains(&ext.to_ascii_lowercase()) + }) + && self.whitelist.contains(&FileType::classify(path)) } async fn get_block( @@ -247,6 +259,7 @@ pub(crate) struct LocalCacheConfig { max_size: Option, block_size: u64, whitelist: HashSet, + excluded_extensions: HashSet, } impl LocalCacheConfig { @@ -313,6 +326,14 @@ impl LocalCacheConfig { max_size, block_size, whitelist: FileType::parse_whitelist(whitelist), + excluded_extensions: options + .get(CatalogOptions::LOCAL_CACHE_EXCLUDE_EXTENSIONS) + .map(String::as_str) + .unwrap_or("") + .split(',') + .map(|value| value.trim().trim_start_matches('.').to_ascii_lowercase()) + .filter(|value| !value.is_empty()) + .collect(), })) } } @@ -473,6 +494,7 @@ mod tests { max_size: None, block_size: 4, whitelist: HashSet::from([FileType::Meta]), + excluded_extensions: HashSet::new(), }) .unwrap(); let path = "s3://bucket/table/snapshot/snapshot-1"; @@ -494,6 +516,7 @@ mod tests { max_size: None, block_size: 4, whitelist: HashSet::from([FileType::Meta]), + excluded_extensions: HashSet::new(), }; let first = LocalCache::new(config()).unwrap(); let second = LocalCache::new(config()).unwrap(); @@ -516,6 +539,7 @@ mod tests { max_size: None, block_size: 4, whitelist: HashSet::from([FileType::Meta]), + excluded_extensions: HashSet::new(), }; let first = LocalCache::new(config()).unwrap(); let second = LocalCache::new(config()).unwrap(); @@ -539,10 +563,12 @@ mod tests { max_size: None, block_size: 4, whitelist: HashSet::from([FileType::Meta]), + excluded_extensions: HashSet::from(["blob".to_string(), "snapshot-1".to_string()]), }) .unwrap(); assert!(cache.is_cacheable("s3://bucket/table/snapshot/snapshot-1")); + assert!(!cache.is_cacheable("s3://bucket/table/snapshot/snapshot-1.blob")); assert!(!cache.is_cacheable("s3://bucket/table/data/data-1.parquet")); assert!(!cache.is_cacheable("s3://bucket/table/snapshot/LATEST")); assert!(!cache.is_cacheable("s3://bucket/table/tag/tag-production")); @@ -566,6 +592,7 @@ mod tests { max_size: None, block_size: 4, whitelist: HashSet::from([FileType::Meta]), + excluded_extensions: HashSet::new(), }) .unwrap(); @@ -586,6 +613,7 @@ mod tests { max_size: Some(8), block_size: 4, whitelist: HashSet::from([FileType::Meta]), + excluded_extensions: HashSet::new(), }) .unwrap(); let first_token = cache.read_token("snapshot-1"); diff --git a/crates/paimon/src/io/cache/reader.rs b/crates/paimon/src/io/cache/reader.rs index 240b163d1..6ea1cfb00 100644 --- a/crates/paimon/src/io/cache/reader.rs +++ b/crates/paimon/src/io/cache/reader.rs @@ -275,6 +275,7 @@ mod tests { max_size: None, block_size: 4, whitelist: std::collections::HashSet::from([FileType::Meta]), + excluded_extensions: std::collections::HashSet::new(), }) .unwrap(), ); @@ -325,6 +326,7 @@ mod tests { max_size: None, block_size: 4, whitelist: std::collections::HashSet::from([FileType::Meta]), + excluded_extensions: std::collections::HashSet::new(), }) .unwrap(), ); @@ -355,6 +357,7 @@ mod tests { max_size: None, block_size: 4, whitelist: std::collections::HashSet::from([FileType::Meta]), + excluded_extensions: std::collections::HashSet::new(), }; let first_reader = CachedFileReader::new( delegate.clone(), @@ -402,6 +405,7 @@ mod tests { max_size: None, block_size: 4, whitelist: std::collections::HashSet::from([FileType::Meta]), + excluded_extensions: std::collections::HashSet::new(), }) .unwrap(), ); @@ -451,6 +455,7 @@ mod tests { max_size: None, block_size: 4, whitelist: std::collections::HashSet::from([FileType::Meta]), + excluded_extensions: std::collections::HashSet::new(), }) .unwrap(), ); @@ -461,6 +466,7 @@ mod tests { max_size: None, block_size: 4, whitelist: std::collections::HashSet::from([FileType::Meta]), + excluded_extensions: std::collections::HashSet::new(), }) .unwrap(), ); @@ -512,6 +518,7 @@ mod tests { max_size: None, block_size: 4, whitelist: std::collections::HashSet::from([FileType::Meta]), + excluded_extensions: std::collections::HashSet::new(), }; let cache_a = Arc::new(LocalCache::new(config()).unwrap()); let cache_b = Arc::new(LocalCache::new(config()).unwrap()); @@ -658,6 +665,7 @@ mod tests { max_size: None, block_size: 4, whitelist: std::collections::HashSet::from([FileType::Meta]), + excluded_extensions: std::collections::HashSet::new(), }) .unwrap(), ); @@ -692,6 +700,7 @@ mod tests { max_size: None, block_size: 4, whitelist: std::collections::HashSet::from([FileType::Meta]), + excluded_extensions: std::collections::HashSet::new(), }) .unwrap(), ); @@ -713,6 +722,7 @@ mod tests { max_size: None, block_size: 4, whitelist: std::collections::HashSet::from([FileType::Meta]), + excluded_extensions: std::collections::HashSet::new(), }) .unwrap(), ); diff --git a/crates/paimon/src/io/file_io.rs b/crates/paimon/src/io/file_io.rs index bc8e2fb56..f844988e8 100644 --- a/crates/paimon/src/io/file_io.rs +++ b/crates/paimon/src/io/file_io.rs @@ -1625,6 +1625,78 @@ mod input_output_test { .unwrap() } + #[tokio::test] + async fn test_data_cache_excludes_blob_extension() { + for disk in [false, true] { + for whitelist in ["meta,global-index", "data"] { + for exclusion in ["", " .BLOB, "] { + let directory = tempfile::tempdir().unwrap(); + let mut options = Options::new(); + options.set(CatalogOptions::LOCAL_CACHE_ENABLED, "true"); + options.set(CatalogOptions::LOCAL_CACHE_WHITELIST, whitelist); + options.set(CatalogOptions::LOCAL_CACHE_EXCLUDE_EXTENSIONS, exclusion); + options.set(CatalogOptions::LOCAL_CACHE_BLOCK_SIZE, "4"); + if disk { + options.set( + CatalogOptions::LOCAL_CACHE_DIR, + directory.path().to_string_lossy(), + ); + } + let cache = Arc::new( + LocalCache::new(LocalCacheConfig::from_options(&options).unwrap().unwrap()) + .unwrap(), + ); + let file_io = FileIOBuilder::new("memory") + .with_local_cache(cache) + .build() + .unwrap(); + for name in [ + "data.parquet", + "data.blob", + "data.orc", + "data.avro", + "data.parquet.index", + "snapshot-1", + "snapshot-1.blob", + ] { + let path = format!("memory:/{name}"); + file_io + .new_output(&path) + .unwrap() + .write(Bytes::from_static(b"abcdefgh")) + .await + .unwrap(); + let input = file_io.new_input(&path).unwrap(); + let reader = input.reader().await.unwrap(); + assert_eq!(reader.read(1..7).await.unwrap(), b"bcdefg"[..]); + drop(reader); + let (op, relative_path, _) = input.source.resolve(&path).await.unwrap(); + op.delete(&relative_path).await.unwrap(); + let cached = match whitelist { + "data" => { + name.starts_with("data.") + && !(name.ends_with(".blob") && !exclusion.is_empty()) + && !name.ends_with(".index") + } + _ => { + name == "snapshot-1" + || (name == "snapshot-1.blob" && exclusion.is_empty()) + } + }; + let result = input.read().await; + if cached { + assert_eq!(result.unwrap(), b"abcdefgh"[..], "{whitelist}: {name}"); + let reader = input.reader().await.unwrap(); + assert_eq!(reader.read(1..7).await.unwrap(), b"bcdefg"[..]); + } else { + assert!(result.is_err(), "{whitelist}: {name}"); + } + } + } + } + } + } + async fn common_test_output_file_write_and_read(file_io: &FileIO, path: &str) { let output = file_io.new_output(path).unwrap(); let mut writer = output.writer().await.unwrap(); diff --git a/docs/src/getting-started.md b/docs/src/getting-started.md index de526dff4..68a1c1ae2 100644 --- a/docs/src/getting-started.md +++ b/docs/src/getting-started.md @@ -325,6 +325,9 @@ let catalog = CatalogFactory::create(options).await?; | `local-cache.max-size` | unlimited | Maximum cache size. Memory caches count payload bytes; disk caches count encoded bytes. Values accept byte units such as `512 MiB` or `20 GiB`. | | `local-cache.block-size` | `1 MiB` | Block size used for cached range reads. | | `local-cache.whitelist` | `meta,global-index` | Comma-separated eligible types: `meta`, `global-index`, `bucket-index`, `data`, and `file-index`. | +| `local-cache.exclude-extensions` | empty | Comma-separated file extensions to bypass, even when whitelisted. | + +Use `meta,global-index,data` with `local-cache.exclude-extensions=blob` to cache data files except BLOB bodies. Each catalog owns its in-memory cache for the catalog's lifetime. Disk caches are reused after process restarts. Cache keys include a catalog-configuration