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};