This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new b417ea83 feat(blob): add standalone BlobDescriptor reader (#762)
b417ea83 is described below

commit b417ea83af077253065bea5eb09e08b3cdeccde6
Author: XiaoHongbo <[email protected]>
AuthorDate: Sun Aug 30 21:01:30 2026 +0800

    feat(blob): add standalone BlobDescriptor reader (#762)
---
 .../paimon/src/catalog/rest/rest_token_file_io.rs  |  30 ++
 crates/paimon/src/lib.rs                           |   6 +-
 crates/paimon/src/table/blob_resolver.rs           | 331 +++++++++++++++++++++
 crates/paimon/src/table/mod.rs                     |   1 +
 4 files changed, 365 insertions(+), 3 deletions(-)

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<Arc<AtomicUsize>>) -> 
Json<GetTableTokenResponse> {
         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<String, String>,
+    file_io: Option<FileIO>,
+}
+
+impl BlobReader {
+    pub fn new(storage_options: HashMap<String, String>) -> 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<u8>]) -> 
Result<Vec<Vec<u8>>> {
+        let mut by_uri = HashMap::<String, Vec<(usize, 
BlobDescriptor)>>::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<Vec<(usize, Vec<u8>)>> = stream::iter(by_uri)
+            .map(|(uri, entries)| {
+                let limiter = limiter.clone();
+                async move {
+                    let indices = entries.iter().map(|(index, _)| 
*index).collect::<Vec<_>>();
+                    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::<Vec<_>>();
+        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<std::sync::atomic::AtomicUsize>,
         max_in_flight: std::sync::Arc<std::sync::atomic::AtomicUsize>,
+        ranges: std::sync::Arc<std::sync::Mutex<Vec<std::ops::Range<u64>>>>,
     }
 
     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<std::ops::Range<u64>> {
+            self.ranges.lock().unwrap().clone()
+        }
     }
 
     #[async_trait::async_trait]
     impl FileRead for TrackingFileRead {
         async fn read(&self, range: std::ops::Range<u64>) -> 
crate::Result<Bytes> {
+            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<u64>) -> 
crate::Result<Bytes> {
+            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<u8> {
+        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<u8> {
+        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:[email protected]/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};

Reply via email to