diff --git a/crates/paimon/src/arrow/format/blob.rs b/crates/paimon/src/arrow/format/blob.rs index a08deb81a..720e4d5dc 100644 --- a/crates/paimon/src/arrow/format/blob.rs +++ b/crates/paimon/src/arrow/format/blob.rs @@ -33,6 +33,7 @@ use async_stream::try_stream; use async_trait::async_trait; use bytes::Bytes; use futures::{StreamExt, TryStreamExt}; +use lru::LruCache; use std::ops::Range; use std::sync::Arc; @@ -60,7 +61,7 @@ impl BlobFormatReader { pub(crate) struct IndexedBlobReader { reader: Box, - index: BlobFileIndex, + index: Arc, descriptor_mode: bool, file_path: String, blob_parallelism: usize, @@ -92,7 +93,7 @@ impl IndexedBlobReader { blob_parallelism: usize, ) -> crate::Result { debug_assert!(blob_parallelism > 0); - let index = BlobFileIndex::load(reader.as_ref(), file_size).await?; + let index = BlobFileIndex::load_cached(reader.as_ref(), file_size, &file_path).await?; Ok(Self { reader, index, @@ -162,6 +163,20 @@ pub(crate) enum BlobReadValue { const BLOB_FOOTER_SIZE: u64 = 5; const BLOB_FORMAT_VERSION: u8 = 1; +const BLOB_INDEX_CACHE_CAPACITY: usize = 16; +#[derive(Clone, Debug, Eq, Hash, PartialEq)] +struct BlobIndexCacheKey { + namespace: usize, + file_path: String, +} + +static BLOB_INDEX_CACHE: std::sync::LazyLock< + std::sync::Mutex>>, +> = std::sync::LazyLock::new(|| { + std::sync::Mutex::new(LruCache::new( + std::num::NonZeroUsize::new(BLOB_INDEX_CACHE_CAPACITY).unwrap(), + )) +}); const BLOB_MAGIC_NUMBER: i32 = 1481511375; const BLOB_MAGIC_NUMBER_BYTES: [u8; 4] = BLOB_MAGIC_NUMBER.to_le_bytes(); const BLOB_INLINE_HEADER_SIZE: u64 = 4; @@ -1652,12 +1667,43 @@ struct BlobArrayLayout { element_index_range: Range, } -#[derive(Debug, Clone)] +#[derive(Debug)] struct BlobFileIndex { entries: Vec, } impl BlobFileIndex { + async fn load_cached( + reader: &dyn FileRead, + file_size: u64, + file_path: &str, + ) -> crate::Result> { + let cache_key = reader + .cache_namespace() + .filter(|_| !file_path.is_empty()) + .map(|namespace| BlobIndexCacheKey { + namespace, + file_path: file_path.to_string(), + }); + if let Some(cache_key) = &cache_key { + let mut cache = BLOB_INDEX_CACHE + .lock() + .unwrap_or_else(|error| error.into_inner()); + if let Some(index) = cache.get(cache_key) { + return Ok(index.clone()); + } + } + + let index = Arc::new(Self::load(reader, file_size).await?); + if let Some(cache_key) = cache_key { + BLOB_INDEX_CACHE + .lock() + .unwrap_or_else(|error| error.into_inner()) + .put(cache_key, index.clone()); + } + Ok(index) + } + async fn load(reader: &dyn FileRead, file_size: u64) -> crate::Result { if file_size < BLOB_FOOTER_SIZE { return Err(Error::DataInvalid { @@ -2229,6 +2275,7 @@ fn encode_varint(value: i64, out: &mut Vec) { mod tests { use super::*; use crate::btree::test_util::BytesFileRead; + use crate::io::{FileIO, FileIOBuilder}; use crate::spec::{ArrayType, BlobType, MapType, VarCharType}; use arrow_array::Array; use bytes::Bytes; @@ -2298,6 +2345,86 @@ mod tests { ); } + #[tokio::test] + async fn test_blob_reader_reuses_cached_index() { + let file_path = "file:///blob-index-cache-test/data.blob"; + let file_bytes = load_blob_fixture("blob-basic.blob"); + let first = + TrackingFileRead::new(Bytes::from(file_bytes.clone())).with_cache_namespace(usize::MAX); + let second = + TrackingFileRead::new(Bytes::from(file_bytes.clone())).with_cache_namespace(usize::MAX); + + let first_reader = IndexedBlobReader::open( + Box::new(first.clone()), + file_bytes.len() as u64, + file_path.to_string(), + true, + ) + .await + .unwrap(); + let second_reader = IndexedBlobReader::open( + Box::new(second.clone()), + file_bytes.len() as u64, + file_path.to_string(), + true, + ) + .await + .unwrap(); + + assert_eq!(first_reader.num_rows(), second_reader.num_rows()); + assert_eq!(first.ranges().len(), 2); + assert!(second.ranges().is_empty()); + } + + #[tokio::test] + async fn test_blob_index_cache_isolated_by_file_io() { + let path = "memory:///blob-index-cache-namespace/data.blob"; + let value = b"value"; + let first_bytes = blob_test_utils::build_blob_file_bytes(&[None, Some(value.as_slice())]); + let second_bytes = blob_test_utils::build_blob_file_bytes(&[Some(value.as_slice()), None]); + assert_eq!(first_bytes.len(), second_bytes.len()); + + let first_io = FileIOBuilder::new("memory").build().unwrap(); + let second_io = FileIOBuilder::new("memory").build().unwrap(); + first_io + .new_output(path) + .unwrap() + .write(Bytes::from(first_bytes)) + .await + .unwrap(); + second_io + .new_output(path) + .unwrap() + .write(Bytes::from(second_bytes)) + .await + .unwrap(); + + let first_namespace = first_io + .new_input(path) + .unwrap() + .reader() + .await + .unwrap() + .cache_namespace(); + let second_namespace = second_io + .new_input(path) + .unwrap() + .reader() + .await + .unwrap() + .cache_namespace(); + assert_ne!(first_namespace, second_namespace); + + assert_eq!( + read_scalar_blob_file(&first_io, path).await, + vec![None, Some(value.to_vec())] + ); + assert_eq!( + read_scalar_blob_file(&second_io, path).await, + vec![Some(value.to_vec()), None] + ); + } + #[tokio::test] async fn test_blob_array_reader_reads_java_fixture() { let read_fields = vec![DataField::new( @@ -2367,6 +2494,7 @@ mod tests { #[tokio::test] async fn test_blob_map_reader_returns_inline_values_and_descriptors() { + let file_path = "file:///tmp/map-values-and-descriptors.blob"; let payload = build_blob_map_payload(&[ ("video", Some(b"alpha")), ("thumbnail", None), @@ -2375,7 +2503,7 @@ mod tests { let file_bytes = blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice()), None]); let fields = blob_map_read_fields(); - let inline = BlobFormatReader::new("file:///tmp/map.blob".to_string(), false) + let inline = BlobFormatReader::new(file_path.to_string(), false) .read_batch_stream( Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))), file_bytes.len() as u64, @@ -2401,7 +2529,7 @@ mod tests { ] ); - let descriptors = BlobFormatReader::new("file:///tmp/map.blob".to_string(), true) + let descriptors = BlobFormatReader::new(file_path.to_string(), true) .read_batch_stream( Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))), file_bytes.len() as u64, @@ -2418,7 +2546,7 @@ mod tests { let rows = collect_blob_map_values(&descriptors[0]); let entries = rows[0].as_ref().unwrap(); let video = BlobDescriptor::deserialize(entries[0].1.as_ref().unwrap()).unwrap(); - assert_eq!(video.uri(), "file:///tmp/map.blob"); + assert_eq!(video.uri(), file_path); assert_eq!(video.length(), 5); assert!(entries[1].1.is_none()); let empty = BlobDescriptor::deserialize(entries[2].1.as_ref().unwrap()).unwrap(); @@ -2468,7 +2596,7 @@ mod tests { #[tokio::test] async fn test_blob_map_descriptor_read_skips_values() { - let file_path = "file:///tmp/map.blob"; + let file_path = "file:///tmp/map-descriptor-skip-values.blob"; let payload = build_blob_map_payload(&[("first", Some(b"alpha")), ("second", Some(b"beta"))]); let file_bytes = blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice())]); @@ -3495,6 +3623,30 @@ mod tests { Ok(batches.iter().flat_map(collect_binary_values).collect()) } + async fn read_scalar_blob_file(file_io: &FileIO, path: &str) -> Vec>> { + let input = file_io.new_input(path).unwrap(); + let file_size = input.metadata().await.unwrap().size; + let batches = BlobFormatReader::new(path.to_string(), false) + .read_batch_stream( + Box::new(input.reader().await.unwrap()), + file_size, + &[DataField::new( + 0, + "payload".to_string(), + DataType::Blob(BlobType::new()), + )], + None, + None, + None, + ) + .await + .unwrap() + .try_collect::>() + .await + .unwrap(); + batches.iter().flat_map(collect_binary_values).collect() + } + fn rewrite_first_blob_entry_crc(file_bytes: &mut [u8], payload_length: usize) { let crc_offset = BLOB_INLINE_HEADER_SIZE as usize + payload_length + size_of::(); let mut hasher = crc32fast::Hasher::new(); @@ -3587,6 +3739,7 @@ mod tests { #[derive(Clone)] struct TrackingFileRead { bytes: Bytes, + cache_namespace: Option, in_flight: Arc, max_in_flight: Arc, ranges: Arc>>>, @@ -3596,12 +3749,18 @@ mod tests { fn new(bytes: Bytes) -> Self { Self { bytes, + cache_namespace: None, in_flight: Arc::new(AtomicUsize::new(0)), max_in_flight: Arc::new(AtomicUsize::new(0)), ranges: Arc::new(Mutex::new(Vec::new())), } } + fn with_cache_namespace(mut self, cache_namespace: usize) -> Self { + self.cache_namespace = Some(cache_namespace); + self + } + fn max_in_flight(&self) -> usize { self.max_in_flight.load(Ordering::SeqCst) } @@ -3621,6 +3780,10 @@ mod tests { self.in_flight.fetch_sub(1, Ordering::SeqCst); Ok(self.bytes.slice(range.start as usize..range.end as usize)) } + + fn cache_namespace(&self) -> Option { + self.cache_namespace + } } struct SparseFileRead { diff --git a/crates/paimon/src/io/file_io.rs b/crates/paimon/src/io/file_io.rs index 720225ebc..a7a08273e 100644 --- a/crates/paimon/src/io/file_io.rs +++ b/crates/paimon/src/io/file_io.rs @@ -20,6 +20,7 @@ use std::collections::HashMap; use std::future::Future; use std::ops::Range; use std::pin::Pin; +use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; use std::task::{Context, Poll}; use std::time::SystemTime; @@ -90,6 +91,13 @@ enum FileIOBackend { pub struct FileIO { backend: FileIOBackend, cache: Option>, + cache_namespace: usize, +} + +static NEXT_FILE_IO_CACHE_NAMESPACE: AtomicUsize = AtomicUsize::new(1); + +fn next_file_io_cache_namespace() -> usize { + NEXT_FILE_IO_CACHE_NAMESPACE.fetch_add(1, Ordering::Relaxed) } impl std::fmt::Debug for FileIO { @@ -126,6 +134,7 @@ impl FileIO { /// subsequently created by [`Self::new_input`] and [`Self::new_output`]. pub fn with_provider(mut self, provider: Arc) -> Self { self.backend = FileIOBackend::Provider(provider); + self.cache_namespace = next_file_io_cache_namespace(); self } @@ -211,6 +220,7 @@ impl FileIO { Ok(InputFile { source: self.file_source(path)?, path: path.to_string(), + cache_namespace: self.cache_namespace, cache: self .cache .as_ref() @@ -227,6 +237,7 @@ impl FileIO { Ok(OutputFile { source: self.file_source(path)?, path: path.to_string(), + cache_namespace: self.cache_namespace, cache: self .cache .as_ref() @@ -651,13 +662,22 @@ impl FileIOBuilder { } else { FileIOBackend::Storage(Arc::new(Storage::build(self)?)) }; - Ok(FileIO { backend, cache }) + Ok(FileIO { + backend, + cache, + cache_namespace: next_file_io_cache_namespace(), + }) } } #[async_trait::async_trait] pub trait FileRead: Send + Sync + Unpin + 'static { async fn read(&self, range: Range) -> crate::Result; + + #[doc(hidden)] + fn cache_namespace(&self) -> Option { + None + } } #[async_trait::async_trait] @@ -668,18 +688,24 @@ impl FileRead for opendal::Reader { } enum InputFileReader { - Direct(opendal::Reader), - Cached(CachedFileReader), + Direct(opendal::Reader, usize), + Cached(CachedFileReader, usize), } #[async_trait::async_trait] impl FileRead for InputFileReader { async fn read(&self, range: Range) -> crate::Result { match self { - Self::Direct(reader) => FileRead::read(reader, range).await, - Self::Cached(reader) => FileRead::read(reader, range).await, + Self::Direct(reader, _) => FileRead::read(reader, range).await, + Self::Cached(reader, _) => FileRead::read(reader, range).await, } } + + fn cache_namespace(&self) -> Option { + Some(match self { + Self::Direct(_, namespace) | Self::Cached(_, namespace) => *namespace, + }) + } } #[async_trait::async_trait] @@ -816,6 +842,7 @@ impl FileSource { pub struct InputFile { source: FileSource, path: String, + cache_namespace: usize, cache: Option>, } @@ -866,7 +893,7 @@ impl InputFile { let (op, relative_path, cache_path) = self.source.resolve(&self.path).await?; let reader = op.reader(&relative_path).await?; let Some(cache) = &self.cache else { - return Ok(InputFileReader::Direct(reader)); + return Ok(InputFileReader::Direct(reader, self.cache_namespace)); }; let read_token = cache.read_token(&cache_path); let size = if let Some(size) = cache.file_size(&cache_path, &read_token).await { @@ -876,13 +903,16 @@ impl InputFile { cache.put_file_size(&cache_path, size, &read_token).await; size }; - Ok(InputFileReader::Cached(CachedFileReader::new_with_token( - Arc::new(reader), - &cache_path, - size, - cache.clone(), - read_token, - ))) + Ok(InputFileReader::Cached( + CachedFileReader::new_with_token( + Arc::new(reader), + &cache_path, + size, + cache.clone(), + read_token, + ), + self.cache_namespace, + )) } } @@ -890,6 +920,7 @@ impl InputFile { pub struct OutputFile { source: FileSource, path: String, + cache_namespace: usize, cache: Option>, } @@ -908,6 +939,7 @@ impl OutputFile { InputFile { source: self.source, path: self.path, + cache_namespace: self.cache_namespace, cache, } } diff --git a/crates/paimon/src/table/data_file_reader.rs b/crates/paimon/src/table/data_file_reader.rs index cf429efc5..3f1940f32 100644 --- a/crates/paimon/src/table/data_file_reader.rs +++ b/crates/paimon/src/table/data_file_reader.rs @@ -105,6 +105,10 @@ impl FileRead for TimedFileRead { self.timing.add_file_read(start.elapsed()); result } + + fn cache_namespace(&self) -> Option { + self.inner.cache_namespace() + } } /// Reads data from Parquet files. @@ -1748,6 +1752,21 @@ mod tests { ); } + #[tokio::test] + async fn timed_file_read_preserves_cache_namespace() { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + let input = file_io.new_input("memory:/timed-file").unwrap(); + let reader = input.reader().await.unwrap(); + let cache_namespace = reader.cache_namespace(); + let timed = TimedFileRead { + inner: Box::new(reader), + timing: Arc::new(DataFileReadTiming::default()), + }; + + assert!(cache_namespace.is_some()); + assert_eq!(timed.cache_namespace(), cache_namespace); + } + #[test] fn merge_row_selection_skips_only_unfiltered_full_coverage() { let full = [RowRange::new(0, 9)];