Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 28 additions & 0 deletions crates/paimon/src/catalog/rest/rest_token_file_io.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}
}
12 changes: 6 additions & 6 deletions crates/paimon/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand Down
183 changes: 180 additions & 3 deletions crates/paimon/src/table/blob_resolver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};

Expand All @@ -33,17 +34,25 @@ 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 {
pub fn new(storage_options: HashMap<String, String>) -> Self {
Self {
storage_options,
file_io: None,
limiter: BlobReadLimiter::new(),
}
}

Expand All @@ -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();
Expand All @@ -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();
Expand Down Expand Up @@ -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],
Expand Down Expand Up @@ -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>,
Expand Down Expand Up @@ -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());
}
}
2 changes: 1 addition & 1 deletion crates/paimon/src/table/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down
Loading