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 7ceb2f8  fix(blob): support unknown-length BlobDescriptor ranges (#521)
7ceb2f8 is described below

commit 7ceb2f882446a67173222477823b3b19361da525
Author: QuakeWang <[email protected]>
AuthorDate: Wed Jul 15 13:08:06 2026 +0800

    fix(blob): support unknown-length BlobDescriptor ranges (#521)
---
 crates/integrations/datafusion/tests/blob_tests.rs | 179 ++++++++++++++++++++-
 crates/paimon/src/arrow/format/blob.rs             | 116 +++++++++++--
 crates/paimon/src/spec/blob_descriptor.rs          | 112 +++++++++++++
 crates/paimon/src/table/blob_resolver.rs           | 141 +++++++++++-----
 docs/src/sql.md                                    |   4 +
 5 files changed, 500 insertions(+), 52 deletions(-)

diff --git a/crates/integrations/datafusion/tests/blob_tests.rs 
b/crates/integrations/datafusion/tests/blob_tests.rs
index 84fad62..b844f65 100644
--- a/crates/integrations/datafusion/tests/blob_tests.rs
+++ b/crates/integrations/datafusion/tests/blob_tests.rs
@@ -22,7 +22,7 @@
 mod common;
 
 use arrow_array::{Array, BinaryArray, Int32Array, RecordBatch, StringArray};
-use common::{create_sql_context, create_test_env, exec};
+use common::{assert_sql_error, create_sql_context, create_test_env, exec};
 use paimon::catalog::Identifier;
 use paimon::spec::{BlobDescriptor, BlobViewStruct};
 use paimon::table::BranchManager;
@@ -577,6 +577,66 @@ async fn test_blob_resolve_descriptor_with_offset() {
     assert_eq!(rows[0], (1, "Partial".into(), Some(b"PAYLOAD".to_vec())));
 }
 
+#[tokio::test]
+async fn test_blob_resolve_unknown_length_descriptor() {
+    let (tmp, sql_context) = setup(BLOB_TABLE_DDL).await;
+
+    let source_data = b"HEADER_PAYLOAD_TRAILER";
+    let source_path = tmp.path().join("blob_unknown_length.bin");
+    std::fs::write(&source_path, source_data).unwrap();
+
+    let uri = format!("file://{}", source_path.display());
+    let full_hex = to_hex(&BlobDescriptor::new(uri.clone(), 0, 
-1).serialize());
+    let suffix_hex = to_hex(&BlobDescriptor::new(uri.clone(), 7, 
-1).serialize());
+    let eof_hex =
+        to_hex(&BlobDescriptor::new(uri.clone(), source_data.len() as i64, 
-1).serialize());
+    let past_eof_hex =
+        to_hex(&BlobDescriptor::new(uri, source_data.len() as i64 + 5, 
-1).serialize());
+
+    let sql = format!(
+        "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \
+         (1, 'Full', X'{full_hex}'), \
+         (2, 'Suffix', X'{suffix_hex}'), \
+         (3, 'Eof', X'{eof_hex}'), \
+         (4, 'PastEof', X'{past_eof_hex}'), \
+         (5, 'Raw', X'524157'), \
+         (6, 'Null', NULL)"
+    );
+    exec(&sql_context, &sql).await;
+
+    let rows = query_id_name_picture(
+        &sql_context,
+        "SELECT id, name, picture FROM paimon.test_db.t ORDER BY id",
+    )
+    .await;
+    assert_eq!(
+        rows,
+        vec![
+            (1, "Full".into(), Some(source_data.to_vec())),
+            (2, "Suffix".into(), Some(b"PAYLOAD_TRAILER".to_vec())),
+            (3, "Eof".into(), Some(Vec::new())),
+            (4, "PastEof".into(), Some(Vec::new())),
+            (5, "Raw".into(), Some(b"RAW".to_vec())),
+            (6, "Null".into(), None),
+        ]
+    );
+}
+
+#[tokio::test]
+async fn test_blob_descriptor_short_read_returns_error() {
+    let (tmp, sql_context) = setup(BLOB_TABLE_DDL).await;
+
+    let source_path = tmp.path().join("blob_short_read.bin");
+    std::fs::write(&source_path, b"short").unwrap();
+    let uri = format!("file://{}", source_path.display());
+    let desc_hex = to_hex(&BlobDescriptor::new(uri, 0, 6).serialize());
+    let sql = format!(
+        "INSERT INTO paimon.test_db.t (id, name, picture) VALUES (1, 'Short', 
X'{desc_hex}')"
+    );
+
+    assert_sql_error(&sql_context, &sql, "Failed to read 
BlobDescriptor").await;
+}
+
 /// Blob files roll independently when `blob.target-file-size` is small.
 #[tokio::test]
 async fn test_blob_rolling() {
@@ -709,6 +769,123 @@ async fn 
test_blob_descriptor_field_resolve_descriptor_value() {
     );
 }
 
+#[tokio::test]
+async fn test_blob_descriptor_field_resolve_unknown_length_descriptor() {
+    let (tmp, sql_context) = setup(
+        "CREATE TABLE paimon.test_db.t (\
+            id INT, \
+            name STRING, \
+            picture BLOB \
+         ) WITH (\
+            'data-evolution.enabled' = 'true', \
+            'row-tracking.enabled' = 'true', \
+            'blob-descriptor-field' = 'picture'\
+         )",
+    )
+    .await;
+
+    let source_data = b"HEADER_PAYLOAD_TRAILER";
+    let source_path = tmp.path().join("descriptor_unknown_length.bin");
+    std::fs::write(&source_path, source_data).unwrap();
+
+    let uri = format!("file://{}", source_path.display());
+    let bounded_hex = to_hex(&BlobDescriptor::new(uri.clone(), 0, 
6).serialize());
+    let suffix_hex = to_hex(&BlobDescriptor::new(uri.clone(), 7, 
-1).serialize());
+    let eof_hex =
+        to_hex(&BlobDescriptor::new(uri.clone(), source_data.len() as i64, 
-1).serialize());
+    let past_eof_hex =
+        to_hex(&BlobDescriptor::new(uri, source_data.len() as i64 + 5, 
-1).serialize());
+    let sql = format!(
+        "INSERT INTO paimon.test_db.t (id, name, picture) VALUES \
+         (1, 'Bounded', X'{bounded_hex}'), \
+         (2, 'Suffix', X'{suffix_hex}'), \
+         (3, 'Eof', X'{eof_hex}'), \
+         (4, 'PastEof', X'{past_eof_hex}'), \
+         (5, 'Raw', X'524157'), \
+         (6, 'Null', NULL)"
+    );
+    exec(&sql_context, &sql).await;
+
+    let rows = query_id_name_picture(
+        &sql_context,
+        "SELECT id, name, picture FROM paimon.test_db.t ORDER BY id",
+    )
+    .await;
+    assert_eq!(
+        rows,
+        vec![
+            (1, "Bounded".into(), Some(b"HEADER".to_vec())),
+            (2, "Suffix".into(), Some(b"PAYLOAD_TRAILER".to_vec())),
+            (3, "Eof".into(), Some(Vec::new())),
+            (4, "PastEof".into(), Some(Vec::new())),
+            (5, "Raw".into(), Some(b"RAW".to_vec())),
+            (6, "Null".into(), None),
+        ]
+    );
+}
+
+#[tokio::test]
+async fn test_blob_descriptor_field_rejects_invalid_length() {
+    let (tmp, sql_context) = setup(
+        "CREATE TABLE paimon.test_db.t (\
+            id INT, \
+            name STRING, \
+            picture BLOB \
+         ) WITH (\
+            'data-evolution.enabled' = 'true', \
+            'row-tracking.enabled' = 'true', \
+            'blob-descriptor-field' = 'picture'\
+         )",
+    )
+    .await;
+
+    let uri = format!("file://{}", tmp.path().join("unused.bin").display());
+    let desc_hex = to_hex(&BlobDescriptor::new(uri, 0, -2).serialize());
+    let sql = format!(
+        "INSERT INTO paimon.test_db.t (id, name, picture) VALUES (1, 
'Invalid', X'{desc_hex}')"
+    );
+    exec(&sql_context, &sql).await;
+
+    assert_sql_error(
+        &sql_context,
+        "SELECT id, name, picture FROM paimon.test_db.t",
+        "length must be -1 or non-negative",
+    )
+    .await;
+}
+
+#[tokio::test]
+async fn test_blob_descriptor_field_short_read_returns_error() {
+    let (tmp, sql_context) = setup(
+        "CREATE TABLE paimon.test_db.t (\
+            id INT, \
+            name STRING, \
+            picture BLOB \
+         ) WITH (\
+            'data-evolution.enabled' = 'true', \
+            'row-tracking.enabled' = 'true', \
+            'blob-descriptor-field' = 'picture'\
+         )",
+    )
+    .await;
+
+    let source_path = tmp.path().join("descriptor_short_read.bin");
+    std::fs::write(&source_path, b"short").unwrap();
+    let uri = format!("file://{}", source_path.display());
+    let desc_hex = to_hex(&BlobDescriptor::new(uri, 0, 6).serialize());
+    let sql = format!(
+        "INSERT INTO paimon.test_db.t (id, name, picture) VALUES (1, 'Short', 
X'{desc_hex}')"
+    );
+    exec(&sql_context, &sql).await;
+
+    assert_sql_error(
+        &sql_context,
+        "SELECT id, name, picture FROM paimon.test_db.t",
+        "Failed to read BlobDescriptor",
+    )
+    .await;
+}
+
 #[tokio::test]
 async fn 
test_blob_descriptor_filter_before_resolve_skips_filtered_bad_descriptor() {
     let (tmp, sql_context) = setup(
diff --git a/crates/paimon/src/arrow/format/blob.rs 
b/crates/paimon/src/arrow/format/blob.rs
index 136170e..c401a06 100644
--- a/crates/paimon/src/arrow/format/blob.rs
+++ b/crates/paimon/src/arrow/format/blob.rs
@@ -617,6 +617,23 @@ impl BlobFormatWriter {
 
 const BLOB_WRITE_BUFFER_SIZE: u64 = 8 * 1024 * 1024; // 8 MB
 
+fn checked_blob_entry_length(payload_len: u64) -> crate::Result<i64> {
+    let entry_length = payload_len
+        .checked_add(BLOB_ENTRY_OVERHEAD)
+        .ok_or_else(|| Error::DataInvalid {
+            message: format!(
+                "Blob entry length overflows u64: 
payload_length={payload_len}, overhead={BLOB_ENTRY_OVERHEAD}"
+            ),
+            source: None,
+        })?;
+    i64::try_from(entry_length).map_err(|e| Error::DataInvalid {
+        message: format!(
+            "Blob entry length exceeds i64: payload_length={payload_len}, 
entry_length={entry_length}"
+        ),
+        source: Some(Box::new(e)),
+    })
+}
+
 #[async_trait]
 impl FormatFileWriter for BlobFormatWriter {
     async fn write(&mut self, batch: &RecordBatch) -> crate::Result<()> {
@@ -643,9 +660,7 @@ impl FormatFileWriter for BlobFormatWriter {
 
             if BlobDescriptor::is_blob_descriptor(value) {
                 let desc = BlobDescriptor::deserialize(value)?;
-                let payload_len = desc.length() as u64;
-                let entry_length = (payload_len + BLOB_ENTRY_OVERHEAD) as i64;
-                self.lengths.push(entry_length);
+                let range = desc.range_spec()?;
 
                 let file_io = self.file_io.as_ref().ok_or_else(|| 
Error::DataInvalid {
                     message:
@@ -654,7 +669,47 @@ impl FormatFileWriter for BlobFormatWriter {
                     source: None,
                 })?;
                 let input = file_io.new_input(desc.uri())?;
-                let reader = input.reader().await?;
+                let offset = range.offset();
+                let payload_len = match range.length() {
+                    Some(length) => length,
+                    None => input
+                        .metadata()
+                        .await
+                        .map_err(|e| Error::UnexpectedError {
+                            message: format!(
+                                "Failed to read metadata for BlobDescriptor 
'{}': {e}",
+                                desc.uri()
+                            ),
+                            source: Some(Box::new(e)),
+                        })?
+                        .size
+                        .saturating_sub(offset),
+                };
+                let end = offset
+                    .checked_add(payload_len)
+                    .ok_or_else(|| Error::DataInvalid {
+                        message: format!(
+                            "BlobDescriptor range overflows u64: 
offset={offset}, length={payload_len}"
+                        ),
+                        source: None,
+                    })?;
+                let entry_length = checked_blob_entry_length(payload_len)?;
+                let entry_length_u64 = entry_length as u64;
+                let bytes_written = self
+                    .bytes_written
+                    .checked_add(entry_length_u64)
+                    .ok_or_else(|| Error::DataInvalid {
+                        message: format!(
+                            "Blob file size overflows u64: current_size={}, 
entry_length={entry_length_u64}",
+                            self.bytes_written
+                        ),
+                        source: None,
+                    })?;
+                let reader = if payload_len == 0 {
+                    None
+                } else {
+                    Some(input.reader().await?)
+                };
 
                 let mut hasher = crc32fast::Hasher::new();
 
@@ -664,15 +719,34 @@ impl FormatFileWriter for BlobFormatWriter {
                     .await?;
 
                 // Stream payload in chunks to avoid loading entire blob into 
memory
-                let start = desc.offset() as u64;
-                let end = start + payload_len;
-                let mut pos = start;
-                while pos < end {
-                    let chunk_end = (pos + BLOB_WRITE_BUFFER_SIZE).min(end);
-                    let chunk = reader.read(pos..chunk_end).await?;
-                    hasher.update(&chunk);
-                    self.writer.write(chunk).await?;
-                    pos = chunk_end;
+                if let Some(reader) = reader.as_ref() {
+                    let mut pos = offset;
+                    while pos < end {
+                        let chunk_end = 
pos.saturating_add(BLOB_WRITE_BUFFER_SIZE).min(end);
+                        let chunk = 
reader.read(pos..chunk_end).await.map_err(|e| {
+                            Error::UnexpectedError {
+                                message: format!(
+                                    "Failed to read BlobDescriptor '{}' range 
{pos}..{chunk_end}: {e}",
+                                    desc.uri()
+                                ),
+                                source: Some(Box::new(e)),
+                            }
+                        })?;
+                        let actual_len = chunk.len() as u64;
+                        let expected_len = chunk_end - pos;
+                        if actual_len != expected_len {
+                            return Err(Error::DataInvalid {
+                                message: format!(
+                                    "Failed to read BlobDescriptor '{}': short 
read for range {pos}..{chunk_end}, expected={expected_len} bytes, 
actual={actual_len} bytes",
+                                    desc.uri()
+                                ),
+                                source: None,
+                            });
+                        }
+                        hasher.update(&chunk);
+                        self.writer.write(chunk).await?;
+                        pos = chunk_end;
+                    }
                 }
 
                 let entry_length_bytes = entry_length.to_le_bytes();
@@ -685,7 +759,8 @@ impl FormatFileWriter for BlobFormatWriter {
                     
.write(Bytes::copy_from_slice(&hasher.finalize().to_le_bytes()))
                     .await?;
 
-                self.bytes_written += entry_length as u64;
+                self.lengths.push(entry_length);
+                self.bytes_written = bytes_written;
             } else {
                 let entry_length = (value.len() + BLOB_ENTRY_OVERHEAD as 
usize) as i64;
                 self.lengths.push(entry_length);
@@ -1050,6 +1125,19 @@ mod tests {
         assert_eq!(decoded, values);
     }
 
+    #[test]
+    fn test_checked_blob_entry_length() {
+        assert_eq!(
+            checked_blob_entry_length(0).unwrap(),
+            BLOB_ENTRY_OVERHEAD as i64
+        );
+
+        let max_payload = i64::MAX as u64 - BLOB_ENTRY_OVERHEAD;
+        assert_eq!(checked_blob_entry_length(max_payload).unwrap(), i64::MAX);
+        assert!(checked_blob_entry_length(max_payload + 1).is_err());
+        assert!(checked_blob_entry_length(u64::MAX).is_err());
+    }
+
     fn basic_blob_rows() -> [Option<&'static [u8]>; 4] {
         [
             Some(&b"hello"[..]),
diff --git a/crates/paimon/src/spec/blob_descriptor.rs 
b/crates/paimon/src/spec/blob_descriptor.rs
index 8924ca6..93c8449 100644
--- a/crates/paimon/src/spec/blob_descriptor.rs
+++ b/crates/paimon/src/spec/blob_descriptor.rs
@@ -28,6 +28,22 @@ pub struct BlobDescriptor {
     length: i64,
 }
 
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub(crate) struct BlobRangeSpec {
+    offset: u64,
+    length: Option<u64>,
+}
+
+impl BlobRangeSpec {
+    pub(crate) fn offset(self) -> u64 {
+        self.offset
+    }
+
+    pub(crate) fn length(self) -> Option<u64> {
+        self.length
+    }
+}
+
 impl BlobDescriptor {
     pub fn new(uri: String, offset: i64, length: i64) -> Self {
         Self {
@@ -50,6 +66,46 @@ impl BlobDescriptor {
         self.length
     }
 
+    pub(crate) fn range_spec(&self) -> crate::Result<BlobRangeSpec> {
+        if self.offset < 0 {
+            return Err(Error::DataInvalid {
+                message: format!(
+                    "BlobDescriptor offset must be non-negative: {}",
+                    self.offset
+                ),
+                source: None,
+            });
+        }
+        if self.length < -1 {
+            return Err(Error::DataInvalid {
+                message: format!(
+                    "BlobDescriptor length must be -1 or non-negative: {}",
+                    self.length
+                ),
+                source: None,
+            });
+        }
+
+        let offset = self.offset as u64;
+        let length = if self.length == -1 {
+            None
+        } else {
+            Some(self.length as u64)
+        };
+        if let Some(length) = length {
+            offset
+                .checked_add(length)
+                .ok_or_else(|| Error::DataInvalid {
+                    message: format!(
+                        "BlobDescriptor range overflows u64: offset={offset}, 
length={length}"
+                    ),
+                    source: None,
+                })?;
+        }
+
+        Ok(BlobRangeSpec { offset, length })
+    }
+
     pub fn serialize(&self) -> Vec<u8> {
         let uri_bytes = self.uri.as_bytes();
         let uri_length = uri_bytes.len();
@@ -198,4 +254,60 @@ mod tests {
         bytes[1] = 0xFF;
         assert!(BlobDescriptor::deserialize(&bytes).is_err());
     }
+
+    #[test]
+    fn test_range_spec_supports_unknown_length() {
+        let full = BlobDescriptor::new("x".to_string(), 0, -1)
+            .range_spec()
+            .unwrap();
+        assert_eq!(full.offset(), 0);
+        assert_eq!(full.length(), None);
+
+        let range = BlobDescriptor::new("x".to_string(), 7, -1)
+            .range_spec()
+            .unwrap();
+        assert_eq!(range.offset(), 7);
+        assert_eq!(range.length(), None);
+    }
+
+    #[test]
+    fn test_range_spec_supports_bounded_and_empty_ranges() {
+        let bounded = BlobDescriptor::new("x".to_string(), 7, 11)
+            .range_spec()
+            .unwrap();
+        assert_eq!(bounded.offset(), 7);
+        assert_eq!(bounded.length(), Some(11));
+
+        let empty = BlobDescriptor::new("x".to_string(), 7, 0)
+            .range_spec()
+            .unwrap();
+        assert_eq!(empty.offset(), 7);
+        assert_eq!(empty.length(), Some(0));
+    }
+
+    #[test]
+    fn test_range_spec_rejects_invalid_signed_values() {
+        let err = BlobDescriptor::new("x".to_string(), -1, 1)
+            .range_spec()
+            .unwrap_err();
+        assert!(
+            matches!(err, Error::DataInvalid { message, .. } if 
message.contains("offset must be non-negative"))
+        );
+
+        let err = BlobDescriptor::new("x".to_string(), 0, -2)
+            .range_spec()
+            .unwrap_err();
+        assert!(
+            matches!(err, Error::DataInvalid { message, .. } if 
message.contains("length must be -1 or non-negative"))
+        );
+    }
+
+    #[test]
+    fn test_range_spec_handles_largest_bounded_range() {
+        let range = BlobDescriptor::new("x".to_string(), i64::MAX, i64::MAX)
+            .range_spec()
+            .unwrap();
+        assert_eq!(range.offset(), i64::MAX as u64);
+        assert_eq!(range.length(), Some(i64::MAX as u64));
+    }
 }
diff --git a/crates/paimon/src/table/blob_resolver.rs 
b/crates/paimon/src/table/blob_resolver.rs
index bbbf24b..d3f5a82 100644
--- a/crates/paimon/src/table/blob_resolver.rs
+++ b/crates/paimon/src/table/blob_resolver.rs
@@ -46,7 +46,7 @@ pub(crate) async fn resolve_blob_column(
     }
 
     let mut cells = Vec::with_capacity(col.len());
-    let mut requests_by_uri: HashMap<String, Vec<BlobReadRequest>> = 
HashMap::new();
+    let mut requests_by_uri: HashMap<String, Vec<BlobReadRequestSpec>> = 
HashMap::new();
     let mut value_capacity = 0usize;
 
     for row in 0..col.len() {
@@ -58,36 +58,14 @@ pub(crate) async fn resolve_blob_column(
         let value = col.value(row);
         if BlobDescriptor::is_blob_descriptor(value) {
             let desc = BlobDescriptor::deserialize(value)?;
-            let offset = u64::try_from(desc.offset()).map_err(|e| 
crate::Error::DataInvalid {
-                message: format!(
-                    "BlobDescriptor offset must be non-negative: {}",
-                    desc.offset()
-                ),
-                source: Some(Box::new(e)),
-            })?;
-            let length = u64::try_from(desc.length()).map_err(|e| 
crate::Error::DataInvalid {
-                message: format!(
-                    "BlobDescriptor length must be non-negative: {}",
-                    desc.length()
-                ),
-                source: Some(Box::new(e)),
-            })?;
-            offset
-                .checked_add(length)
-                .ok_or_else(|| crate::Error::DataInvalid {
-                    message: format!(
-                        "BlobDescriptor range overflows u64: offset={offset}, 
length={length}"
-                    ),
-                    source: None,
-                })?;
-            value_capacity = value_capacity.saturating_add(length as usize);
+            let range = desc.range_spec()?;
             requests_by_uri
                 .entry(desc.uri().to_string())
                 .or_default()
-                .push(BlobReadRequest {
+                .push(BlobReadRequestSpec {
                     row,
-                    offset,
-                    length,
+                    offset: range.offset(),
+                    length: range.length(),
                 });
             cells.push(ResolvedBlobCell::Null);
         } else {
@@ -96,19 +74,101 @@ pub(crate) async fn resolve_blob_column(
         }
     }
 
-    let mut readers: HashMap<String, Box<dyn FileRead>> = HashMap::new();
     for (uri, requests) in requests_by_uri {
-        if !readers.contains_key(&uri) {
-            let input = file_io.new_input(&uri)?;
-            let reader = input.reader().await?;
-            readers.insert(uri.clone(), Box::new(reader));
+        let input = file_io.new_input(&uri)?;
+        let file_size = if requests.iter().any(|request| 
request.length.is_none()) {
+            input
+                .metadata()
+                .await
+                .map_err(|e| crate::Error::UnexpectedError {
+                    message: format!("Failed to read metadata for 
BlobDescriptor URI '{uri}': {e}"),
+                    source: Some(Box::new(e)),
+                })?
+                .size
+        } else {
+            0
+        };
+        let mut bounded_requests = Vec::with_capacity(requests.len());
+        for request in requests {
+            let length = request
+                .length
+                .unwrap_or_else(|| file_size.saturating_sub(request.offset));
+            request
+                .offset
+                .checked_add(length)
+                .ok_or_else(|| crate::Error::DataInvalid {
+                    message: format!(
+                        "BlobDescriptor range overflows u64: offset={}, 
length={length}",
+                        request.offset
+                    ),
+                    source: None,
+                })?;
+            value_capacity = value_capacity.saturating_add(length as usize);
+            if length == 0 {
+                cells[request.row] = ResolvedBlobCell::Value(Bytes::new());
+                continue;
+            }
+            bounded_requests.push(BlobReadRequest {
+                row: request.row,
+                offset: request.offset,
+                length,
+            });
         }
-        let reader = readers.get(&uri).unwrap();
-        for merged in merge_blob_read_requests(requests) {
-            let data = reader.read(merged.start..merged.end).await?;
+
+        if bounded_requests.is_empty() {
+            continue;
+        }
+
+        let reader = input.reader().await?;
+        for merged in merge_blob_read_requests(bounded_requests) {
+            let data = reader.read(merged.start..merged.end).await.map_err(|e| 
{
+                crate::Error::UnexpectedError {
+                    message: format!(
+                        "Failed to read BlobDescriptor URI '{uri}' range 
{}..{}: {e}",
+                        merged.start, merged.end
+                    ),
+                    source: Some(Box::new(e)),
+                }
+            })?;
+            let expected_len = merged.end - merged.start;
+            let actual_len = data.len() as u64;
+            if actual_len != expected_len {
+                return Err(crate::Error::DataInvalid {
+                    message: format!(
+                        "Failed to read BlobDescriptor URI '{uri}': short read 
for range {}..{}, expected={expected_len} bytes, actual={actual_len} bytes",
+                        merged.start, merged.end
+                    ),
+                    source: None,
+                });
+            }
             for request in merged.requests {
-                let start = (request.offset - merged.start) as usize;
-                let end = start + request.length as usize;
+                let start = usize::try_from(request.offset - 
merged.start).map_err(|e| {
+                    crate::Error::DataInvalid {
+                        message: format!(
+                            "BlobDescriptor slice offset exceeds usize: 
offset={}, merged_start={}",
+                            request.offset, merged.start
+                        ),
+                        source: Some(Box::new(e)),
+                    }
+                })?;
+                let length =
+                    usize::try_from(request.length).map_err(|e| 
crate::Error::DataInvalid {
+                        message: format!(
+                            "BlobDescriptor slice length exceeds usize: {}",
+                            request.length
+                        ),
+                        source: Some(Box::new(e)),
+                    })?;
+                let end = start
+                    .checked_add(length)
+                    .filter(|end| *end <= data.len())
+                    .ok_or_else(|| crate::Error::DataInvalid {
+                        message: format!(
+                            "BlobDescriptor slice exceeds read data: 
start={start}, length={length}, actual={}",
+                            data.len()
+                        ),
+                        source: None,
+                    })?;
                 cells[request.row] = 
ResolvedBlobCell::Value(data.slice(start..end));
             }
         }
@@ -130,6 +190,13 @@ enum ResolvedBlobCell {
     Value(Bytes),
 }
 
+#[derive(Debug)]
+struct BlobReadRequestSpec {
+    row: usize,
+    offset: u64,
+    length: Option<u64>,
+}
+
 #[derive(Debug, Clone, PartialEq, Eq)]
 struct BlobReadRequest {
     row: usize,
diff --git a/docs/src/sql.md b/docs/src/sql.md
index 0207062..d550c73 100644
--- a/docs/src/sql.md
+++ b/docs/src/sql.md
@@ -294,6 +294,10 @@ the supported directives is rejected.
 | `__BLOB_DESCRIPTOR_FIELD` | `blob-descriptor-field` | Store serialized 
`BlobDescriptor` bytes inline |
 | `__BLOB_VIEW_FIELD` | `blob-view-field` | Store serialized `BlobViewStruct` 
bytes inline |
 
+For serialized `BlobDescriptor` values supplied by another Paimon engine,
+`length = -1` means reading from `offset` to the end of the referenced object.
+The offset must be non-negative, and lengths below `-1` are invalid.
+
 The same directives are supported by `ALTER TABLE ... ADD COLUMN`.
 
 ### Blob View

Reply via email to