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 403a2b2e perf(vindex): size native batches by active indexes (#709)
403a2b2e is described below

commit 403a2b2e9bfc4ea66cd7e633619f1460efd18bc8
Author: jerry <[email protected]>
AuthorDate: Mon Aug 17 10:09:06 2026 +0800

    perf(vindex): size native batches by active indexes (#709)
---
 crates/paimon/src/table/vector_search_builder.rs |  39 +++-
 crates/paimon/src/vindex/range_reader.rs         |  36 ++++
 crates/paimon/src/vindex/reader.rs               | 238 +++++++++++++++++++----
 3 files changed, 274 insertions(+), 39 deletions(-)

diff --git a/crates/paimon/src/table/vector_search_builder.rs 
b/crates/paimon/src/table/vector_search_builder.rs
index 21c308bf..f5352206 100644
--- a/crates/paimon/src/table/vector_search_builder.rs
+++ b/crates/paimon/src/table/vector_search_builder.rs
@@ -107,7 +107,7 @@ async fn execute_vindex_searches<S: SeekRead + 'static, G: 
Send + 'static>(
     vector_searches: Vec<VectorSearch>,
     source: S,
     file_name: String,
-    shard_concurrency: usize,
+    index_parallelism: usize,
     guard: G,
 ) -> crate::Result<Vec<Option<HashMap<u64, f32>>>> {
     let panic_context = if vector_searches.len() > 1 {
@@ -117,7 +117,7 @@ async fn execute_vindex_searches<S: SeekRead + 'static, G: 
Send + 'static>(
     };
     execute_global_index_with_guard(panic_context, guard, move || {
         let mut reader = VindexVectorGlobalIndexReader::new(io_meta, options)
-            .with_batch_shard_concurrency(shard_concurrency);
+            .with_batch_index_parallelism(index_parallelism);
         reader
             .visit_batch_vector_search(&vector_searches, |_| Ok(source))
             .map_err(|e| crate::Error::DataInvalid {
@@ -135,6 +135,10 @@ fn current_tokio_runtime_handle() -> 
crate::Result<tokio::runtime::Handle> {
     })
 }
 
+fn vindex_index_parallelism(entry_count: usize, max_concurrency: usize) -> 
usize {
+    entry_count.min(max_concurrency).max(1)
+}
+
 pub struct VectorSearchBuilder<'a> {
     table: &'a Table,
     vector_column: Option<String>,
@@ -822,6 +826,16 @@ async fn plan_and_search_pk_candidates_batch(
             source: None,
         }
     })?;
+    let batch_index_parallelism = match backend {
+        VectorIndexBackend::Vindex => vindex_index_parallelism(
+            plan.splits
+                .iter()
+                .map(|split| split.ann_segments.len())
+                .sum(),
+            concurrency,
+        ),
+        VectorIndexBackend::Lumina => 1,
+    };
 
     // Production data-file reader, mirroring 
`table_read.rs::new_data_file_reader`
     // but projecting only the vector column with no predicates.
@@ -907,7 +921,7 @@ async fn plan_and_search_pk_candidates_batch(
                 }
                 (VectorIndexBackend::Vindex, AnnSegmentSource::Vindex(source)) 
=> {
                     let mut reader = 
VindexVectorGlobalIndexReader::new(io_meta, options.clone())
-                        .with_batch_shard_concurrency(concurrency);
+                        .with_batch_index_parallelism(batch_index_parallelism);
                     reader.load_validated(
                         |_| Ok(source),
                         |metadata| {
@@ -1524,6 +1538,13 @@ async fn evaluate_batch_vector_search(
         }
         ensure_global_index_executor_capacity(concurrency);
         let range_read_permits = 
Arc::new(tokio::sync::Semaphore::new(concurrency));
+        let batch_index_parallelism = vindex_index_parallelism(
+            vector_entries
+                .iter()
+                .filter(|entry| 
is_vindex_index_type(&entry.index_file.index_type))
+                .count(),
+            concurrency,
+        );
         let futures: Vec<_> = vector_entries
             .into_iter()
             .map(|entry| {
@@ -1608,7 +1629,7 @@ async fn evaluate_batch_vector_search(
                                         vector_searches,
                                         source,
                                         file_name,
-                                        concurrency,
+                                        batch_index_parallelism,
                                         permit,
                                     )
                                     .await?
@@ -1629,7 +1650,7 @@ async fn evaluate_batch_vector_search(
                                         vector_searches,
                                         Cursor::new(data),
                                         file_name,
-                                        concurrency,
+                                        batch_index_parallelism,
                                         permit,
                                     )
                                     .await?
@@ -3110,6 +3131,14 @@ mod tests {
         VectorSearchMetric::L2.distance_to_score(distance)
     }
 
+    #[test]
+    fn vindex_batch_parallelism_tracks_active_entries() {
+        assert_eq!(vindex_index_parallelism(1, 1), 1);
+        assert_eq!(vindex_index_parallelism(1, 64), 1);
+        assert_eq!(vindex_index_parallelism(8, 4), 4);
+        assert_eq!(vindex_index_parallelism(4, 8), 4);
+    }
+
     fn make_field(id: i32, name: &str) -> DataField {
         DataField::new(id, name.to_string(), DataType::Int(IntType::default()))
     }
diff --git a/crates/paimon/src/vindex/range_reader.rs 
b/crates/paimon/src/vindex/range_reader.rs
index 07431749..d9c5e158 100644
--- a/crates/paimon/src/vindex/range_reader.rs
+++ b/crates/paimon/src/vindex/range_reader.rs
@@ -633,6 +633,42 @@ mod tests {
         assert_eq!(tracking.max_active.load(Ordering::SeqCst), 1);
     }
 
+    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
+    async fn read_many_caps_each_batch_at_32_ranges() {
+        let range_count = RANGE_READ_CONCURRENCY + 1;
+        let stride = RANGE_COALESCE_GAP + 2;
+        let data = Bytes::from(vec![8u8; range_count * stride as usize]);
+        let tracking = Arc::new(ConcurrencyTrackingRead {
+            data: data.clone(),
+            active: AtomicUsize::new(0),
+            max_active: AtomicUsize::new(0),
+        });
+        let source: Arc<dyn FileRead> = tracking.clone();
+        let mut reader = VindexFileReader::new(
+            source,
+            tokio::runtime::Handle::current(),
+            data.len() as u64,
+            "index".to_string(),
+        );
+
+        tokio::task::spawn_blocking(move || {
+            let mut buffers = vec![[0u8; 1]; range_count];
+            let mut requests = buffers
+                .iter_mut()
+                .enumerate()
+                .map(|(index, buffer)| ReadRequest::new(index as u64 * stride, 
buffer))
+                .collect::<Vec<_>>();
+            reader.pread(&mut requests).unwrap();
+        })
+        .await
+        .unwrap();
+
+        assert_eq!(
+            tracking.max_active.load(Ordering::SeqCst),
+            RANGE_READ_CONCURRENCY
+        );
+    }
+
     #[test]
     fn local_fs_read_completes_with_one_host_blocking_thread() {
         let temp_dir = tempfile::tempdir().unwrap();
diff --git a/crates/paimon/src/vindex/reader.rs 
b/crates/paimon/src/vindex/reader.rs
index 21603f17..55112dd8 100644
--- a/crates/paimon/src/vindex/reader.rs
+++ b/crates/paimon/src/vindex/reader.rs
@@ -15,7 +15,6 @@
 // specific language governing permissions and limitations
 // under the License.
 
-use crate::spec::CoreOptions;
 use crate::vector_search::{GlobalIndexIOMeta, VectorSearch};
 use paimon_vindex_core::distance::MetricType;
 use paimon_vindex_core::index::{
@@ -25,10 +24,72 @@ use paimon_vindex_core::io::{ReadRequest, SeekRead, 
SeekReadCapabilities};
 use std::collections::BinaryHeap;
 use std::collections::HashMap;
 use std::io;
+use std::sync::{Condvar, Mutex};
 
 const DEFAULT_NPROBE: usize = 16;
 const NPROBE_PARAMETER: &str = "ivf.nprobe";
-const NATIVE_BATCH_OPERATION_WORKING_SET_BYTES: usize = 64 * 1024 * 1024;
+const NATIVE_BATCH_PROCESS_WORKING_SET_BYTES: usize = 64 * 1024 * 1024;
+// Native searches run on dedicated executor threads, so blocking here does 
not block async I/O.
+static NATIVE_BATCH_MEMORY_POOL: NativeBatchMemoryPool =
+    NativeBatchMemoryPool::new(NATIVE_BATCH_PROCESS_WORKING_SET_BYTES);
+
+struct NativeBatchMemoryPool {
+    capacity: usize,
+    available_bytes: Mutex<usize>,
+    memory_available: Condvar,
+}
+
+impl NativeBatchMemoryPool {
+    const fn new(bytes: usize) -> Self {
+        Self {
+            capacity: bytes,
+            available_bytes: Mutex::new(bytes),
+            memory_available: Condvar::new(),
+        }
+    }
+
+    fn acquire(&self, bytes: usize) -> NativeBatchMemoryPermit<'_> {
+        let bytes = bytes.min(self.capacity);
+        let mut available = self
+            .available_bytes
+            .lock()
+            .unwrap_or_else(|poisoned| poisoned.into_inner());
+        while *available < bytes {
+            available = self
+                .memory_available
+                .wait(available)
+                .unwrap_or_else(|poisoned| poisoned.into_inner());
+        }
+        *available -= bytes;
+        NativeBatchMemoryPermit { pool: self, bytes }
+    }
+}
+
+struct NativeBatchMemoryPermit<'a> {
+    pool: &'a NativeBatchMemoryPool,
+    bytes: usize,
+}
+
+impl Drop for NativeBatchMemoryPermit<'_> {
+    fn drop(&mut self) {
+        let mut available = self
+            .pool
+            .available_bytes
+            .lock()
+            .unwrap_or_else(|poisoned| poisoned.into_inner());
+        *available += self.bytes;
+        drop(available);
+        self.pool.memory_available.notify_all();
+    }
+}
+
+fn acquire_native_batch_memory(bytes: usize) -> 
NativeBatchMemoryPermit<'static> {
+    NATIVE_BATCH_MEMORY_POOL.acquire(bytes)
+}
+
+fn native_batch_memory_reservation(index_parallelism: usize) -> usize {
+    NATIVE_BATCH_PROCESS_WORKING_SET_BYTES / index_parallelism.max(1)
+}
 
 trait ErasedSeekRead: Send {
     fn pread_erased(&mut self, ranges: &mut [ReadRequest<'_>]) -> 
io::Result<()>;
@@ -78,7 +139,7 @@ impl SeekRead for VindexInput {
 pub struct VindexVectorGlobalIndexReader {
     io_meta: GlobalIndexIOMeta,
     options: HashMap<String, String>,
-    batch_shard_concurrency: Option<usize>,
+    batch_index_parallelism: usize,
     reader: Option<VIndexReader<VindexInput>>,
     metadata: Option<VectorIndexMetadata>,
 }
@@ -88,14 +149,14 @@ impl VindexVectorGlobalIndexReader {
         Self {
             io_meta,
             options,
-            batch_shard_concurrency: None,
+            batch_index_parallelism: 1,
             reader: None,
             metadata: None,
         }
     }
 
-    pub(crate) fn with_batch_shard_concurrency(mut self, concurrency: usize) 
-> Self {
-        self.batch_shard_concurrency = Some(concurrency.max(1));
+    pub(crate) fn with_batch_index_parallelism(mut self, parallelism: usize) 
-> Self {
+        self.batch_index_parallelism = parallelism.max(1);
         self
     }
 
@@ -150,10 +211,6 @@ impl VindexVectorGlobalIndexReader {
         &mut self,
         vector_searches: &[VectorSearch],
     ) -> crate::Result<Vec<Option<HashMap<u64, f32>>>> {
-        let shard_concurrency = match self.batch_shard_concurrency {
-            Some(concurrency) => concurrency,
-            None => CoreOptions::new(&self.options).global_index_thread_num()?,
-        };
         let reader = self
             .reader
             .as_mut()
@@ -173,7 +230,7 @@ impl VindexVectorGlobalIndexReader {
             metadata,
             &self.options,
             vector_searches,
-            shard_concurrency,
+            self.batch_index_parallelism,
         )
     }
 
@@ -348,7 +405,7 @@ fn search_batch_vindex(
     metadata: &VectorIndexMetadata,
     options: &HashMap<String, String>,
     vector_searches: &[VectorSearch],
-    shard_concurrency: usize,
+    index_parallelism: usize,
 ) -> crate::Result<Vec<Option<HashMap<u64, f32>>>> {
     let mut results: Vec<Option<HashMap<u64, f32>>> =
         (0..vector_searches.len()).map(|_| None).collect();
@@ -366,7 +423,7 @@ fn search_batch_vindex(
     }
 
     for (prepared, indices) in groups {
-        let chunk_size = native_batch_chunk_size(metadata, &prepared, 
shard_concurrency);
+        let chunk_size = native_batch_chunk_size(metadata, &prepared, 
index_parallelism);
         for indices in indices.chunks(chunk_size) {
             if indices.len() == 1 {
                 let index = indices[0];
@@ -379,6 +436,10 @@ fn search_batch_vindex(
                 continue;
             }
 
+            let reservation =
+                native_batch_chunk_working_set_bytes(metadata, &prepared, 
indices.len());
+            debug_assert!(reservation <= 
native_batch_memory_reservation(index_parallelism));
+            let _memory_permit = acquire_native_batch_memory(reservation);
             let mut queries = Vec::with_capacity(indices.len() * 
metadata.dimension);
             for &index in indices {
                 queries.extend_from_slice(&vector_searches[index].vector);
@@ -431,22 +492,34 @@ fn search_batch_vindex(
 fn native_batch_chunk_size(
     metadata: &VectorIndexMetadata,
     prepared: &PreparedSearch,
-    shard_concurrency: usize,
+    index_parallelism: usize,
 ) -> usize {
-    let per_shard_budget = NATIVE_BATCH_OPERATION_WORKING_SET_BYTES
-        .checked_div(shard_concurrency.max(1))
-        .unwrap_or(0);
-    let filter_bytes = prepared
-        .filter_bytes
-        .as_ref()
-        .map_or(0, |filter| filter.len().saturating_mul(2));
-    let query_budget = per_shard_budget.saturating_sub(filter_bytes);
+    let per_index_budget = native_batch_memory_reservation(index_parallelism);
+    let filter_bytes = native_batch_filter_working_set_bytes(prepared);
+    let query_budget = per_index_budget.saturating_sub(filter_bytes);
     query_budget
         .checked_div(native_batch_query_working_set_bytes(metadata, prepared))
         .unwrap_or(0)
         .max(1)
 }
 
+fn native_batch_chunk_working_set_bytes(
+    metadata: &VectorIndexMetadata,
+    prepared: &PreparedSearch,
+    query_count: usize,
+) -> usize {
+    native_batch_filter_working_set_bytes(prepared).saturating_add(
+        
query_count.saturating_mul(native_batch_query_working_set_bytes(metadata, 
prepared)),
+    )
+}
+
+fn native_batch_filter_working_set_bytes(prepared: &PreparedSearch) -> usize {
+    prepared
+        .filter_bytes
+        .as_ref()
+        .map_or(0, |filter| filter.len().saturating_mul(2))
+}
+
 fn native_batch_query_working_set_bytes(
     metadata: &VectorIndexMetadata,
     prepared: &PreparedSearch,
@@ -669,7 +742,7 @@ mod tests {
             index,
             query_count,
             HashMap::from([(NPROBE_PARAMETER.to_string(), "1".to_string())]),
-            32,
+            1,
         )
         .await
     }
@@ -678,7 +751,7 @@ mod tests {
         index: Bytes,
         query_count: usize,
         options: HashMap<String, String>,
-        shard_concurrency: usize,
+        index_parallelism: usize,
     ) -> (Vec<Option<HashMap<u64, f32>>>, usize) {
         let tracking = TrackingIndexRead::new(index.clone());
         let source: Arc<dyn FileRead> = tracking.clone();
@@ -694,7 +767,7 @@ mod tests {
                 GlobalIndexIOMeta::new("batch.index".to_string(), index.len() 
as u64, Vec::new());
             let searches = vec![query(); query_count];
             let mut reader = VindexVectorGlobalIndexReader::new(io_meta, 
options)
-                .with_batch_shard_concurrency(shard_concurrency);
+                .with_batch_index_parallelism(index_parallelism);
             reader
                 .visit_batch_vector_search(&searches, |_| Ok(source))
                 .unwrap()
@@ -760,23 +833,120 @@ mod tests {
             nprobe: 16,
             filter_bytes: None,
         };
-        let base = native_batch_chunk_size(&base_metadata, &base_prepared, 32);
+        let index_parallelism = 32;
+        let base = native_batch_chunk_size(&base_metadata, &base_prepared, 
index_parallelism);
 
         let mut larger_index = base_metadata.clone();
         larger_index.dimension *= 2;
         larger_index.nlist *= 2;
-        assert!(native_batch_chunk_size(&larger_index, &base_prepared, 32) < 
base);
+        assert!(native_batch_chunk_size(&larger_index, &base_prepared, 
index_parallelism) < base);
 
         let mut larger_top_k = base_prepared.clone();
         larger_top_k.top_k *= 4;
-        assert!(native_batch_chunk_size(&base_metadata, &larger_top_k, 32) < 
base);
+        assert!(native_batch_chunk_size(&base_metadata, &larger_top_k, 
index_parallelism) < base);
 
         let mut pq_metadata = base_metadata.clone();
         pq_metadata.pq_m = Some(64);
         pq_metadata.pq_bits = Some(8);
-        assert!(native_batch_chunk_size(&pq_metadata, &base_prepared, 32) < 
base);
+        assert!(native_batch_chunk_size(&pq_metadata, &base_prepared, 
index_parallelism) < base);
 
-        assert!(native_batch_chunk_size(&base_metadata, &base_prepared, 64) < 
base);
+        let per_index_working_set = 
base.saturating_mul(native_batch_query_working_set_bytes(
+            &base_metadata,
+            &base_prepared,
+        ));
+        assert!(
+            per_index_working_set.saturating_mul(index_parallelism)
+                <= NATIVE_BATCH_PROCESS_WORKING_SET_BYTES
+        );
+        for parallelism in [1, 2, 3, 32, 64] {
+            assert!(
+                
native_batch_memory_reservation(parallelism).saturating_mul(parallelism)
+                    <= NATIVE_BATCH_PROCESS_WORKING_SET_BYTES
+            );
+        }
+    }
+
+    #[test]
+    fn native_batch_chunk_reservation_tracks_actual_chunk() {
+        let metadata = VectorIndexMetadata {
+            index_type: paimon_vindex_core::index::IndexType::IvfFlat,
+            dimension: 128,
+            nlist: 256,
+            metric: MetricType::L2,
+            total_vectors: 8192,
+            pq_m: None,
+            pq_bits: None,
+            rq_bits: None,
+            diskann: None,
+        };
+        let prepared = PreparedSearch {
+            top_k: 10,
+            nprobe: 16,
+            filter_bytes: Some(vec![0; 128]),
+        };
+        let chunk_size = native_batch_chunk_size(&metadata, &prepared, 1);
+        let full_chunk = native_batch_chunk_working_set_bytes(&metadata, 
&prepared, chunk_size);
+        let final_chunk = native_batch_chunk_working_set_bytes(&metadata, 
&prepared, 2);
+
+        assert!(chunk_size > 2);
+        assert!(full_chunk <= native_batch_memory_reservation(1));
+        assert!(final_chunk < full_chunk);
+    }
+
+    #[test]
+    fn native_batch_memory_pool_admits_only_available_bytes() {
+        let pool = NativeBatchMemoryPool::new(64);
+        let large = pool.acquire(48);
+
+        std::thread::scope(|scope| {
+            let pool = &pool;
+            let (fits_tx, fits_rx) = std::sync::mpsc::channel();
+            scope.spawn(move || {
+                let _permit = pool.acquire(16);
+                fits_tx.send(()).unwrap();
+            });
+            fits_rx
+                .recv_timeout(std::time::Duration::from_secs(1))
+                .expect("reservation fitting the available bytes should not 
wait");
+
+            let (blocked_tx, blocked_rx) = std::sync::mpsc::channel();
+            scope.spawn(move || {
+                let _permit = pool.acquire(17);
+                blocked_tx.send(()).unwrap();
+            });
+            assert!(
+                blocked_rx
+                    .recv_timeout(std::time::Duration::from_millis(50))
+                    .is_err(),
+                "reservation exceeding the available bytes should wait"
+            );
+
+            drop(large);
+            blocked_rx
+                .recv_timeout(std::time::Duration::from_secs(1))
+                .expect("waiting reservation should proceed after bytes are 
released");
+        });
+    }
+
+    #[test]
+    fn native_batch_memory_pool_oversized_request_occupies_pool() {
+        let pool = NativeBatchMemoryPool::new(64);
+
+        std::thread::scope(|scope| {
+            let pool = &pool;
+            let (acquired_tx, acquired_rx) = std::sync::mpsc::channel();
+            scope.spawn(move || {
+                let permit = pool.acquire(65);
+                let _ = acquired_tx.send(permit.bytes);
+            });
+            let acquired = 
acquired_rx.recv_timeout(std::time::Duration::from_secs(1));
+            if acquired.is_err() {
+                // Unblock the old behavior so a regression fails instead of 
hanging the test.
+                *pool.available_bytes.lock().unwrap() = 65;
+                pool.memory_available.notify_all();
+            }
+            assert_eq!(acquired.unwrap(), 64);
+        });
     }
 
     #[test]
@@ -1094,7 +1264,7 @@ mod tests {
 
     #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
     async fn homogeneous_batch_chunks_at_working_set_boundary() {
-        let shard_concurrency = 4096;
+        let index_parallelism = 4096;
         let metadata = VectorIndexMetadata {
             index_type: paimon_vindex_core::index::IndexType::IvfFlat,
             dimension: TEST_DIMENSION,
@@ -1113,7 +1283,7 @@ mod tests {
         let prepared = prepare_search(&metadata, &options, &query())
             .unwrap()
             .unwrap();
-        let chunk_size = native_batch_chunk_size(&metadata, &prepared, 
shard_concurrency);
+        let chunk_size = native_batch_chunk_size(&metadata, &prepared, 
index_parallelism);
         assert!(chunk_size > 16);
 
         let index = build_ivf_flat_index();
@@ -1121,11 +1291,11 @@ mod tests {
             index.clone(),
             chunk_size,
             options.clone(),
-            shard_concurrency,
+            index_parallelism,
         )
         .await;
         let (over_results, over_bytes) =
-            tracked_batch_search_with_options(index, chunk_size + 1, options, 
shard_concurrency)
+            tracked_batch_search_with_options(index, chunk_size + 1, options, 
index_parallelism)
                 .await;
 
         assert_eq!(within_results.len(), chunk_size);

Reply via email to