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.