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 f5ee1c01 [rust] Reduce I/O requests for MAP BLOB descriptor reads
(#920)
f5ee1c01 is described below
commit f5ee1c01673d30cd56dba3037a54a6d3fa89be84
Author: XiaoHongbo <[email protected]>
AuthorDate: Wed Sep 23 23:30:02 2026 +0800
[rust] Reduce I/O requests for MAP BLOB descriptor reads (#920)
---
crates/paimon/src/arrow/format/blob.rs | 76 ++++++++++++++++++++++++++++------
1 file changed, 64 insertions(+), 12 deletions(-)
diff --git a/crates/paimon/src/arrow/format/blob.rs
b/crates/paimon/src/arrow/format/blob.rs
index a08deb81..8c91cbf1 100644
--- a/crates/paimon/src/arrow/format/blob.rs
+++ b/crates/paimon/src/arrow/format/blob.rs
@@ -1392,6 +1392,7 @@ async fn read_blob_map_entry(
});
}
+ let index_lengths_start = payload_range.end - BLOB_MAP_INDEX_LENGTHS_SIZE;
let header = read_blob_map_range(
reader,
payload_range.start..payload_range.start + BLOB_MAP_HEADER_SIZE,
@@ -1424,7 +1425,6 @@ async fn read_blob_map_entry(
}
let entry_count = entry_count as usize;
- let index_lengths_start = payload_range.end - BLOB_MAP_INDEX_LENGTHS_SIZE;
let index_lengths = read_blob_map_range(
reader,
index_lengths_start..payload_range.end,
@@ -1460,19 +1460,20 @@ async fn read_blob_map_entry(
let value_index_start = index_lengths_start - value_index_length;
let key_index_start = value_index_start - key_index_length;
- let key_index =
- read_blob_map_range(reader, key_index_start..value_index_start, "key
index").await?;
- let value_index = read_blob_map_range(
+ let indexes = read_blob_map_range(
reader,
- value_index_start..index_lengths_start,
- "value index",
+ key_index_start..index_lengths_start,
+ "key/value indexes",
)
.await?;
- let key_lengths = decode_delta_varints(&key_index).map_err(|e|
Error::DataInvalid {
+ let key_index_length = key_index_length as usize;
+ let key_index = &indexes[..key_index_length];
+ let value_index = &indexes[key_index_length..];
+ let key_lengths = decode_delta_varints(key_index).map_err(|e|
Error::DataInvalid {
message: format!("Invalid MAP<X, BLOB> key index: {e}"),
source: Some(Box::new(e)),
})?;
- let value_lengths = decode_delta_varints(&value_index).map_err(|e|
Error::DataInvalid {
+ let value_lengths = decode_delta_varints(value_index).map_err(|e|
Error::DataInvalid {
message: format!("Invalid MAP<X, BLOB> value index: {e}"),
source: Some(Box::new(e)),
})?;
@@ -2499,6 +2500,56 @@ mod tests {
}
}
+ #[tokio::test]
+ async fn test_blob_map_descriptor_merges_adjacent_indexes() {
+ let payload =
+ build_blob_map_payload(&[("first", Some(b"alpha")), ("second",
Some(b"beta"))]);
+ let reader = TrackingFileRead::new(Bytes::from(payload.clone()));
+ let key_type =
DataType::VarChar(VarCharType::new(VarCharType::MAX_LENGTH).unwrap());
+
+ read_blob_map_entry(&reader, 0..payload.len() as u64, "", true,
&key_type)
+ .await
+ .unwrap();
+
+ assert_eq!(reader.ranges().len(), 4);
+ assert_eq!(reader.max_in_flight(), 1);
+ }
+
+ #[tokio::test]
+ async fn test_blob_map_descriptor_honors_configured_parallelism() {
+ let payloads = ["first", "second", "third"]
+ .into_iter()
+ .map(|key| build_blob_map_payload(&[(key, Some(b"value"))]))
+ .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);
+
+ for parallelism in [1, 2] {
+ let reader =
TrackingFileRead::new(Bytes::from(file_bytes.clone()));
+ let batches = BlobFormatReader::new(String::new(), true)
+ .with_blob_parallelism(parallelism)
+ .read_batch_stream(
+ Box::new(reader.clone()),
+ file_bytes.len() as u64,
+ &blob_map_read_fields(),
+ None,
+ Some(rows.len()),
+ None,
+ )
+ .await
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ assert_eq!(collect_blob_map_values(&batches[0]).len(), rows.len());
+ assert_eq!(reader.max_in_flight(), parallelism);
+ }
+ }
+
#[tokio::test]
async fn test_inline_blob_map_reader_rejects_crc_mismatch() {
let payload = build_blob_map_payload(&[("key", Some(b"value"))]);
@@ -2548,7 +2599,7 @@ mod tests {
.await
.unwrap_err();
- assert_eq!(reader.ranges().len(), 4);
+ assert_eq!(reader.ranges().len(), 3);
assert!(!reader.ranges().contains(&blob_entry_range(&payload_range)));
assert_data_invalid(error, "too large");
}
@@ -2563,7 +2614,7 @@ mod tests {
.await
.unwrap_err();
- assert_eq!(reader.ranges().len(), 4);
+ assert_eq!(reader.ranges().len(), 3);
assert!(!reader.ranges().contains(&blob_entry_range(&payload_range)));
assert_data_invalid(error, "fixed-width key length");
}
@@ -3669,14 +3720,15 @@ mod tests {
let mut index_lengths = Vec::with_capacity(BLOB_MAP_INDEX_LENGTHS_SIZE
as usize);
index_lengths.extend_from_slice(&(key_index.len() as
i32).to_le_bytes());
index_lengths.extend_from_slice(&(value_index.len() as
i32).to_le_bytes());
+ let mut indexes = key_index.clone();
+ indexes.extend_from_slice(&value_index);
let reader = SparseFileRead::new(vec![
(
payload_range.start..payload_range.start +
BLOB_MAP_HEADER_SIZE,
Bytes::from(header),
),
(lengths_start..payload_range.end, Bytes::from(index_lengths)),
- (key_index_start..value_index_start, Bytes::from(key_index)),
- (value_index_start..lengths_start, Bytes::from(value_index)),
+ (key_index_start..lengths_start, Bytes::from(indexes)),
(data_start..data_start, Bytes::new()),
]);
(reader, payload_range)