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 7215a561 feat(blob): add seekable descriptor streams (#763)
7215a561 is described below
commit 7215a561d1c9809b320a0e4d598dddd66a2097bd
Author: XiaoHongbo <[email protected]>
AuthorDate: Mon Aug 31 12:58:32 2026 +0800
feat(blob): add seekable descriptor streams (#763)
---
.../paimon/src/catalog/rest/rest_token_file_io.rs | 28 ++++
crates/paimon/src/lib.rs | 12 +-
crates/paimon/src/table/blob_resolver.rs | 183 ++++++++++++++++++++-
crates/paimon/src/table/mod.rs | 2 +-
4 files changed, 215 insertions(+), 10 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 22480799..d7d17f78 100644
--- a/crates/paimon/src/catalog/rest/rest_token_file_io.rs
+++ b/crates/paimon/src/catalog/rest/rest_token_file_io.rs
@@ -322,4 +322,32 @@ mod tests {
assert_eq!(requests.load(Ordering::SeqCst), 2);
server.abort();
}
+
+ #[tokio::test]
+ async fn test_blob_stream_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 mut stream = BlobReader::from_file_io(file_io)
+ .open_blob(&descriptor)
+ .unwrap();
+
+ assert_eq!(stream.read(2).await.unwrap(), b"cd");
+ assert_eq!(stream.read(2).await.unwrap(), b"ef");
+ assert_eq!(requests.load(Ordering::SeqCst), 2);
+ server.abort();
+ }
}
diff --git a/crates/paimon/src/lib.rs b/crates/paimon/src/lib.rs
index fbb56333..5a33f65c 100644
--- a/crates/paimon/src/lib.rs
+++ b/crates/paimon/src/lib.rs
@@ -48,12 +48,12 @@ pub use catalog::CatalogFactory;
pub use catalog::FileSystemCatalog;
pub use table::{
- 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,
+ BlobReader, BlobStream, 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,
};
pub use table::{
diff --git a/crates/paimon/src/table/blob_resolver.rs
b/crates/paimon/src/table/blob_resolver.rs
index f5029ba1..22cee1b0 100644
--- a/crates/paimon/src/table/blob_resolver.rs
+++ b/crates/paimon/src/table/blob_resolver.rs
@@ -23,6 +23,7 @@ use arrow_array::{Array, BinaryArray};
use bytes::Bytes;
use futures::{stream, StreamExt, TryStreamExt};
use std::collections::HashMap;
+use std::io::SeekFrom;
use std::sync::Arc;
use tokio::sync::{OwnedSemaphorePermit, Semaphore};
@@ -33,10 +34,17 @@ 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)]
+#[derive(Clone, Debug)]
pub struct BlobReader {
storage_options: HashMap<String, String>,
file_io: Option<FileIO>,
+ limiter: BlobReadLimiter,
+}
+
+impl Default for BlobReader {
+ fn default() -> Self {
+ Self::new(HashMap::new())
+ }
}
impl BlobReader {
@@ -44,6 +52,7 @@ impl BlobReader {
Self {
storage_options,
file_io: None,
+ limiter: BlobReadLimiter::new(),
}
}
@@ -52,9 +61,34 @@ impl BlobReader {
Self {
storage_options: HashMap::new(),
file_io: Some(file_io),
+ limiter: BlobReadLimiter::new(),
}
}
+ /// Open one descriptor for incremental reads.
+ pub fn open_blob(&self, bytes: &[u8]) -> Result<BlobStream> {
+ let descriptor = BlobDescriptor::deserialize(bytes)
+ .map_err(|error| blob_error_with_context(error, &[0], None))?;
+ let range = descriptor
+ .range_spec()
+ .map_err(|error| blob_error_with_context(error, &[0],
Some(descriptor.uri())))?;
+ let file_io = match &self.file_io {
+ Some(file_io) => file_io.clone(),
+ None => FileIO::from_path(descriptor.uri())
+ .and_then(|builder|
builder.with_props(self.storage_options.iter()).build())
+ .map_err(|error| blob_error_with_context(error, &[0],
Some(descriptor.uri())))?,
+ };
+
+ Ok(BlobStream {
+ file_io,
+ uri: descriptor.uri().to_string(),
+ offset: range.offset(),
+ length: range.length(),
+ position: 0,
+ limiter: self.limiter.clone(),
+ })
+ }
+
/// 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();
@@ -70,7 +104,7 @@ impl BlobReader {
.push((index, descriptor));
}
- let limiter = BlobReadLimiter::new();
+ let limiter = self.limiter.clone();
let groups: Vec<Vec<(usize, Vec<u8>)>> = stream::iter(by_uri)
.map(|(uri, entries)| {
let limiter = limiter.clone();
@@ -121,6 +155,98 @@ impl BlobReader {
}
}
+/// Incremental reader for one serialized [`BlobDescriptor`].
+#[derive(Debug)]
+pub struct BlobStream {
+ file_io: FileIO,
+ uri: String,
+ offset: u64,
+ length: Option<u64>,
+ position: u64,
+ limiter: BlobReadLimiter,
+}
+
+impl BlobStream {
+ /// Read at most `max_bytes`, returning an empty buffer at end of stream.
+ pub async fn read(&mut self, max_bytes: usize) -> Result<Vec<u8>> {
+ self.read_inner(max_bytes)
+ .await
+ .map_err(|error| blob_error_with_context(error, &[0],
Some(&self.uri)))
+ }
+
+ /// Seek within the descriptor range.
+ pub async fn seek(&mut self, from: SeekFrom) -> Result<u64> {
+ self.seek_inner(from)
+ .await
+ .map_err(|error| blob_error_with_context(error, &[0],
Some(&self.uri)))
+ }
+
+ async fn seek_inner(&mut self, from: SeekFrom) -> Result<u64> {
+ let position = match from {
+ SeekFrom::Start(position) => i128::from(position),
+ SeekFrom::Current(offset) => i128::from(self.position) +
i128::from(offset),
+ SeekFrom::End(offset) => i128::from(self.length().await?) +
i128::from(offset),
+ };
+ self.position = u64::try_from(position).map_err(|_|
crate::Error::DataInvalid {
+ message: "invalid BlobDescriptor stream seek".to_string(),
+ source: None,
+ })?;
+ Ok(self.position)
+ }
+
+ async fn length(&mut self) -> Result<u64> {
+ if let Some(length) = self.length {
+ return Ok(length);
+ }
+ let input = self.file_io.new_input(&self.uri)?;
+ let _permit = self.limiter.acquire_request(&self.uri,
"metadata").await?;
+ let length = input.metadata().await?.size.saturating_sub(self.offset);
+ self.length = Some(length);
+ Ok(length)
+ }
+
+ async fn read_inner(&mut self, max_bytes: usize) -> Result<Vec<u8>> {
+ if max_bytes == 0 {
+ return Ok(Vec::new());
+ }
+
+ let remaining = self.length().await?.saturating_sub(self.position);
+ if remaining == 0 {
+ return Ok(Vec::new());
+ }
+
+ let length =
remaining.min(u64::try_from(max_bytes).unwrap_or(u64::MAX));
+ let start =
+ self.offset
+ .checked_add(self.position)
+ .ok_or_else(|| crate::Error::DataInvalid {
+ message: "BlobDescriptor stream position overflows
u64".to_string(),
+ source: None,
+ })?;
+ let end = start
+ .checked_add(length)
+ .ok_or_else(|| crate::Error::DataInvalid {
+ message: "BlobDescriptor stream range overflows
u64".to_string(),
+ source: None,
+ })?;
+ let input = self.file_io.new_input(&self.uri)?;
+ let reader = input.reader().await?;
+ let _permits = self.limiter.acquire_read(length, &self.uri).await?;
+ let bytes = reader.read(start..end).await?;
+ if bytes.len() as u64 != length {
+ return Err(crate::Error::DataInvalid {
+ message: format!(
+ "short read for range {start}..{end}, expected={length}
bytes, actual={} bytes",
+ bytes.len()
+ ),
+ source: None,
+ });
+ }
+ self.position += length;
+ Ok(bytes.to_vec())
+ }
+}
+
fn blob_error_with_context(
error: crate::Error,
indices: &[usize],
@@ -190,7 +316,7 @@ fn sanitize_blob_uri(uri: &str) -> String {
/// The byte semaphore budgets active range I/O only. A single range larger
than
/// the budget consumes every byte permit and runs alone, but can still
allocate
/// more than the configured budget because the complete value is required.
-#[derive(Clone)]
+#[derive(Clone, Debug)]
pub(crate) struct BlobReadLimiter {
requests: Arc<Semaphore>,
bytes: Arc<Semaphore>,
@@ -945,4 +1071,55 @@ mod tests {
.unwrap()
.is_empty());
}
+
+ #[tokio::test]
+ async fn test_blob_stream_reads_and_seeks() {
+ let file = tempfile::NamedTempFile::new().unwrap();
+ std::fs::write(file.path(), b"abcdefghij").unwrap();
+ let uri = file_uri(file.path());
+ let reader = BlobReader::default();
+
+ let mut fixed = reader.open_blob(&java_v2_descriptor(&uri, 2,
5)).unwrap();
+ assert_eq!(fixed.read(2).await.unwrap(), b"cd");
+ assert_eq!(fixed.seek(SeekFrom::Start(1)).await.unwrap(), 1);
+ assert_eq!(fixed.read(2).await.unwrap(), b"de");
+ assert_eq!(fixed.seek(SeekFrom::Current(-1)).await.unwrap(), 2);
+ assert_eq!(fixed.read(2).await.unwrap(), b"ef");
+ assert_eq!(fixed.seek(SeekFrom::End(-2)).await.unwrap(), 3);
+ assert_eq!(fixed.read(8).await.unwrap(), b"fg");
+ assert!(fixed.read(1).await.unwrap().is_empty());
+ assert!(fixed.seek(SeekFrom::Current(-6)).await.is_err());
+
+ let mut to_end = reader.open_blob(&java_v1_descriptor(&uri, 4,
-1)).unwrap();
+ assert_eq!(to_end.seek(SeekFrom::End(-3)).await.unwrap(), 3);
+ assert_eq!(to_end.read(3).await.unwrap(), b"hij");
+ assert_eq!(to_end.seek(SeekFrom::Start(0)).await.unwrap(), 0);
+ assert_eq!(to_end.read(8).await.unwrap(), b"efghij");
+
+ let mut empty = reader.open_blob(&java_v2_descriptor(&uri, 3,
0)).unwrap();
+ assert!(empty.read(1).await.unwrap().is_empty());
+
+ let mut short = reader.open_blob(&java_v2_descriptor(&uri, 8,
4)).unwrap();
+ assert!(short.read(4).await.is_err());
+ }
+
+ #[tokio::test]
+ async fn test_blob_stream_is_lazy_and_validates_input() {
+ let reader = BlobReader::default();
+ let directory = tempfile::tempdir().unwrap();
+ let missing = file_uri(&directory.path().join("missing"));
+ let mut stream = reader
+ .open_blob(&java_v2_descriptor(&missing, 0, -1))
+ .unwrap();
+
+ assert!(stream.read(0).await.unwrap().is_empty());
+ let error = stream.read(1).await.unwrap_err().to_string();
+ assert!(error.contains("input indices [0]"));
+ assert!(error.contains("object not found") || error.contains("storage
I/O failed"));
+
+ assert!(reader.open_blob(&[]).is_err());
+ assert!(reader
+ .open_blob(&BlobDescriptor::new("file:///tmp/a".to_string(), -1,
1).serialize())
+ .is_err());
+ }
}
diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs
index 6f56b58f..3b0f447d 100644
--- a/crates/paimon/src/table/mod.rs
+++ b/crates/paimon/src/table/mod.rs
@@ -109,7 +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 blob_resolver::{BlobReader, BlobStream};
pub use branch_manager::BranchManager;
pub use commit_message::CommitMessage;
pub use cow_writer::{CopyOnWriteMergeWriter, FileInfo};