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 e242cbd1 feat(blob): support ARRAY<BLOB> reads (#566)
e242cbd1 is described below
commit e242cbd16f6cadd9ed875f999371bbae80c7e89c
Author: QuakeWang <[email protected]>
AuthorDate: Tue Jul 21 12:16:35 2026 +0800
feat(blob): support ARRAY<BLOB> reads (#566)
---
crates/paimon/src/arrow/format/blob.rs | 838 ++++++++++++++++++++-
crates/paimon/src/spec/types.rs | 21 +
crates/paimon/src/table/data_evolution_reader.rs | 300 +++++++-
.../table/data_evolution_reader/blob_fallback.rs | 47 +-
crates/paimon/testdata/blob/blob-array.blob | Bin 0 -> 80 bytes
5 files changed, 1158 insertions(+), 48 deletions(-)
diff --git a/crates/paimon/src/arrow/format/blob.rs
b/crates/paimon/src/arrow/format/blob.rs
index f66108c2..5b455735 100644
--- a/crates/paimon/src/arrow/format/blob.rs
+++ b/crates/paimon/src/arrow/format/blob.rs
@@ -21,7 +21,7 @@ use crate::io::{FileRead, FileWrite};
use crate::spec::{BlobDescriptor, DataField, DataType};
use crate::table::{ArrowRecordBatchStream, RowRange};
use crate::Error;
-use arrow_array::builder::BinaryBuilder;
+use arrow_array::builder::{BinaryBuilder, ListBuilder};
use arrow_array::{Array, ArrayRef, RecordBatch, RecordBatchOptions};
use async_stream::try_stream;
use async_trait::async_trait;
@@ -82,11 +82,26 @@ impl IndexedBlobReader {
fetch_blob_values(self.reader.as_ref(), planned_reads).await
}
}
+
+ pub(crate) async fn read_array_positions(
+ &self,
+ positions: &[usize],
+ ) -> crate::Result<Vec<BlobReadValue>> {
+ let planned_reads = plan_blob_array_reads(&self.index, positions)?;
+ fetch_blob_array_values(
+ self.reader.as_ref(),
+ planned_reads,
+ &self.file_path,
+ self.descriptor_mode,
+ )
+ .await
+ }
}
#[derive(Debug)]
pub(crate) enum BlobReadValue {
Value(Bytes),
+ Array(Vec<Option<Bytes>>),
Null,
Placeholder,
}
@@ -98,6 +113,18 @@ const BLOB_TRAILER_SIZE: u64 = 12;
const BLOB_ENTRY_OVERHEAD: u64 = BLOB_INLINE_HEADER_SIZE + BLOB_TRAILER_SIZE;
const DEFAULT_BATCH_SIZE: usize = 128;
const BLOB_READ_CONCURRENCY: usize = 8;
+const BLOB_ARRAY_MAGIC_NUMBER: i32 = 1094861634;
+const BLOB_ARRAY_VERSION: u8 = 1;
+const BLOB_ARRAY_HEADER_SIZE: u64 = 9;
+const BLOB_ARRAY_INDEX_LENGTH_SIZE: u64 = 4;
+const BLOB_ARRAY_MIN_PAYLOAD_SIZE: u64 = BLOB_ARRAY_HEADER_SIZE +
BLOB_ARRAY_INDEX_LENGTH_SIZE;
+const BLOB_ARRAY_NULL_ELEMENT_LENGTH: i64 = -1;
+
+#[derive(Debug, Clone, Copy)]
+enum BlobFieldKind {
+ Scalar,
+ Array,
+}
#[async_trait]
impl FormatFileReader for BlobFormatReader {
@@ -110,7 +137,7 @@ impl FormatFileReader for BlobFormatReader {
batch_size: Option<usize>,
row_selection: Option<Vec<RowRange>>,
) -> crate::Result<ArrowRecordBatchStream> {
- validate_read_fields(read_fields)?;
+ let field_kind = validate_read_fields(read_fields)?;
let target_schema = build_target_arrow_schema(read_fields)?;
let batch_size = batch_size.unwrap_or(DEFAULT_BATCH_SIZE);
@@ -122,15 +149,19 @@ impl FormatFileReader for BlobFormatReader {
)
.await?;
let mut selection = RowSelectionCursor::new(blob_reader.num_rows(),
row_selection)?;
- let project_values = !read_fields.is_empty();
Ok(try_stream! {
while let Some(positions) = selection.next_batch(batch_size) {
- let batch = if project_values {
- let values = blob_reader.read_positions(&positions).await?;
- build_blob_batch(&target_schema, values)?
- } else {
- RecordBatch::try_new_with_options(
+ let batch = match field_kind {
+ Some(BlobFieldKind::Scalar) => {
+ let values =
blob_reader.read_positions(&positions).await?;
+ build_blob_batch(&target_schema, values)?
+ }
+ Some(BlobFieldKind::Array) => {
+ let values =
blob_reader.read_array_positions(&positions).await?;
+ build_blob_array_batch(&target_schema, values)?
+ }
+ None => RecordBatch::try_new_with_options(
target_schema.clone(),
Vec::new(),
&RecordBatchOptions::new().with_row_count(Some(positions.len())),
@@ -138,7 +169,7 @@ impl FormatFileReader for BlobFormatReader {
.map_err(|e| Error::UnexpectedError {
message: format!("Failed to build empty blob
RecordBatch: {e}"),
source: Some(Box::new(e)),
- })?
+ })?,
};
yield batch;
}
@@ -147,7 +178,7 @@ impl FormatFileReader for BlobFormatReader {
}
}
-fn validate_read_fields(read_fields: &[DataField]) -> crate::Result<()> {
+fn validate_read_fields(read_fields: &[DataField]) ->
crate::Result<Option<BlobFieldKind>> {
if read_fields.len() > 1 {
return Err(Error::DataInvalid {
message: format!(
@@ -158,21 +189,23 @@ fn validate_read_fields(read_fields: &[DataField]) ->
crate::Result<()> {
});
}
- if let Some(field) = read_fields.first() {
- match field.data_type() {
- DataType::Blob(_) => Ok(()),
+ read_fields
+ .first()
+ .map(|field| match field.data_type() {
+ DataType::Blob(_) => Ok(BlobFieldKind::Scalar),
+ DataType::Array(array) if matches!(array.element_type(),
DataType::Blob(_)) => {
+ Ok(BlobFieldKind::Array)
+ }
other => Err(Error::DataInvalid {
message: format!(
- ".blob format requires a Blob field, got {:?} for column
'{}'",
+ ".blob format requires a Blob or Array<Blob> field, got
{:?} for column '{}'",
other,
field.name()
),
source: None,
}),
- }?;
- }
-
- Ok(())
+ })
+ .transpose()
}
fn build_descriptor_values(
@@ -209,7 +242,7 @@ fn build_descriptor_values(
.collect()
}
-fn build_blob_batch(
+pub(crate) fn build_blob_batch(
target_schema: &Arc<arrow_schema::Schema>,
values: Vec<BlobReadValue>,
) -> crate::Result<RecordBatch> {
@@ -218,6 +251,12 @@ fn build_blob_batch(
match value {
BlobReadValue::Value(bytes) =>
builder.append_value(bytes.as_ref()),
BlobReadValue::Null | BlobReadValue::Placeholder =>
builder.append_null(),
+ BlobReadValue::Array(_) => {
+ return Err(Error::UnexpectedError {
+ message: "Scalar BLOB reader produced an ARRAY<BLOB>
value".to_string(),
+ source: None,
+ });
+ }
}
}
@@ -228,6 +267,50 @@ fn build_blob_batch(
})
}
+pub(crate) fn build_blob_array_batch(
+ target_schema: &Arc<arrow_schema::Schema>,
+ values: Vec<BlobReadValue>,
+) -> crate::Result<RecordBatch> {
+ let element_field = match target_schema.field(0).data_type() {
+ arrow_schema::DataType::List(element_field) => element_field.clone(),
+ other => {
+ return Err(Error::UnexpectedError {
+ message: format!(
+ "Expected Array<Blob> to map to Arrow List<Binary>, got
{other:?}"
+ ),
+ source: None,
+ });
+ }
+ };
+ let mut builder =
ListBuilder::new(BinaryBuilder::new()).with_field(element_field);
+ for value in values {
+ match value {
+ BlobReadValue::Array(elements) => {
+ for element in elements {
+ match element {
+ Some(bytes) =>
builder.values().append_value(bytes.as_ref()),
+ None => builder.values().append_null(),
+ }
+ }
+ builder.append(true);
+ }
+ BlobReadValue::Null | BlobReadValue::Placeholder =>
builder.append(false),
+ BlobReadValue::Value(_) => {
+ return Err(Error::UnexpectedError {
+ message: "ARRAY<BLOB> reader produced a scalar BLOB
value".to_string(),
+ source: None,
+ });
+ }
+ }
+ }
+
+ let columns: Vec<ArrayRef> = vec![Arc::new(builder.finish())];
+ RecordBatch::try_new(target_schema.clone(), columns).map_err(|e|
Error::UnexpectedError {
+ message: format!("Failed to build ARRAY<BLOB> RecordBatch: {e}"),
+ source: Some(Box::new(e)),
+ })
+}
+
fn plan_blob_reads(
blob_index: &BlobFileIndex,
positions: &[usize],
@@ -272,6 +355,279 @@ async fn fetch_blob_values(
.await
}
+fn plan_blob_array_reads(
+ blob_index: &BlobFileIndex,
+ positions: &[usize],
+) -> crate::Result<Vec<PlannedBlobArrayRead>> {
+ positions
+ .iter()
+ .map(|&position| {
+ let entry = blob_index
+ .entry(position)
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!(
+ "Blob row selection referenced out-of-range position
{position} for {} rows",
+ blob_index.num_rows()
+ ),
+ source: None,
+ })?;
+
+ Ok(match entry {
+ BlobEntry::Value(range) =>
PlannedBlobArrayRead::Read(range.clone()),
+ BlobEntry::Null => PlannedBlobArrayRead::Null,
+ BlobEntry::Placeholder => PlannedBlobArrayRead::Placeholder,
+ })
+ })
+ .collect()
+}
+
+async fn fetch_blob_array_values(
+ reader: &dyn FileRead,
+ planned_reads: Vec<PlannedBlobArrayRead>,
+ file_path: &str,
+ descriptor_mode: bool,
+) -> crate::Result<Vec<BlobReadValue>> {
+ futures::stream::iter(planned_reads.into_iter().map(|planned_read| async
move {
+ match planned_read {
+ PlannedBlobArrayRead::Null => Ok(BlobReadValue::Null),
+ PlannedBlobArrayRead::Placeholder =>
Ok(BlobReadValue::Placeholder),
+ PlannedBlobArrayRead::Read(payload_range) => {
+ let metadata = read_blob_array_metadata(reader,
payload_range).await?;
+ if descriptor_mode {
+ build_blob_array_descriptors(metadata, file_path)
+ } else {
+ read_inline_blob_array(reader, metadata).await
+ }
+ }
+ }
+ }))
+ .buffered(BLOB_READ_CONCURRENCY)
+ .try_collect()
+ .await
+}
+
+async fn read_blob_array_metadata(
+ reader: &dyn FileRead,
+ payload_range: Range<u64>,
+) -> crate::Result<BlobArrayMetadata> {
+ let payload_length = payload_range
+ .end
+ .checked_sub(payload_range.start)
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!("Invalid ARRAY<BLOB> payload range:
{payload_range:?}"),
+ source: None,
+ })?;
+ if payload_length < BLOB_ARRAY_MIN_PAYLOAD_SIZE {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "ARRAY<BLOB> payload is too small: expected at least
{BLOB_ARRAY_MIN_PAYLOAD_SIZE} bytes, got {payload_length}"
+ ),
+ source: None,
+ });
+ }
+
+ let header_end = payload_range.start + BLOB_ARRAY_HEADER_SIZE;
+ let header = read_blob_array_range(reader,
payload_range.start..header_end, "header").await?;
+ let magic = i32::from_le_bytes(header[..4].try_into().unwrap());
+ if magic != BLOB_ARRAY_MAGIC_NUMBER {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Invalid ARRAY<BLOB> payload magic number: expected
{BLOB_ARRAY_MAGIC_NUMBER}, got {magic}"
+ ),
+ source: None,
+ });
+ }
+ if header[4] != BLOB_ARRAY_VERSION {
+ return Err(Error::Unsupported {
+ message: format!(
+ "Unsupported ARRAY<BLOB> payload version: expected
{BLOB_ARRAY_VERSION}, got {}",
+ header[4]
+ ),
+ });
+ }
+ let element_count = i32::from_le_bytes(header[5..9].try_into().unwrap());
+ if element_count < 0 {
+ return Err(Error::DataInvalid {
+ message: format!("Invalid ARRAY<BLOB> element count:
{element_count}"),
+ source: None,
+ });
+ }
+
+ let index_length_position = payload_range.end -
BLOB_ARRAY_INDEX_LENGTH_SIZE;
+ let index_length_bytes = read_blob_array_range(
+ reader,
+ index_length_position..payload_range.end,
+ "index length",
+ )
+ .await?;
+ let index_length =
i32::from_le_bytes(index_length_bytes[..4].try_into().unwrap());
+ let maximum_index_length = payload_length - BLOB_ARRAY_MIN_PAYLOAD_SIZE;
+ if index_length < 0 || index_length as u64 > maximum_index_length {
+ return Err(Error::DataInvalid {
+ message: format!("Invalid ARRAY<BLOB> element index length:
{index_length}"),
+ source: None,
+ });
+ }
+ let index_length = index_length as u64;
+ if element_count as u64 > index_length {
+ return Err(Error::DataInvalid {
+ message: "ARRAY<BLOB> element count exceeds element index
length".to_string(),
+ source: None,
+ });
+ }
+
+ let index_start = index_length_position - index_length;
+ let index_bytes = if index_length == 0 {
+ Bytes::new()
+ } else {
+ read_blob_array_range(reader, index_start..index_length_position,
"element index").await?
+ };
+ let encoded_lengths =
+ decode_delta_varints(index_bytes.as_ref()).map_err(|e|
Error::DataInvalid {
+ message: format!("Invalid ARRAY<BLOB> element index: {e}"),
+ source: Some(Box::new(e)),
+ })?;
+ if encoded_lengths.len() != element_count as usize {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "ARRAY<BLOB> element count {element_count} does not match
index value count {}",
+ encoded_lengths.len()
+ ),
+ source: None,
+ });
+ }
+
+ let element_data_range = header_end..index_start;
+ let mut remaining_data_length = element_data_range.end -
element_data_range.start;
+ let mut element_lengths = Vec::with_capacity(encoded_lengths.len());
+ for encoded_length in encoded_lengths {
+ if encoded_length == BLOB_ARRAY_NULL_ELEMENT_LENGTH {
+ element_lengths.push(None);
+ continue;
+ }
+ let element_length = u64::try_from(encoded_length).map_err(|e|
Error::DataInvalid {
+ message: format!("Invalid ARRAY<BLOB> element length:
{encoded_length}"),
+ source: Some(Box::new(e)),
+ })?;
+ if element_length > remaining_data_length {
+ return Err(Error::DataInvalid {
+ message: "ARRAY<BLOB> element lengths exceed the payload data
length".to_string(),
+ source: None,
+ });
+ }
+ remaining_data_length -= element_length;
+ element_lengths.push(Some(element_length));
+ }
+ if remaining_data_length != 0 {
+ return Err(Error::DataInvalid {
+ message: "ARRAY<BLOB> element lengths do not match the payload
data length".to_string(),
+ source: None,
+ });
+ }
+
+ Ok(BlobArrayMetadata {
+ element_data_range,
+ element_lengths,
+ })
+}
+
+async fn read_blob_array_range(
+ reader: &dyn FileRead,
+ range: Range<u64>,
+ part: &str,
+) -> crate::Result<Bytes> {
+ let expected_length = range.end - range.start;
+ let bytes = reader
+ .read(range.clone())
+ .await
+ .map_err(|e| Error::UnexpectedError {
+ message: format!("Failed to read ARRAY<BLOB> {part} range
{range:?}: {e}"),
+ source: Some(Box::new(e)),
+ })?;
+ if bytes.len() as u64 != expected_length {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Short read for ARRAY<BLOB> {part} range {range:?}: expected
{expected_length} bytes, got {}",
+ bytes.len()
+ ),
+ source: None,
+ });
+ }
+ Ok(bytes)
+}
+
+async fn read_inline_blob_array(
+ reader: &dyn FileRead,
+ metadata: BlobArrayMetadata,
+) -> crate::Result<BlobReadValue> {
+ let data_length = metadata.element_data_range.end -
metadata.element_data_range.start;
+ if data_length > i32::MAX as u64 {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "ARRAY<BLOB> inline element data is too large for Arrow
Binary: {data_length} bytes"
+ ),
+ source: None,
+ });
+ }
+ let data = if data_length == 0 {
+ Bytes::new()
+ } else {
+ read_blob_array_range(reader, metadata.element_data_range, "element
data").await?
+ };
+
+ let mut offset = 0usize;
+ let mut elements = Vec::with_capacity(metadata.element_lengths.len());
+ for element_length in metadata.element_lengths {
+ match element_length {
+ None => elements.push(None),
+ Some(element_length) => {
+ let element_length = element_length as usize;
+ let end = offset + element_length;
+ elements.push(Some(data.slice(offset..end)));
+ offset = end;
+ }
+ }
+ }
+ Ok(BlobReadValue::Array(elements))
+}
+
+fn build_blob_array_descriptors(
+ metadata: BlobArrayMetadata,
+ file_path: &str,
+) -> crate::Result<BlobReadValue> {
+ let mut element_offset = metadata.element_data_range.start;
+ let mut elements = Vec::with_capacity(metadata.element_lengths.len());
+ for element_length in metadata.element_lengths {
+ match element_length {
+ None => elements.push(None),
+ Some(element_length) => {
+ let descriptor_offset =
+ i64::try_from(element_offset).map_err(|e|
Error::DataInvalid {
+ message: format!(
+ "ARRAY<BLOB> descriptor offset exceeds i64:
{element_offset}"
+ ),
+ source: Some(Box::new(e)),
+ })?;
+ let descriptor_length =
+ i64::try_from(element_length).map_err(|e|
Error::DataInvalid {
+ message: format!(
+ "ARRAY<BLOB> descriptor length exceeds i64:
{element_length}"
+ ),
+ source: Some(Box::new(e)),
+ })?;
+ let descriptor = BlobDescriptor::new(
+ file_path.to_string(),
+ descriptor_offset,
+ descriptor_length,
+ );
+ elements.push(Some(Bytes::from(descriptor.serialize())));
+ element_offset += element_length;
+ }
+ }
+ }
+ Ok(BlobReadValue::Array(elements))
+}
+
#[derive(Debug, Clone)]
enum PlannedBlobRead {
Null,
@@ -280,6 +636,19 @@ enum PlannedBlobRead {
Read(Range<u64>),
}
+#[derive(Debug, Clone)]
+enum PlannedBlobArrayRead {
+ Null,
+ Placeholder,
+ Read(Range<u64>),
+}
+
+#[derive(Debug)]
+struct BlobArrayMetadata {
+ element_data_range: Range<u64>,
+ element_lengths: Vec<Option<u64>>,
+}
+
#[derive(Debug, Clone)]
struct BlobFileIndex {
entries: Vec<BlobEntry>,
@@ -859,11 +1228,12 @@ fn encode_varint(value: i64, out: &mut Vec<u8>) {
mod tests {
use super::*;
use crate::btree::test_util::BytesFileRead;
- use crate::spec::BlobType;
+ use crate::spec::{ArrayType, BlobType};
use arrow_array::Array;
use bytes::Bytes;
use futures::TryStreamExt;
use std::sync::atomic::{AtomicUsize, Ordering};
+ use std::sync::Mutex;
use std::time::Duration;
#[allow(dead_code)]
@@ -929,6 +1299,195 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn test_blob_array_reader_reads_java_fixture() {
+ let read_fields = vec![DataField::new(
+ 0,
+ "payloads".to_string(),
+ DataType::Array(ArrayType::new(DataType::Blob(BlobType::new()))),
+ )];
+ let file_bytes = load_blob_fixture("blob-array.blob");
+
+ let batches = BlobFormatReader::new(String::new(), false)
+ .read_batch_stream(
+ Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))),
+ file_bytes.len() as u64,
+ &read_fields,
+ None,
+ Some(2),
+ None,
+ )
+ .await
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ assert_eq!(batches.len(), 2);
+ assert_eq!(
+ collect_blob_array_values(&batches[0]),
+ vec![
+ Some(vec![Some(b"hello".to_vec()), None,
Some(b"world".to_vec())]),
+ None,
+ ]
+ );
+ assert_eq!(
+ collect_blob_array_values(&batches[1]),
+ vec![None, Some(Vec::new())]
+ );
+
+ let selected = BlobFormatReader::new(String::new(), false)
+ .read_batch_stream(
+ Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))),
+ file_bytes.len() as u64,
+ &read_fields,
+ None,
+ Some(1),
+ Some(vec![RowRange::new(0, 0), RowRange::new(3, 3)]),
+ )
+ .await
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ assert_eq!(selected.len(), 2);
+ assert_eq!(
+ collect_blob_array_values(&selected[0]),
+ vec![Some(vec![
+ Some(b"hello".to_vec()),
+ None,
+ Some(b"world".to_vec()),
+ ])]
+ );
+ assert_eq!(
+ collect_blob_array_values(&selected[1]),
+ vec![Some(Vec::new())]
+ );
+ }
+
+ #[tokio::test]
+ async fn
test_blob_array_reader_builds_exact_descriptors_without_payload_reads() {
+ let file_path = "file:///tmp/blob-array.blob";
+ let file_bytes = load_blob_fixture("blob-array.blob");
+ let reader = TrackingFileRead::new(Bytes::from(file_bytes.clone()));
+
+ let batches = BlobFormatReader::new(file_path.to_string(), true)
+ .read_batch_stream(
+ Box::new(reader.clone()),
+ file_bytes.len() as u64,
+ &blob_array_read_fields(),
+ None,
+ Some(8),
+ None,
+ )
+ .await
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ let rows = collect_blob_array_values(&batches[0]);
+ let first = rows[0].as_ref().unwrap();
+ assert!(first[1].is_none());
+ let hello =
BlobDescriptor::deserialize(first[0].as_ref().unwrap()).unwrap();
+ let world =
BlobDescriptor::deserialize(first[2].as_ref().unwrap()).unwrap();
+ assert_eq!(
+ (hello.uri(), hello.offset(), hello.length()),
+ (file_path, 13, 5)
+ );
+ assert_eq!(
+ (world.uri(), world.offset(), world.length()),
+ (file_path, 18, 5)
+ );
+ assert_eq!(rows[1], None);
+ assert_eq!(rows[2], None);
+ assert_eq!(rows[3], Some(Vec::new()));
+
+ assert!(
+ !reader
+ .ranges()
+ .iter()
+ .any(|range| range.start < 23 && range.end > 13),
+ "descriptor mode must not read element payload bytes"
+ );
+ }
+
+ #[tokio::test]
+ async fn test_blob_array_reader_reads_each_row_element_data_once() {
+ let file_bytes = load_blob_fixture("blob-array.blob");
+ let reader = TrackingFileRead::new(Bytes::from(file_bytes.clone()));
+
+ BlobFormatReader::new(String::new(), false)
+ .read_batch_stream(
+ Box::new(reader.clone()),
+ file_bytes.len() as u64,
+ &blob_array_read_fields(),
+ None,
+ Some(8),
+ None,
+ )
+ .await
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ assert_eq!(
+ reader
+ .ranges()
+ .iter()
+ .filter(|range| **range == (13..23))
+ .count(),
+ 1
+ );
+ assert!(
+ !reader
+ .ranges()
+ .iter()
+ .any(|range| (13..23).contains(&range.start) && *range !=
(13..23)),
+ "inline mode must not issue per-element payload reads"
+ );
+ }
+
+ #[tokio::test]
+ async fn test_blob_array_reader_preserves_order_with_bounded_parallelism()
{
+ let payloads = (0_u8..12)
+ .map(|value| build_blob_array_payload(&[value], &[1]))
+ .collect::<Vec<_>>();
+ let rows = payloads
+ .iter()
+ .map(|payload| Some(payload.as_slice()))
+ .collect::<Vec<_>>();
+ let file_bytes = blob_test_utils::build_blob_file_bytes(&rows);
+ let reader = TrackingFileRead::new(Bytes::from(file_bytes.clone()));
+
+ let batches = BlobFormatReader::new(String::new(), false)
+ .read_batch_stream(
+ Box::new(reader.clone()),
+ file_bytes.len() as u64,
+ &blob_array_read_fields(),
+ None,
+ Some(12),
+ None,
+ )
+ .await
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ assert_eq!(batches.len(), 1);
+ assert_eq!(
+ collect_blob_array_values(&batches[0]),
+ (0_u8..12)
+ .map(|value| Some(vec![Some(vec![value])]))
+ .collect::<Vec<_>>()
+ );
+ assert!(reader.max_in_flight() > 1);
+ assert!(reader.max_in_flight() <= BLOB_READ_CONCURRENCY);
+ }
+
#[tokio::test]
async fn test_blob_reader_treats_java_placeholders_as_null() {
let read_fields = vec![DataField::new(
@@ -1020,6 +1579,22 @@ mod tests {
assert_eq!(generated, load_blob_fixture("blob-placeholder.blob"));
}
+ #[test]
+ fn test_blob_array_fixture_matches_java_writer_layout() {
+ use blob_test_utils::BlobFixtureValue::{Null, Placeholder, Value};
+
+ let first = build_blob_array_payload(b"helloworld", &[5, -1, 5]);
+ let empty = build_blob_array_payload(b"", &[]);
+ let generated = blob_test_utils::build_blob_file_bytes_with_values(&[
+ Value(first.as_slice()),
+ Null,
+ Placeholder,
+ Value(empty.as_slice()),
+ ]);
+
+ assert_eq!(generated, load_blob_fixture("blob-array.blob"));
+ }
+
#[tokio::test]
async fn test_blob_reader_supports_empty_projection() {
let reader = BlobFormatReader::new(String::new(), false);
@@ -1095,7 +1670,143 @@ mod tests {
.await;
assert!(
- matches!(result, Err(Error::DataInvalid { message, .. }) if
message.contains("Blob field"))
+ matches!(result, Err(Error::DataInvalid { message, .. }) if
message.contains("Blob or Array<Blob> field"))
+ );
+ }
+
+ #[tokio::test]
+ async fn test_blob_array_reader_rejects_nested_array() {
+ let file_bytes = load_blob_fixture("blob-array.blob");
+ let read_fields = vec![DataField::new(
+ 0,
+ "payloads".to_string(),
+ DataType::Array(ArrayType::new(DataType::Array(ArrayType::new(
+ DataType::Blob(BlobType::new()),
+ )))),
+ )];
+
+ let result = BlobFormatReader::new(String::new(), false)
+ .read_batch_stream(
+ Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))),
+ file_bytes.len() as u64,
+ &read_fields,
+ None,
+ None,
+ None,
+ )
+ .await;
+
+ assert!(
+ matches!(result, Err(Error::DataInvalid { message, .. }) if
message.contains("Blob or Array<Blob>"))
+ );
+ }
+
+ #[tokio::test]
+ async fn test_blob_array_reader_rejects_invalid_header() {
+ let mut invalid_magic = build_blob_array_payload(b"a", &[1]);
+ invalid_magic[..4].copy_from_slice(&0_i32.to_le_bytes());
+ assert_data_invalid(
+ read_blob_array_payload_error(invalid_magic).await,
+ "magic number",
+ );
+
+ let mut unsupported_version = build_blob_array_payload(b"a", &[1]);
+ unsupported_version[4] = 2;
+ assert!(
+ matches!(read_blob_array_payload_error(unsupported_version).await,
Error::Unsupported { message } if message.contains("payload version"))
+ );
+
+ let mut negative_count = build_blob_array_payload(b"a", &[1]);
+ negative_count[5..9].copy_from_slice(&(-1_i32).to_le_bytes());
+ assert_data_invalid(
+ read_blob_array_payload_error(negative_count).await,
+ "element count",
+ );
+
+ assert_data_invalid(
+ read_blob_array_payload_error(vec![0; BLOB_ARRAY_MIN_PAYLOAD_SIZE
as usize - 1]).await,
+ "too small",
+ );
+ }
+
+ #[tokio::test]
+ async fn test_blob_array_reader_rejects_invalid_index() {
+ let mut negative_length = build_blob_array_payload(b"a", &[1]);
+ set_blob_array_index_length(&mut negative_length, -1);
+ assert_data_invalid(
+ read_blob_array_payload_error(negative_length).await,
+ "index length",
+ );
+
+ let mut oversized_length = build_blob_array_payload(b"a", &[1]);
+ set_blob_array_index_length(&mut oversized_length, 3);
+ assert_data_invalid(
+ read_blob_array_payload_error(oversized_length).await,
+ "index length",
+ );
+
+ let mut count_exceeds_index = build_blob_array_payload(b"a", &[1]);
+ count_exceeds_index[5..9].copy_from_slice(&2_i32.to_le_bytes());
+ assert_data_invalid(
+ read_blob_array_payload_error(count_exceeds_index).await,
+ "count exceeds",
+ );
+
+ let mut truncated_varint = build_blob_array_payload(b"a", &[1]);
+ let index_position = truncated_varint.len() -
BLOB_ARRAY_INDEX_LENGTH_SIZE as usize - 1;
+ truncated_varint[index_position] = 0x80;
+ assert_data_invalid(
+ read_blob_array_payload_error(truncated_varint).await,
+ "element index",
+ );
+
+ let mut count_mismatch = build_blob_array_payload(b"ab", &[1, 1]);
+ count_mismatch[5..9].copy_from_slice(&1_i32.to_le_bytes());
+ assert_data_invalid(
+ read_blob_array_payload_error(count_mismatch).await,
+ "does not match index value count",
+ );
+ }
+
+ #[tokio::test]
+ async fn test_blob_array_reader_rejects_invalid_element_bounds() {
+ assert_data_invalid(
+ read_blob_array_payload_error(build_blob_array_payload(b"",
&[-2])).await,
+ "element length",
+ );
+ assert_data_invalid(
+ read_blob_array_payload_error(build_blob_array_payload(b"a",
&[2])).await,
+ "exceed the payload data length",
+ );
+ assert_data_invalid(
+ read_blob_array_payload_error(build_blob_array_payload(b"ab",
&[1])).await,
+ "do not match the payload data length",
+ );
+ }
+
+ #[tokio::test]
+ async fn test_blob_array_reader_preserves_empty_and_null_elements() {
+ let payload = build_blob_array_payload(b"", &[0, -1]);
+ let file_bytes =
blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice())]);
+
+ let batches = BlobFormatReader::new(String::new(), false)
+ .read_batch_stream(
+ Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))),
+ file_bytes.len() as u64,
+ &blob_array_read_fields(),
+ None,
+ None,
+ None,
+ )
+ .await
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ assert_eq!(
+ collect_blob_array_values(&batches[0]),
+ vec![Some(vec![Some(Vec::new()), None])]
);
}
@@ -1207,6 +1918,58 @@ mod tests {
]
}
+ fn blob_array_read_fields() -> Vec<DataField> {
+ vec![DataField::new(
+ 0,
+ "payloads".to_string(),
+ DataType::Array(ArrayType::new(DataType::Blob(BlobType::new()))),
+ )]
+ }
+
+ fn build_blob_array_payload(element_data: &[u8], element_lengths: &[i64])
-> Vec<u8> {
+ let index = encode_delta_varints_write(element_lengths);
+ let mut payload = Vec::with_capacity(
+ BLOB_ARRAY_MIN_PAYLOAD_SIZE as usize + element_data.len() +
index.len(),
+ );
+ payload.extend_from_slice(&BLOB_ARRAY_MAGIC_NUMBER.to_le_bytes());
+ payload.push(BLOB_ARRAY_VERSION);
+ payload.extend_from_slice(&(element_lengths.len() as
i32).to_le_bytes());
+ payload.extend_from_slice(element_data);
+ payload.extend_from_slice(&index);
+ payload.extend_from_slice(&(index.len() as i32).to_le_bytes());
+ payload
+ }
+
+ fn set_blob_array_index_length(payload: &mut [u8], index_length: i32) {
+ let index_length_position = payload.len() -
BLOB_ARRAY_INDEX_LENGTH_SIZE as usize;
+
payload[index_length_position..].copy_from_slice(&index_length.to_le_bytes());
+ }
+
+ async fn read_blob_array_payload_error(payload: Vec<u8>) -> Error {
+ let file_bytes =
blob_test_utils::build_blob_file_bytes(&[Some(payload.as_slice())]);
+ BlobFormatReader::new(String::new(), false)
+ .read_batch_stream(
+ Box::new(BytesFileRead(Bytes::from(file_bytes.clone()))),
+ file_bytes.len() as u64,
+ &blob_array_read_fields(),
+ None,
+ None,
+ None,
+ )
+ .await
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap_err()
+ }
+
+ fn assert_data_invalid(error: Error, expected_message: &str) {
+ assert!(
+ matches!(error, Error::DataInvalid { message, .. } if
message.contains(expected_message)),
+ "expected DataInvalid containing '{expected_message}'"
+ );
+ }
+
fn collect_binary_values(batch: &RecordBatch) -> Vec<Option<Vec<u8>>> {
let array = batch
.column(0)
@@ -1218,6 +1981,32 @@ mod tests {
.collect()
}
+ fn collect_blob_array_values(batch: &RecordBatch) ->
Vec<Option<Vec<Option<Vec<u8>>>>> {
+ let array = batch
+ .column(0)
+ .as_any()
+ .downcast_ref::<arrow_array::ListArray>()
+ .unwrap();
+ (0..array.len())
+ .map(|row_idx| {
+ if array.is_null(row_idx) {
+ return None;
+ }
+
+ let values = array.value(row_idx);
+ let values = values
+ .as_any()
+ .downcast_ref::<arrow_array::BinaryArray>()
+ .unwrap();
+ Some(
+ (0..values.len())
+ .map(|idx| (!values.is_null(idx)).then(||
values.value(idx).to_vec()))
+ .collect(),
+ )
+ })
+ .collect()
+ }
+
fn load_blob_fixture(name: &str) -> Vec<u8> {
let path = format!("{}/testdata/blob/{name}",
env!("CARGO_MANIFEST_DIR"));
std::fs::read(&path).unwrap_or_else(|e| panic!("Failed to read {path}:
{e}"))
@@ -1228,6 +2017,7 @@ mod tests {
bytes: Bytes,
in_flight: Arc<AtomicUsize>,
max_in_flight: Arc<AtomicUsize>,
+ ranges: Arc<Mutex<Vec<Range<u64>>>>,
}
impl TrackingFileRead {
@@ -1236,17 +2026,23 @@ mod tests {
bytes,
in_flight: Arc::new(AtomicUsize::new(0)),
max_in_flight: Arc::new(AtomicUsize::new(0)),
+ ranges: Arc::new(Mutex::new(Vec::new())),
}
}
fn max_in_flight(&self) -> usize {
self.max_in_flight.load(Ordering::SeqCst)
}
+
+ fn ranges(&self) -> Vec<Range<u64>> {
+ self.ranges.lock().unwrap().clone()
+ }
}
#[async_trait::async_trait]
impl FileRead for TrackingFileRead {
async fn read(&self, range: Range<u64>) -> crate::Result<Bytes> {
+ self.ranges.lock().unwrap().push(range.clone());
let in_flight = self.in_flight.fetch_add(1, Ordering::SeqCst) + 1;
self.max_in_flight.fetch_max(in_flight, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(10)).await;
diff --git a/crates/paimon/src/spec/types.rs b/crates/paimon/src/spec/types.rs
index b1e5ebbb..ba5272fd 100644
--- a/crates/paimon/src/spec/types.rs
+++ b/crates/paimon/src/spec/types.rs
@@ -128,6 +128,15 @@ impl DataType {
matches!(self, DataType::Blob(_))
}
+ /// Returns whether this field is stored in a dedicated `.blob` file.
+ pub(crate) fn is_blob_file_field(&self) -> bool {
+ match self {
+ DataType::Blob(_) => true,
+ DataType::Array(array) => array.element_type().is_blob_type(),
+ _ => false,
+ }
+ }
+
/// Returns whether this type is nullable.
pub fn is_nullable(&self) -> bool {
match self {
@@ -1957,6 +1966,18 @@ mod tests {
use super::*;
use pretty_assertions::assert_eq;
+ #[test]
+ fn test_blob_file_field_classification() {
+ let blob = DataType::Blob(BlobType::new());
+ let array_blob = DataType::Array(ArrayType::new(blob.clone()));
+ let nested_array_blob =
DataType::Array(ArrayType::new(array_blob.clone()));
+
+ assert!(blob.is_blob_file_field());
+ assert!(array_blob.is_blob_file_field());
+ assert!(!nested_array_blob.is_blob_file_field());
+ assert!(!DataType::Int(IntType::new()).is_blob_file_field());
+ }
+
fn load_fixture(name: &str) -> String {
let workdir =
std::env::current_dir().unwrap_or_else(|err| panic!("current_dir
must exist: {err}"));
diff --git a/crates/paimon/src/table/data_evolution_reader.rs
b/crates/paimon/src/table/data_evolution_reader.rs
index ae99a5fc..3756025b 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -1704,7 +1704,7 @@ fn build_source_plan(
let mut column_plan = Vec::with_capacity(read_type.len());
for field in read_type {
- let source_idx = if matches!(field.data_type(), DataType::Blob(_))
+ let source_idx = if field.data_type().is_blob_file_field()
&& !blob_descriptor_fields.contains(field.name())
{
blob_source_indices.get(&field.id()).copied()
@@ -2275,12 +2275,13 @@ mod tests {
use crate::io::FileIOBuilder;
use crate::spec::stats::BinaryTableStats;
use crate::spec::{
- BinaryRow, BlobType, Datum, FloatType, IntType, PredicateBuilder,
Schema, TableSchema,
- VectorType,
+ ArrayType, BinaryRow, BlobType, Datum, FloatType, IntType,
PredicateBuilder, Schema,
+ TableSchema, VectorType,
};
use crate::table::{DataSplitBuilder, DeletionFile, Table, TableRead};
use arrow_array::{
- Array, BinaryArray, FixedSizeListArray, Float32Array, Int32Array,
Int64Array, RecordBatch,
+ Array, BinaryArray, FixedSizeListArray, Float32Array, Int32Array,
Int64Array, ListArray,
+ RecordBatch,
};
use bytes::Bytes;
use futures::TryStreamExt;
@@ -3291,6 +3292,243 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn test_table_read_merges_java_array_blob_file() {
+ let tempdir = tempdir().unwrap();
+ let table_path = local_file_path(tempdir.path());
+ let bucket_dir = tempdir.path().join("bucket-0");
+ fs::create_dir_all(&bucket_dir).unwrap();
+
+ let parquet_path = bucket_dir.join("data.parquet");
+ write_int_parquet_file(&parquet_path, vec![("id", vec![1, 2, 3, 4])],
None);
+
+ let blob_path = bucket_dir.join("payloads.blob");
+ copy_blob_fixture("blob-array.blob", &blob_path);
+
+ let file_io = FileIOBuilder::new("file").build().unwrap();
+ let table_schema = TableSchema::new(
+ 0,
+ &Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column(
+ "payloads",
+
DataType::Array(ArrayType::new(DataType::Blob(BlobType::new()))),
+ )
+ .option("data-evolution.enabled", "true")
+ .build()
+ .unwrap(),
+ );
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "blob_array_t"),
+ table_path,
+ table_schema,
+ None,
+ );
+
+ let split = DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(local_file_path(&bucket_dir))
+ .with_total_buckets(1)
+ .with_data_files(vec![
+ data_file_meta_with_path(
+ "data.parquet",
+ 0,
+ 4,
+ 1,
+ parquet_path.metadata().unwrap().len() as i64,
+ Some(vec!["id"]),
+ ),
+ data_file_meta_with_path(
+ "payloads.blob",
+ 0,
+ 4,
+ 1,
+ blob_path.metadata().unwrap().len() as i64,
+ Some(vec!["payloads"]),
+ ),
+ ])
+ .build()
+ .unwrap();
+
+ let read = TableRead::new(&table, table.schema().fields().to_vec(),
Vec::new());
+ let batches = read
+ .to_arrow(std::slice::from_ref(&split))
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ assert_eq!(collect_int_values(&batches, "id"), vec![1, 2, 3, 4]);
+ assert_eq!(
+ collect_blob_array_values(&batches, "payloads"),
+ vec![
+ Some(vec![Some(b"hello".to_vec()), None,
Some(b"world".to_vec())]),
+ None,
+ None,
+ Some(Vec::new()),
+ ]
+ );
+
+ let descriptor_table = table.copy_with_options(HashMap::from([(
+ "blob-as-descriptor".to_string(),
+ "true".to_string(),
+ )]));
+ let descriptor_read = TableRead::new(
+ &descriptor_table,
+ descriptor_table.schema().fields().to_vec(),
+ Vec::new(),
+ );
+ let descriptor_batches = descriptor_read
+ .to_arrow(std::slice::from_ref(&split))
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ let descriptor_rows = collect_blob_array_values(&descriptor_batches,
"payloads");
+ let first = descriptor_rows[0].as_ref().unwrap();
+ let hello =
BlobDescriptor::deserialize(first[0].as_deref().unwrap()).unwrap();
+ let world =
BlobDescriptor::deserialize(first[2].as_deref().unwrap()).unwrap();
+ assert!(hello.uri().ends_with("payloads.blob"));
+ assert_eq!((hello.offset(), hello.length()), (13, 5));
+ assert_eq!((world.offset(), world.length()), (18, 5));
+ assert_eq!(descriptor_rows[1], None);
+ assert_eq!(descriptor_rows[2], None);
+ assert_eq!(descriptor_rows[3], Some(Vec::new()));
+ }
+
+ #[tokio::test]
+ async fn test_table_read_falls_back_across_array_blob_sequence_groups() {
+ use BlobFixtureValue::{Null, Placeholder, Value};
+
+ let tempdir = tempdir().unwrap();
+ let table_path = local_file_path(tempdir.path());
+ let bucket_dir = tempdir.path().join("bucket-0");
+ fs::create_dir_all(&bucket_dir).unwrap();
+
+ let parquet_path = bucket_dir.join("data.parquet");
+ write_int_parquet_file(&parquet_path, vec![("id", vec![1, 2, 3, 4])],
None);
+
+ let old_0 = build_blob_array_payload(&[Some(&b"old-0"[..]), None]);
+ let old_1 = build_blob_array_payload(&[Some(&b"old-1"[..])]);
+ let old_2 = build_blob_array_payload(&[Some(&b"old-2"[..])]);
+ let empty = build_blob_array_payload(&[]);
+ let new_2 = build_blob_array_payload(&[Some(&b"new-2"[..])]);
+
+ let base_path = bucket_dir.join("array-base.blob");
+ write_blob_file_with_values(
+ &base_path,
+ &[Value(&old_0), Value(&old_1), Value(&old_2), Value(&empty)],
+ );
+ let latest_path = bucket_dir.join("array-latest.blob");
+ write_blob_file_with_values(
+ &latest_path,
+ &[Placeholder, Null, Value(&new_2), Placeholder],
+ );
+
+ let file_io = FileIOBuilder::new("file").build().unwrap();
+ let table_schema = TableSchema::new(
+ 0,
+ &Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column(
+ "payloads",
+
DataType::Array(ArrayType::new(DataType::Blob(BlobType::new()))),
+ )
+ .option("data-evolution.enabled", "true")
+ .build()
+ .unwrap(),
+ );
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "blob_array_fallback_t"),
+ table_path,
+ table_schema,
+ None,
+ );
+
+ let split = DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(local_file_path(&bucket_dir))
+ .with_total_buckets(1)
+ .with_data_files(vec![
+ data_file_meta_with_path(
+ "data.parquet",
+ 0,
+ 4,
+ 1,
+ parquet_path.metadata().unwrap().len() as i64,
+ Some(vec!["id"]),
+ ),
+ data_file_meta_with_path(
+ "array-base.blob",
+ 0,
+ 4,
+ 1,
+ base_path.metadata().unwrap().len() as i64,
+ Some(vec!["payloads"]),
+ ),
+ data_file_meta_with_path(
+ "array-latest.blob",
+ 0,
+ 4,
+ 2,
+ latest_path.metadata().unwrap().len() as i64,
+ Some(vec!["payloads"]),
+ ),
+ ])
+ .build()
+ .unwrap();
+
+ let read = TableRead::new(&table, table.schema().fields().to_vec(),
Vec::new());
+ let batches = read
+ .to_arrow(std::slice::from_ref(&split))
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ assert_eq!(
+ collect_blob_array_values(&batches, "payloads"),
+ vec![
+ Some(vec![Some(b"old-0".to_vec()), None]),
+ None,
+ Some(vec![Some(b"new-2".to_vec())]),
+ Some(Vec::new()),
+ ]
+ );
+
+ let descriptor_table = table.copy_with_options(HashMap::from([(
+ "blob-as-descriptor".to_string(),
+ "true".to_string(),
+ )]));
+ let descriptor_read = TableRead::new(
+ &descriptor_table,
+ descriptor_table.schema().fields().to_vec(),
+ Vec::new(),
+ );
+ let descriptor_batches = descriptor_read
+ .to_arrow(std::slice::from_ref(&split))
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ let descriptor_rows = collect_blob_array_values(&descriptor_batches,
"payloads");
+ let old_0 = descriptor_rows[0].as_ref().unwrap();
+ let old_0 =
BlobDescriptor::deserialize(old_0[0].as_deref().unwrap()).unwrap();
+ let new_2 = descriptor_rows[2].as_ref().unwrap();
+ let new_2 =
BlobDescriptor::deserialize(new_2[0].as_deref().unwrap()).unwrap();
+ assert!(old_0.uri().ends_with("array-base.blob"));
+ assert_eq!(descriptor_rows[0].as_ref().unwrap()[1], None);
+ assert_eq!(descriptor_rows[1], None);
+ assert!(new_2.uri().ends_with("array-latest.blob"));
+ assert_eq!(descriptor_rows[3], Some(Vec::new()));
+ }
+
#[tokio::test]
async fn test_table_read_falls_back_across_blob_sequence_groups() {
use BlobFixtureValue::{Null, Placeholder, Value};
@@ -5039,6 +5277,60 @@ mod tests {
.collect()
}
+ fn collect_blob_array_values(
+ batches: &[RecordBatch],
+ column_name: &str,
+ ) -> Vec<Option<Vec<Option<Vec<u8>>>>> {
+ batches
+ .iter()
+ .flat_map(|batch| {
+ let idx = batch.schema().index_of(column_name).unwrap();
+ let array = batch
+ .column(idx)
+ .as_any()
+ .downcast_ref::<ListArray>()
+ .unwrap();
+ (0..array.len())
+ .map(|row| {
+ if array.is_null(row) {
+ return None;
+ }
+ let values = array.value(row);
+ let values =
values.as_any().downcast_ref::<BinaryArray>().unwrap();
+ Some(
+ (0..values.len())
+ .map(|idx| {
+ (!values.is_null(idx)).then(||
values.value(idx).to_vec())
+ })
+ .collect(),
+ )
+ })
+ .collect::<Vec<_>>()
+ })
+ .collect()
+ }
+
+ fn build_blob_array_payload(elements: &[Option<&[u8]>]) -> Vec<u8> {
+ const ARRAY_MAGIC_NUMBER: i32 = 1094861634;
+ const ARRAY_VERSION: u8 = 1;
+
+ let lengths = elements
+ .iter()
+ .map(|element| element.map_or(-1, |bytes| bytes.len() as i64))
+ .collect::<Vec<_>>();
+ let index = blob_test_utils::encode_delta_varints(&lengths);
+ let mut payload = Vec::new();
+ payload.extend_from_slice(&ARRAY_MAGIC_NUMBER.to_le_bytes());
+ payload.push(ARRAY_VERSION);
+ payload.extend_from_slice(&(elements.len() as i32).to_le_bytes());
+ for bytes in elements.iter().flatten() {
+ payload.extend_from_slice(bytes);
+ }
+ payload.extend_from_slice(&index);
+ payload.extend_from_slice(&(index.len() as i32).to_le_bytes());
+ payload
+ }
+
fn two_col_evolution_table(table_path: String) -> Table {
let file_io = FileIOBuilder::new("file").build().unwrap();
let table_schema = TableSchema::new(
diff --git a/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs
b/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs
index 105ce4e4..4f778b24 100644
--- a/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs
+++ b/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs
@@ -20,12 +20,13 @@ use super::{
DeletionVectorContext,
};
use crate::arrow::build_target_arrow_schema;
-use crate::arrow::format::blob::{BlobReadValue, IndexedBlobReader};
+use crate::arrow::format::blob::{
+ build_blob_array_batch, build_blob_batch, BlobReadValue, IndexedBlobReader,
+};
use crate::io::FileIO;
-use crate::spec::DataField;
+use crate::spec::{DataField, DataType};
use crate::table::{ArrowRecordBatchStream, RowRange};
use crate::{DataSplit, Error};
-use arrow_array::builder::BinaryBuilder;
use arrow_array::RecordBatch;
use async_stream::try_stream;
use futures::StreamExt;
@@ -49,6 +50,7 @@ impl LazyBlobFile {
positions: &[usize],
file_io: &FileIO,
blob_as_descriptor: bool,
+ array_field: bool,
) -> crate::Result<Vec<BlobReadValue>> {
if self.reader.is_none() {
let file_size = u64::try_from(self.file_size).map_err(|e|
Error::DataInvalid {
@@ -88,11 +90,15 @@ impl LazyBlobFile {
self.reader = Some(reader);
}
- self.reader
+ let reader = self
+ .reader
.as_ref()
- .expect("blob reader is initialized above")
- .read_positions(positions)
- .await
+ .expect("blob reader is initialized above");
+ if array_field {
+ reader.read_array_positions(positions).await
+ } else {
+ reader.read_positions(positions).await
+ }
}
fn release_reader(&mut self) {
@@ -109,14 +115,15 @@ pub(super) fn read(
blob_as_descriptor: bool,
anchor_deletion_vector: Option<DeletionVectorContext>,
) -> crate::Result<ArrowRecordBatchStream> {
- if read_fields.len() != 1 || !read_fields[0].data_type().is_blob_type() {
+ if read_fields.len() != 1 ||
!read_fields[0].data_type().is_blob_file_field() {
return Err(Error::DataInvalid {
- message: "Blob bunch should provide exactly one BLOB
field".to_string(),
+ message: "Blob bunch should provide exactly one BLOB or
ARRAY<BLOB> field".to_string(),
source: None,
});
}
let target_schema = build_target_arrow_schema(&read_fields)?;
+ let array_field = matches!(read_fields[0].data_type(), DataType::Array(_));
let split = split.clone();
Ok(try_stream! {
@@ -162,6 +169,7 @@ pub(super) fn read(
target_schema.clone(),
&file_io,
blob_as_descriptor,
+ array_field,
).await?;
}
}
@@ -174,6 +182,7 @@ async fn resolve_batch(
target_schema: Arc<arrow_schema::Schema>,
file_io: &FileIO,
blob_as_descriptor: bool,
+ array_field: bool,
) -> crate::Result<RecordBatch> {
let mut resolved = (0..row_ids.len())
.map(|_| BlobReadValue::Placeholder)
@@ -221,7 +230,7 @@ async fn resolve_batch(
if !file_positions.is_empty() {
let values = file
- .read_positions(&file_positions, file_io,
blob_as_descriptor)
+ .read_positions(&file_positions, file_io,
blob_as_descriptor, array_field)
.await?;
for (output_position, value) in
output_positions.into_iter().zip(values) {
if !matches!(&value, BlobReadValue::Placeholder) {
@@ -250,19 +259,11 @@ async fn resolve_batch(
}
}
- let mut builder = BinaryBuilder::new();
- for value in resolved {
- match value {
- BlobReadValue::Value(bytes) => builder.append_value(bytes),
- BlobReadValue::Null | BlobReadValue::Placeholder =>
builder.append_null(),
- }
+ if array_field {
+ build_blob_array_batch(&target_schema, resolved)
+ } else {
+ build_blob_batch(&target_schema, resolved)
}
- RecordBatch::try_new(target_schema,
vec![Arc::new(builder.finish())]).map_err(|e| {
- Error::UnexpectedError {
- message: format!("Failed to build blob fallback RecordBatch: {e}"),
- source: Some(Box::new(e)),
- }
- })
}
struct RowIdBatchCursor {
@@ -411,7 +412,7 @@ mod tests {
VecDeque::from([oldest]),
];
let file_io = crate::io::FileIOBuilder::new("file").build().unwrap();
- let batch = resolve_batch(&mut groups, &[0, 1, 2, 3], schema,
&file_io, false)
+ let batch = resolve_batch(&mut groups, &[0, 1, 2, 3], schema,
&file_io, false, false)
.await
.unwrap();
let values = batch
diff --git a/crates/paimon/testdata/blob/blob-array.blob
b/crates/paimon/testdata/blob/blob-array.blob
new file mode 100644
index 00000000..e35cfc5e
Binary files /dev/null and b/crates/paimon/testdata/blob/blob-array.blob differ