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 ae53bbf  perf(blob): parallelize descriptor range reads (#534)
ae53bbf is described below

commit ae53bbf0ac51bc560fcc55992d899417bb89c960
Author: Jingsong Lee <[email protected]>
AuthorDate: Fri Jul 17 14:41:22 2026 +0800

    perf(blob): parallelize descriptor range reads (#534)
---
 crates/paimon/src/table/blob_resolver.rs         | 409 ++++++++++++++++++++---
 crates/paimon/src/table/data_evolution_reader.rs | 115 ++++++-
 2 files changed, 471 insertions(+), 53 deletions(-)

diff --git a/crates/paimon/src/table/blob_resolver.rs 
b/crates/paimon/src/table/blob_resolver.rs
index d3f5a82..673efb9 100644
--- a/crates/paimon/src/table/blob_resolver.rs
+++ b/crates/paimon/src/table/blob_resolver.rs
@@ -21,10 +21,90 @@ use crate::Result;
 use arrow_array::builder::BinaryBuilder;
 use arrow_array::{Array, BinaryArray};
 use bytes::Bytes;
+use futures::{stream, StreamExt, TryStreamExt};
 use std::collections::HashMap;
+use std::sync::Arc;
+use tokio::sync::{OwnedSemaphorePermit, Semaphore};
 
 const BLOB_RANGE_MERGE_GAP: u64 = 64 * 1024;
 const BLOB_RANGE_MERGE_MAX_SPAN: u64 = 8 * 1024 * 1024;
+pub(crate) const BLOB_DESCRIPTOR_READ_CONCURRENCY: usize = 8;
+const BLOB_DESCRIPTOR_READ_BYTE_UNIT: u64 = 1024 * 1024;
+const BLOB_DESCRIPTOR_READ_MAX_IN_FLIGHT_BYTES: u64 = 64 * 1024 * 1024;
+
+/// Shared admission control for external descriptor metadata and range reads.
+///
+/// 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)]
+pub(crate) struct BlobReadLimiter {
+    requests: Arc<Semaphore>,
+    bytes: Arc<Semaphore>,
+    byte_unit: u64,
+    max_byte_permits: u32,
+}
+
+impl BlobReadLimiter {
+    pub(crate) fn new() -> Self {
+        Self::with_limits(
+            BLOB_DESCRIPTOR_READ_CONCURRENCY,
+            BLOB_DESCRIPTOR_READ_MAX_IN_FLIGHT_BYTES,
+            BLOB_DESCRIPTOR_READ_BYTE_UNIT,
+        )
+    }
+
+    fn with_limits(request_limit: usize, byte_budget: u64, byte_unit: u64) -> 
Self {
+        assert!(request_limit > 0);
+        assert!(byte_budget > 0);
+        assert!(byte_unit > 0);
+        let max_byte_permits = byte_budget.div_ceil(byte_unit);
+        assert!(max_byte_permits <= u32::MAX as u64);
+        Self {
+            requests: Arc::new(Semaphore::new(request_limit)),
+            bytes: Arc::new(Semaphore::new(max_byte_permits as usize)),
+            byte_unit,
+            max_byte_permits: max_byte_permits as u32,
+        }
+    }
+
+    async fn acquire_read(
+        &self,
+        length: u64,
+        uri: &str,
+    ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
+        let request_permit = self.acquire_request(uri, "range read").await?;
+        let byte_permits = length
+            .div_ceil(self.byte_unit)
+            .max(1)
+            .min(self.max_byte_permits as u64) as u32;
+        let byte_permit = self
+            .bytes
+            .clone()
+            .acquire_many_owned(byte_permits)
+            .await
+            .map_err(|e| crate::Error::UnexpectedError {
+                message: format!(
+                    "Failed to acquire BlobDescriptor byte permits for URI 
'{uri}': {e}"
+                ),
+                source: Some(Box::new(e)),
+            })?;
+        Ok((request_permit, byte_permit))
+    }
+
+    async fn acquire_request(&self, uri: &str, operation: &str) -> 
Result<OwnedSemaphorePermit> {
+        self.requests
+            .clone()
+            .acquire_owned()
+            .await
+            .map_err(|e| crate::Error::UnexpectedError {
+                message: format!(
+                    "Failed to acquire BlobDescriptor {operation} permit for 
URI '{uri}': {e}"
+                ),
+                source: Some(Box::new(e)),
+            })
+    }
+}
 
 /// For each row in a blob column, if the value is a serialized 
`BlobDescriptor`,
 /// resolve it by reading the actual data from the referenced 
URI+offset+length.
@@ -32,6 +112,7 @@ const BLOB_RANGE_MERGE_MAX_SPAN: u64 = 8 * 1024 * 1024;
 pub(crate) async fn resolve_blob_column(
     col: &BinaryArray,
     file_io: &FileIO,
+    limiter: BlobReadLimiter,
 ) -> Result<BinaryArray> {
     let mut needs_resolve = false;
     for i in 0..col.len() {
@@ -74,9 +155,11 @@ pub(crate) async fn resolve_blob_column(
         }
     }
 
+    let mut read_groups = Vec::with_capacity(requests_by_uri.len());
     for (uri, requests) in requests_by_uri {
         let input = file_io.new_input(&uri)?;
         let file_size = if requests.iter().any(|request| 
request.length.is_none()) {
+            let _metadata_permit = limiter.acquire_request(&uri, 
"metadata").await?;
             input
                 .metadata()
                 .await
@@ -119,58 +202,44 @@ pub(crate) async fn resolve_blob_column(
             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 {
+        let reader: Arc<dyn FileRead> = Arc::new(input.reader().await?);
+        read_groups.push(BlobReadGroup {
+            uri,
+            reader,
+            reads: merge_blob_read_requests(bounded_requests),
+        });
+    }
+
+    for ResolvedMergedBlobRead { merged, data } in 
read_blob_groups(read_groups, limiter).await? {
+        for request in merged.requests {
+            let start = usize::try_from(request.offset - 
merged.start).map_err(|e| {
+                crate::Error::DataInvalid {
                     message: format!(
-                        "Failed to read BlobDescriptor URI '{uri}' range 
{}..{}: {e}",
-                        merged.start, merged.end
+                        "BlobDescriptor slice offset exceeds usize: offset={}, 
merged_start={}",
+                        request.offset, merged.start
                     ),
                     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 {
+            let length =
+                usize::try_from(request.length).map_err(|e| 
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
+                        "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,
-                });
-            }
-            for request in merged.requests {
-                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));
-            }
+            cells[request.row] = 
ResolvedBlobCell::Value(data.slice(start..end));
         }
     }
 
@@ -211,6 +280,79 @@ struct MergedBlobRead {
     requests: Vec<BlobReadRequest>,
 }
 
+struct ResolvedMergedBlobRead {
+    merged: MergedBlobRead,
+    data: Bytes,
+}
+
+struct BlobReadGroup {
+    uri: String,
+    reader: Arc<dyn FileRead>,
+    reads: Vec<MergedBlobRead>,
+}
+
+async fn read_blob_groups(
+    groups: Vec<BlobReadGroup>,
+    limiter: BlobReadLimiter,
+) -> Result<Vec<ResolvedMergedBlobRead>> {
+    let grouped_results: Vec<Vec<ResolvedMergedBlobRead>> =
+        stream::iter(groups)
+            .map(|group| {
+                let limiter = limiter.clone();
+                async move {
+                    read_merged_blob_ranges(&group.uri, group.reader, 
group.reads, limiter).await
+                }
+            })
+            .buffer_unordered(BLOB_DESCRIPTOR_READ_CONCURRENCY)
+            .try_collect()
+            .await?;
+    Ok(grouped_results.into_iter().flatten().collect())
+}
+
+async fn read_merged_blob_ranges(
+    uri: &str,
+    reader: Arc<dyn FileRead>,
+    reads: Vec<MergedBlobRead>,
+    limiter: BlobReadLimiter,
+) -> Result<Vec<ResolvedMergedBlobRead>> {
+    stream::iter(reads)
+        .map(|merged| {
+            let uri = uri.to_string();
+            let reader = reader.clone();
+            let limiter = limiter.clone();
+            async move {
+                let _permits = limiter
+                    .acquire_read(merged.end - merged.start, &uri)
+                    .await?;
+                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,
+                    });
+                }
+                Ok(ResolvedMergedBlobRead { merged, data })
+            }
+        })
+        .buffer_unordered(BLOB_DESCRIPTOR_READ_CONCURRENCY)
+        .try_collect()
+        .await
+}
+
 fn merge_blob_read_requests(mut requests: Vec<BlobReadRequest>) -> 
Vec<MergedBlobRead> {
     if requests.is_empty() {
         return Vec::new();
@@ -252,6 +394,187 @@ fn merge_blob_read_requests(mut requests: 
Vec<BlobReadRequest>) -> Vec<MergedBlo
 mod tests {
     use super::*;
 
+    #[derive(Clone)]
+    struct TrackingFileRead {
+        bytes: Bytes,
+        in_flight: std::sync::Arc<std::sync::atomic::AtomicUsize>,
+        max_in_flight: std::sync::Arc<std::sync::atomic::AtomicUsize>,
+    }
+
+    impl TrackingFileRead {
+        fn new(bytes: Bytes) -> Self {
+            Self {
+                bytes,
+                in_flight: 
std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)),
+                max_in_flight: 
std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)),
+            }
+        }
+
+        fn with_counters(
+            bytes: Bytes,
+            in_flight: std::sync::Arc<std::sync::atomic::AtomicUsize>,
+            max_in_flight: std::sync::Arc<std::sync::atomic::AtomicUsize>,
+        ) -> Self {
+            Self {
+                bytes,
+                in_flight,
+                max_in_flight,
+            }
+        }
+
+        fn max_in_flight(&self) -> usize {
+            self.max_in_flight.load(std::sync::atomic::Ordering::SeqCst)
+        }
+    }
+
+    #[async_trait::async_trait]
+    impl FileRead for TrackingFileRead {
+        async fn read(&self, range: std::ops::Range<u64>) -> 
crate::Result<Bytes> {
+            let in_flight = self
+                .in_flight
+                .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
+                + 1;
+            self.max_in_flight
+                .fetch_max(in_flight, std::sync::atomic::Ordering::SeqCst);
+            tokio::time::sleep(std::time::Duration::from_millis(10)).await;
+            self.in_flight
+                .fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
+            Ok(self.bytes.slice(range.start as usize..range.end as usize))
+        }
+    }
+
+    #[tokio::test]
+    async fn test_blob_range_reads_use_bounded_parallelism() {
+        let reader = 
TrackingFileRead::new(Bytes::from_static(b"abcdefghijkl"));
+        let reads = (0..12)
+            .map(|row| MergedBlobRead {
+                start: row,
+                end: row + 1,
+                requests: vec![BlobReadRequest {
+                    row: row as usize,
+                    offset: row,
+                    length: 1,
+                }],
+            })
+            .collect();
+
+        let results = read_merged_blob_ranges(
+            "memory:/blob.bin",
+            std::sync::Arc::new(reader.clone()),
+            reads,
+            BlobReadLimiter::new(),
+        )
+        .await
+        .unwrap();
+
+        assert_eq!(results.len(), 12);
+        assert!(reader.max_in_flight() > 1);
+        assert!(reader.max_in_flight() <= BLOB_DESCRIPTOR_READ_CONCURRENCY);
+    }
+
+    #[tokio::test]
+    async fn test_blob_range_reads_apply_byte_budget_and_preserve_rows() {
+        let reader = TrackingFileRead::new(Bytes::from_static(b"abcdefgh"));
+        let reads = vec![
+            MergedBlobRead {
+                start: 4,
+                end: 8,
+                requests: vec![BlobReadRequest {
+                    row: 0,
+                    offset: 4,
+                    length: 4,
+                }],
+            },
+            MergedBlobRead {
+                start: 0,
+                end: 4,
+                requests: vec![BlobReadRequest {
+                    row: 1,
+                    offset: 0,
+                    length: 4,
+                }],
+            },
+        ];
+
+        let results = read_merged_blob_ranges(
+            "memory:/blob.bin",
+            std::sync::Arc::new(reader.clone()),
+            reads,
+            BlobReadLimiter::with_limits(8, 4, 1),
+        )
+        .await
+        .unwrap();
+
+        let mut by_row = results
+            .into_iter()
+            .map(|result| (result.merged.requests[0].row, result.data))
+            .collect::<Vec<_>>();
+        by_row.sort_by_key(|(row, _)| *row);
+        assert_eq!(by_row[0], (0, Bytes::from_static(b"efgh")));
+        assert_eq!(by_row[1], (1, Bytes::from_static(b"abcd")));
+        assert_eq!(reader.max_in_flight(), 1);
+    }
+
+    #[tokio::test]
+    async fn test_blob_range_reads_overlap_across_uris() {
+        let in_flight = 
std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
+        let max_in_flight = 
std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
+        let groups = b"ab"
+            .iter()
+            .copied()
+            .enumerate()
+            .map(|(row, value)| BlobReadGroup {
+                uri: format!("memory:/blob-{row}.bin"),
+                reader: std::sync::Arc::new(TrackingFileRead::with_counters(
+                    Bytes::from(vec![value]),
+                    in_flight.clone(),
+                    max_in_flight.clone(),
+                )),
+                reads: vec![MergedBlobRead {
+                    start: 0,
+                    end: 1,
+                    requests: vec![BlobReadRequest {
+                        row,
+                        offset: 0,
+                        length: 1,
+                    }],
+                }],
+            })
+            .collect();
+
+        let results = read_blob_groups(groups, BlobReadLimiter::new())
+            .await
+            .unwrap();
+
+        assert_eq!(results.len(), 2);
+        assert_eq!(max_in_flight.load(std::sync::atomic::Ordering::SeqCst), 2);
+    }
+
+    #[tokio::test]
+    async fn test_blob_read_limiter_is_shared_by_metadata_and_ranges() {
+        let limiter = BlobReadLimiter::with_limits(1, 4, 1);
+        let metadata_permit = limiter
+            .acquire_request("memory:/blob.bin", "metadata")
+            .await
+            .unwrap();
+
+        assert!(tokio::time::timeout(
+            std::time::Duration::from_millis(10),
+            limiter.acquire_read(1, "memory:/blob.bin")
+        )
+        .await
+        .is_err());
+
+        drop(metadata_permit);
+        let _permits = tokio::time::timeout(
+            std::time::Duration::from_secs(1),
+            limiter.acquire_read(1, "memory:/blob.bin"),
+        )
+        .await
+        .unwrap()
+        .unwrap();
+    }
+
     #[test]
     fn test_merge_blob_read_requests_merges_nearby_ranges() {
         let merged = merge_blob_read_requests(vec![
diff --git a/crates/paimon/src/table/data_evolution_reader.rs 
b/crates/paimon/src/table/data_evolution_reader.rs
index 513d4db..1f0e9af 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -15,6 +15,7 @@
 // specific language governing permissions and limitations
 // under the License.
 
+use super::blob_resolver::{BlobReadLimiter, BLOB_DESCRIPTOR_READ_CONCURRENCY};
 use super::data_file_reader::{
     append_null_row_id_column, attach_row_id, expand_selected_row_ids, 
insert_column_at,
     DataFileReader,
@@ -33,9 +34,10 @@ use crate::table::{ArrowRecordBatchStream, RESTEnv, 
RowRange};
 use crate::{DataSplit, Error};
 use arrow_array::{Array, BinaryArray, Int64Array, RecordBatch};
 use async_stream::try_stream;
-use futures::StreamExt;
+use futures::{StreamExt, TryStreamExt};
 use roaring::RoaringBitmap;
 use std::collections::{HashMap, HashSet};
+use std::future::Future;
 use std::sync::Arc;
 
 /// Whether a file name denotes a dedicated vector-store file 
(`*.vector.<format>`).
@@ -106,6 +108,7 @@ pub(crate) struct DataEvolutionReader {
     blob_view_fields: HashSet<String>,
     blob_view_resolve_enabled: bool,
     blob_view_rest_env: Option<RESTEnv>,
+    blob_read_limiter: BlobReadLimiter,
 }
 
 impl DataEvolutionReader {
@@ -165,6 +168,7 @@ impl DataEvolutionReader {
             blob_view_fields,
             blob_view_resolve_enabled,
             blob_view_rest_env,
+            blob_read_limiter: BlobReadLimiter::new(),
         })
     }
 
@@ -411,7 +415,13 @@ impl DataEvolutionReader {
 
         batch = self.resolve_blob_view_columns(batch, blob_view_lookup)?;
         let mut batch = if !self.blob_as_descriptor && 
!descriptor_fields.is_empty() {
-            resolve_descriptor_columns(batch, descriptor_fields, 
&self.file_io).await?
+            resolve_descriptor_columns(
+                batch,
+                descriptor_fields,
+                &self.file_io,
+                &self.blob_read_limiter,
+            )
+            .await?
         } else {
             batch
         };
@@ -726,10 +736,28 @@ async fn resolve_descriptor_columns(
     batch: RecordBatch,
     blob_descriptor_fields: &HashSet<String>,
     file_io: &FileIO,
+    limiter: &BlobReadLimiter,
 ) -> crate::Result<RecordBatch> {
+    resolve_descriptor_columns_with(batch, blob_descriptor_fields, |column| {
+        let file_io = file_io.clone();
+        let limiter = limiter.clone();
+        async move { super::blob_resolver::resolve_blob_column(&column, 
&file_io, limiter).await }
+    })
+    .await
+}
+
+async fn resolve_descriptor_columns_with<F, Fut>(
+    batch: RecordBatch,
+    blob_descriptor_fields: &HashSet<String>,
+    resolve: F,
+) -> crate::Result<RecordBatch>
+where
+    F: Fn(BinaryArray) -> Fut,
+    Fut: Future<Output = crate::Result<BinaryArray>>,
+{
     let schema = batch.schema();
-    let mut columns: Vec<Arc<dyn arrow_array::Array>> = 
Vec::with_capacity(batch.num_columns());
-    let mut changed = false;
+    let mut columns = batch.columns().to_vec();
+    let mut descriptor_columns = Vec::new();
 
     for (idx, field) in schema.fields().iter().enumerate() {
         if blob_descriptor_fields.contains(field.name()) {
@@ -738,19 +766,28 @@ async fn resolve_descriptor_columns(
                 .as_any()
                 .downcast_ref::<arrow_array::BinaryArray>()
             {
-                let resolved = 
super::blob_resolver::resolve_blob_column(bin_col, file_io).await?;
-                columns.push(Arc::new(resolved));
-                changed = true;
-                continue;
+                descriptor_columns.push((idx, bin_col.clone()));
             }
         }
-        columns.push(batch.column(idx).clone());
     }
 
-    if !changed {
+    if descriptor_columns.is_empty() {
         return Ok(batch);
     }
 
+    let resolve = &resolve;
+    let resolved_columns: Vec<(usize, BinaryArray)> = 
futures::stream::iter(descriptor_columns)
+        .map(move |(idx, column)| {
+            let future = resolve(column);
+            async move { future.await.map(|resolved| (idx, resolved)) }
+        })
+        .buffer_unordered(BLOB_DESCRIPTOR_READ_CONCURRENCY)
+        .try_collect()
+        .await?;
+    for (idx, resolved) in resolved_columns {
+        columns[idx] = Arc::new(resolved);
+    }
+
     RecordBatch::try_new(schema, columns).map_err(|e| Error::UnexpectedError {
         message: format!("Failed to rebuild RecordBatch after resolving blob 
descriptors: {e}"),
         source: Some(Box::new(e)),
@@ -2088,6 +2125,64 @@ mod tests {
     use blob_test_utils::write_blob_file;
     use test_utils::{local_file_path, write_int_parquet_file};
 
+    #[tokio::test]
+    async fn test_descriptor_columns_resolve_concurrently_and_preserve_order() 
{
+        let schema = Arc::new(arrow_schema::Schema::new(vec![
+            arrow_schema::Field::new("blob_a", arrow_schema::DataType::Binary, 
true),
+            arrow_schema::Field::new("id", arrow_schema::DataType::Int32, 
false),
+            arrow_schema::Field::new("blob_b", arrow_schema::DataType::Binary, 
true),
+        ]));
+        let batch = RecordBatch::try_new(
+            schema,
+            vec![
+                Arc::new(BinaryArray::from(vec![Some(b"a".as_slice())])),
+                Arc::new(Int32Array::from(vec![7])),
+                Arc::new(BinaryArray::from(vec![Some(b"b".as_slice())])),
+            ],
+        )
+        .unwrap();
+        let fields = HashSet::from(["blob_a".to_string(), 
"blob_b".to_string()]);
+        let in_flight = Arc::new(std::sync::atomic::AtomicUsize::new(0));
+        let max_in_flight = Arc::new(std::sync::atomic::AtomicUsize::new(0));
+
+        let resolved = resolve_descriptor_columns_with(batch, &fields, 
|column| {
+            let in_flight = in_flight.clone();
+            let max_in_flight = max_in_flight.clone();
+            async move {
+                let current = in_flight.fetch_add(1, 
std::sync::atomic::Ordering::SeqCst) + 1;
+                max_in_flight.fetch_max(current, 
std::sync::atomic::Ordering::SeqCst);
+                tokio::time::sleep(std::time::Duration::from_millis(10)).await;
+                in_flight.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
+                Ok(column)
+            }
+        })
+        .await
+        .unwrap();
+
+        assert_eq!(resolved.schema().field(0).name(), "blob_a");
+        assert_eq!(resolved.schema().field(1).name(), "id");
+        assert_eq!(resolved.schema().field(2).name(), "blob_b");
+        assert_eq!(
+            resolved
+                .column(0)
+                .as_any()
+                .downcast_ref::<BinaryArray>()
+                .unwrap()
+                .value(0),
+            b"a"
+        );
+        assert_eq!(
+            resolved
+                .column(2)
+                .as_any()
+                .downcast_ref::<BinaryArray>()
+                .unwrap()
+                .value(0),
+            b"b"
+        );
+        assert_eq!(max_in_flight.load(std::sync::atomic::Ordering::SeqCst), 2);
+    }
+
     #[test]
     fn test_build_source_plan_aggregates_same_key_vector_segments() {
         // Two contiguous vector segments, same key -> ONE VectorBunch, files 
in sorted order.

Reply via email to