diff --git a/crates/paimon/src/catalog/rest/rest_token_file_io.rs b/crates/paimon/src/catalog/rest/rest_token_file_io.rs index c34967a9..22480799 100644 --- a/crates/paimon/src/catalog/rest/rest_token_file_io.rs +++ b/crates/paimon/src/catalog/rest/rest_token_file_io.rs @@ -185,6 +185,8 @@ mod tests { use super::*; use crate::api::GetTableTokenResponse; use crate::io::cache::create_local_cache; + use crate::spec::BlobDescriptor; + use crate::BlobReader; async fn token(State(requests): State>) -> Json { let request = requests.fetch_add(1, Ordering::SeqCst); @@ -292,4 +294,32 @@ mod tests { assert_eq!(requests.load(Ordering::SeqCst), 2); server.abort(); } + + #[tokio::test] + async fn test_blob_reader_reuses_refreshing_file_io() { + let table_directory = tempfile::tempdir().unwrap(); + let file_path = table_directory.path().join("blob"); + std::fs::write(&file_path, b"abcdefghij").unwrap(); + let (options, api, requests, server) = token_api().await; + let token_file_io = Arc::new(RESTTokenFileIO::new( + Identifier::new("database", "table"), + table_directory.path().to_string_lossy().into_owned(), + options, + api, + None, + )); + + let file_io = token_file_io.build_file_io().await.unwrap(); + assert_eq!(requests.load(Ordering::SeqCst), 1); + let uri = url::Url::from_file_path(file_path).unwrap().to_string(); + let descriptor = BlobDescriptor::new(uri, 2, 4).serialize(); + let values = BlobReader::from_file_io(file_io) + .read_blobs(&[descriptor]) + .await + .unwrap(); + + assert_eq!(values, vec![b"cdef".to_vec()]); + assert_eq!(requests.load(Ordering::SeqCst), 2); + server.abort(); + } } diff --git a/crates/paimon/src/lib.rs b/crates/paimon/src/lib.rs index a86a95e9..fbb56333 100644 --- a/crates/paimon/src/lib.rs +++ b/crates/paimon/src/lib.rs @@ -48,9 +48,9 @@ pub use catalog::CatalogFactory; pub use catalog::FileSystemCatalog; pub use table::{ - CommitMessage, DataEvolutionDeleteWriter, DataEvolutionWriter, DataSplit, DataSplitBuilder, - DeletionFile, IncrementalPlan, IncrementalScan, IncrementalScanMode, IncrementalSplit, - PartitionBucket, Plan, PostponeBucketPlan, PostponeFixedBucketTableCommit, + BlobReader, CommitMessage, DataEvolutionDeleteWriter, DataEvolutionWriter, DataSplit, + DataSplitBuilder, DeletionFile, IncrementalPlan, IncrementalScan, IncrementalScanMode, + IncrementalSplit, PartitionBucket, Plan, PostponeBucketPlan, PostponeFixedBucketTableCommit, PostponeFixedBucketTableWrite, RESTEnv, RESTSnapshotCommit, ReadBuilder, RenamingSnapshotCommit, RowRange, ScanTrace, SnapshotCommit, SnapshotManager, Table, TableCommit, TableRead, TableScan, TableUpdate, TableWrite, TagManager, WriteBuilder, diff --git a/crates/paimon/src/table/blob_resolver.rs b/crates/paimon/src/table/blob_resolver.rs index 673efb9f..f5029ba1 100644 --- a/crates/paimon/src/table/blob_resolver.rs +++ b/crates/paimon/src/table/blob_resolver.rs @@ -32,6 +32,159 @@ pub(crate) const BLOB_DESCRIPTOR_READ_CONCURRENCY: usize = 8; const BLOB_DESCRIPTOR_READ_BYTE_UNIT: u64 = 1024 * 1024; const BLOB_DESCRIPTOR_READ_MAX_IN_FLIGHT_BYTES: u64 = 64 * 1024 * 1024; +/// Reads serialized [`BlobDescriptor`] values without requiring a table. +#[derive(Clone, Debug, Default)] +pub struct BlobReader { + storage_options: HashMap, + file_io: Option, +} + +impl BlobReader { + pub fn new(storage_options: HashMap) -> Self { + Self { + storage_options, + file_io: None, + } + } + + /// Create a reader that reuses an existing FileIO. + pub fn from_file_io(file_io: FileIO) -> Self { + Self { + storage_options: HashMap::new(), + file_io: Some(file_io), + } + } + + /// Read a descriptor batch in input order. + pub async fn read_blobs(&self, descriptors: &[Vec]) -> Result>> { + let mut by_uri = HashMap::>::new(); + for (index, bytes) in descriptors.iter().enumerate() { + let descriptor = BlobDescriptor::deserialize(bytes) + .map_err(|error| blob_error_with_context(error, &[index], None))?; + descriptor.range_spec().map_err(|error| { + blob_error_with_context(error, &[index], Some(descriptor.uri())) + })?; + by_uri + .entry(descriptor.uri().to_string()) + .or_default() + .push((index, descriptor)); + } + + let limiter = BlobReadLimiter::new(); + let groups: Vec)>> = stream::iter(by_uri) + .map(|(uri, entries)| { + let limiter = limiter.clone(); + async move { + let indices = entries.iter().map(|(index, _)| *index).collect::>(); + let file_io = match &self.file_io { + Some(file_io) => file_io.clone(), + None => FileIO::from_path(&uri) + .and_then(|builder| { + builder.with_props(self.storage_options.iter()).build() + }) + .map_err(|error| { + blob_error_with_context(error, &indices, Some(&uri)) + })?, + }; + let mut builder = BinaryBuilder::with_capacity(entries.len(), 0); + for (_, descriptor) in &entries { + builder.append_value( + BlobDescriptor::new( + descriptor.uri().to_string(), + descriptor.offset(), + descriptor.length(), + ) + .serialize(), + ); + } + let resolved = resolve_blob_column(&builder.finish(), &file_io, limiter) + .await + .map_err(|error| blob_error_with_context(error, &indices, Some(&uri)))?; + Ok::<_, crate::Error>( + entries + .into_iter() + .enumerate() + .map(|(position, (index, _))| { + (index, resolved.value(position).to_vec()) + }) + .collect(), + ) + } + }) + .buffer_unordered(BLOB_DESCRIPTOR_READ_CONCURRENCY) + .try_collect() + .await?; + + let mut values = groups.into_iter().flatten().collect::>(); + values.sort_unstable_by_key(|(index, _)| *index); + Ok(values.into_iter().map(|(_, value)| value).collect()) + } +} + +fn blob_error_with_context( + error: crate::Error, + indices: &[usize], + uri: Option<&str>, +) -> crate::Error { + let location = match uri { + Some(uri) => format!( + "input indices {indices:?}, URI '{}'", + sanitize_blob_uri(uri) + ), + None => format!("input indices {indices:?}, URI unavailable"), + }; + let sanitize = |message: String| match uri { + Some(uri) => message.replace(uri, &sanitize_blob_uri(uri)), + None => message, + }; + match error { + crate::Error::Unsupported { message } => crate::Error::Unsupported { + message: format!("BlobDescriptor {location}: {}", sanitize(message)), + }, + crate::Error::IoUnsupported { message } => crate::Error::IoUnsupported { + message: format!("BlobDescriptor {location}: {}", sanitize(message)), + }, + crate::Error::ConfigInvalid { .. } => crate::Error::ConfigInvalid { + message: format!("BlobDescriptor {location}: invalid storage URI or options"), + }, + crate::Error::DataInvalid { message, .. } => crate::Error::DataInvalid { + message: format!("BlobDescriptor {location}: {}", sanitize(message)), + source: None, + }, + crate::Error::IoUnexpected { source, .. } + if source.kind() == opendal::ErrorKind::NotFound => + { + crate::Error::UnexpectedError { + message: format!("BlobDescriptor {location}: object not found"), + source: None, + } + } + crate::Error::IoUnexpected { .. } => crate::Error::UnexpectedError { + message: format!("BlobDescriptor {location}: storage I/O failed"), + source: None, + }, + crate::Error::UnexpectedError { message, .. } => crate::Error::UnexpectedError { + message: format!("BlobDescriptor {location}: {}", sanitize(message)), + source: None, + }, + _ => crate::Error::UnexpectedError { + message: format!("BlobDescriptor {location}: operation failed"), + source: None, + }, + } +} + +fn sanitize_blob_uri(uri: &str) -> String { + if let Ok(mut url) = url::Url::parse(uri) { + let _ = url.set_username(""); + let _ = url.set_password(None); + url.set_query(None); + url.set_fragment(None); + return url.to_string(); + } + uri.split(['?', '#']).next().unwrap_or(uri).to_string() +} + /// Shared admission control for external descriptor metadata and range reads. /// /// The byte semaphore budgets active range I/O only. A single range larger than @@ -399,6 +552,7 @@ mod tests { bytes: Bytes, in_flight: std::sync::Arc, max_in_flight: std::sync::Arc, + ranges: std::sync::Arc>>>, } impl TrackingFileRead { @@ -407,6 +561,7 @@ mod tests { bytes, in_flight: std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)), max_in_flight: std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)), + ranges: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())), } } @@ -419,17 +574,23 @@ mod tests { bytes, in_flight, max_in_flight, + ranges: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())), } } fn max_in_flight(&self) -> usize { self.max_in_flight.load(std::sync::atomic::Ordering::SeqCst) } + + fn ranges(&self) -> Vec> { + self.ranges.lock().unwrap().clone() + } } #[async_trait::async_trait] impl FileRead for TrackingFileRead { async fn read(&self, range: std::ops::Range) -> crate::Result { + self.ranges.lock().unwrap().push(range.clone()); let in_flight = self .in_flight .fetch_add(1, std::sync::atomic::Ordering::SeqCst) @@ -443,6 +604,15 @@ mod tests { } } + struct ShortFileRead; + + #[async_trait::async_trait] + impl FileRead for ShortFileRead { + async fn read(&self, _range: std::ops::Range) -> crate::Result { + Ok(Bytes::from_static(b"x")) + } + } + #[tokio::test] async fn test_blob_range_reads_use_bounded_parallelism() { let reader = TrackingFileRead::new(Bytes::from_static(b"abcdefghijkl")); @@ -575,6 +745,63 @@ mod tests { .unwrap(); } + #[tokio::test] + async fn test_merged_descriptors_issue_one_underlying_read() { + let reader = TrackingFileRead::new(Bytes::from_static(b"abcdefghijkl")); + let reads = merge_blob_read_requests(vec![ + BlobReadRequest { + row: 0, + offset: 0, + length: 4, + }, + BlobReadRequest { + row: 1, + offset: 4, + length: 4, + }, + BlobReadRequest { + row: 2, + offset: 2, + length: 6, + }, + ]); + + let results = read_merged_blob_ranges( + "memory:/blob.bin", + Arc::new(reader.clone()), + reads, + BlobReadLimiter::new(), + ) + .await + .unwrap(); + + assert_eq!(results.len(), 1); + assert_eq!(reader.ranges(), vec![0..8]); + } + + #[tokio::test] + async fn test_blob_range_read_rejects_short_data() { + let error = read_merged_blob_ranges( + "memory:/blob.bin", + Arc::new(ShortFileRead), + vec![MergedBlobRead { + start: 4, + end: 8, + requests: vec![BlobReadRequest { + row: 0, + offset: 4, + length: 4, + }], + }], + BlobReadLimiter::new(), + ) + .await + .err() + .expect("short read must fail"); + + assert!(error.to_string().contains("short read")); + } + #[test] fn test_merge_blob_read_requests_merges_nearby_ranges() { let merged = merge_blob_read_requests(vec![ @@ -614,4 +841,108 @@ mod tests { assert_eq!(merged[1].start, BLOB_RANGE_MERGE_MAX_SPAN + 1); assert_eq!(merged[1].end, BLOB_RANGE_MERGE_MAX_SPAN + 5); } + + fn java_v2_descriptor(uri: &str, offset: i64, length: i64) -> Vec { + let mut bytes = Vec::new(); + bytes.push(2); + bytes.extend_from_slice(&0x424C4F4244455343_u64.to_le_bytes()); + bytes.extend_from_slice(&(uri.len() as i32).to_le_bytes()); + bytes.extend_from_slice(uri.as_bytes()); + bytes.extend_from_slice(&offset.to_le_bytes()); + bytes.extend_from_slice(&length.to_le_bytes()); + bytes + } + + fn java_v1_descriptor(uri: &str, offset: i64, length: i64) -> Vec { + let mut bytes = vec![1]; + bytes.extend_from_slice(&(uri.len() as i32).to_le_bytes()); + bytes.extend_from_slice(uri.as_bytes()); + bytes.extend_from_slice(&offset.to_le_bytes()); + bytes.extend_from_slice(&length.to_le_bytes()); + bytes + } + + fn file_uri(path: &std::path::Path) -> String { + url::Url::from_file_path(path).unwrap().to_string() + } + + #[tokio::test] + async fn test_standalone_blob_reader_reads_ranges_in_input_order() { + let first = tempfile::NamedTempFile::new().unwrap(); + let second = tempfile::NamedTempFile::new().unwrap(); + std::fs::write(first.path(), b"abcdefghij").unwrap(); + std::fs::write(second.path(), b"UVWXYZ").unwrap(); + let first_uri = file_uri(first.path()); + let second_uri = file_uri(second.path()); + let descriptors = vec![ + java_v2_descriptor(&second_uri, 1, 3), + java_v2_descriptor(&first_uri, 3, -1), + java_v2_descriptor(&first_uri, 5, 0), + java_v2_descriptor(&first_uri, 2, 4), + java_v2_descriptor(&first_uri, 2, 4), + ]; + + let values = BlobReader::default() + .read_blobs(&descriptors) + .await + .unwrap(); + + assert_eq!( + values, + vec![ + b"VWX".to_vec(), + b"defghij".to_vec(), + Vec::new(), + b"cdef".to_vec(), + b"cdef".to_vec(), + ] + ); + } + + #[tokio::test] + async fn test_standalone_blob_reader_reads_java_v1_descriptor() { + let file = tempfile::NamedTempFile::new().unwrap(); + std::fs::write(file.path(), b"abcdefghij").unwrap(); + let descriptor = java_v1_descriptor(&file_uri(file.path()), 2, 4); + + let values = BlobReader::default() + .read_blobs(&[descriptor]) + .await + .unwrap(); + + assert_eq!(values, vec![b"cdef".to_vec()]); + } + + #[tokio::test] + async fn test_standalone_blob_reader_validates_input() { + let reader = BlobReader::default(); + + let error = reader.read_blobs(&[Vec::new()]).await.unwrap_err(); + assert!(error.to_string().contains("input indices [0]")); + + let error = reader + .read_blobs(&[BlobDescriptor::new("file:///tmp/a".to_string(), -1, 1).serialize()]) + .await + .unwrap_err(); + assert!(error.to_string().contains("offset must be non-negative")); + + let secret_uri = "ftp://access-key:secret@example.com/a?token=sensitive"; + let error = reader + .read_blobs(&[BlobDescriptor::new(secret_uri.to_string(), 0, 1).serialize()]) + .await + .unwrap_err(); + let message = error.to_string(); + assert!(message.contains("ftp://example.com/a")); + assert!(!message.contains("access-key")); + assert!(!message.contains("sensitive")); + } + + #[tokio::test] + async fn test_standalone_blob_reader_empty_batch() { + assert!(BlobReader::default() + .read_blobs(&[]) + .await + .unwrap() + .is_empty()); + } } diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs index 35f1e45a..6f56b58f 100644 --- a/crates/paimon/src/table/mod.rs +++ b/crates/paimon/src/table/mod.rs @@ -109,6 +109,7 @@ mod write_builder; use crate::Result; use arrow_array::RecordBatch; pub use audit_log_table::AuditLogTable; +pub use blob_resolver::BlobReader; pub use branch_manager::BranchManager; pub use commit_message::CommitMessage; pub use cow_writer::{CopyOnWriteMergeWriter, FileInfo};