jerry-024 commented on code in PR #99:
URL: 
https://github.com/apache/paimon-vector-index/pull/99#discussion_r4003210941


##########
core/src/ivfflat_io.rs:
##########
@@ -763,6 +765,282 @@ impl<R: SeekRead> IVFFlatIndexReader<R> {
         let filter = decode_roaring_filter(roaring_filter_bytes)?;
         self.search_with_filter(query, k, nprobe, Some(&filter))
     }
+
+    /// Distance range search: returns every member of the **probed lists** 
whose
+    /// distance falls in `[lower, upper)`.
+    ///
+    /// "No cap" is a promise about **not truncating**, not about
+    /// **completeness**: how many in-band rows come back depends on how much 
the
+    /// IVF probe covered, and a smaller `nprobe` returns fewer.
+    ///
+    /// How the contract differs from [`search`]: there is **no limit and no
+    /// cap**, aligned with Faiss's `range_search`. Ordering is the caller's
+    /// business
+    /// (it is SQL's job), so results are neither padded nor sorted.
+    ///
+    /// [`search`]: IVFFlatIndexReader::search
+    pub fn range_search(
+        &mut self,
+        query: &[f32],
+        params: VectorRangeSearchParams,
+    ) -> io::Result<RangeSearchResult> {
+        self.range_search_with_filter(query, params, None)
+    }
+
+    pub fn range_search_with_filter(
+        &mut self,
+        query: &[f32],
+        params: VectorRangeSearchParams,
+        filter: Option<&dyn RowIdFilter>,
+    ) -> io::Result<RangeSearchResult> {
+        self.range_search_batch_with_filter(query, 1, params, filter)
+    }
+
+    /// Range search restricted to the ids in a serialized Roaring allow-list.
+    ///
+    /// The filter is decoded **before** the empty-band shortcut, so a 
malformed
+    /// filter is still rejected. Unlike the top-K Roaring path there is no
+    /// filter-cardinality probe: the only consumer of that count is the 
automatic
+    /// width mode, which range search does not support yet.
+    pub fn range_search_with_roaring_filter(
+        &mut self,
+        query: &[f32],
+        params: VectorRangeSearchParams,
+        roaring_filter_bytes: &[u8],
+    ) -> io::Result<RangeSearchResult> {
+        let filter = decode_roaring_filter(roaring_filter_bytes)?;
+        self.range_search_with_filter(query, params, Some(&filter))
+    }
+
+    /// Batched range search. A unique probed list is still read once and 
fanned
+    /// out to every query that selected it.
+    pub fn range_search_batch(
+        &mut self,
+        queries: &[f32],
+        nq: usize,
+        params: VectorRangeSearchParams,
+    ) -> io::Result<RangeSearchResult> {
+        self.range_search_batch_with_filter(queries, nq, params, None)
+    }
+
+    pub fn range_search_batch_with_roaring_filter(
+        &mut self,
+        queries: &[f32],
+        nq: usize,
+        params: VectorRangeSearchParams,
+        roaring_filter_bytes: &[u8],
+    ) -> io::Result<RangeSearchResult> {
+        // As in the single-query case: decode first, so a malformed filter is
+        // rejected even under an empty band.
+        let filter = decode_roaring_filter(roaring_filter_bytes)?;
+        self.range_search_batch_with_filter(queries, nq, params, Some(&filter))
+    }
+
+    /// The batch range engine. The single-query entry point is its `nq = 1`
+    /// wrapper, so both share one scheduling and accounting path.
+    ///
+    /// Going per-list with `read_inverted_lists(&[one])` would degenerate into
+    /// one I/O per list, whereas the top-K single-query path already batches
+    /// reads and then chooses a parallel or serial arm; sharing the batch 
engine
+    /// keeps that behaviour and means the statistics are wired up in exactly 
one
+    /// place.
+    pub fn range_search_batch_with_filter(
+        &mut self,
+        queries: &[f32],
+        nq: usize,
+        params: VectorRangeSearchParams,
+        filter: Option<&dyn RowIdFilter>,
+    ) -> io::Result<RangeSearchResult> {
+        // This reader method is public and can be called directly, so the
+        // enum layer's validation cannot be relied upon here.
+        validate_queries(queries, nq, self.d)?;
+        // The band's metric must match the index's, otherwise the scan would
+        // compute values under the index's metric while deciding membership 
with
+        // another metric's band, silently returning wrong rows.
+        if params.band().metric() != self.metric {
+            return Err(io::Error::new(
+                io::ErrorKind::InvalidInput,
+                format!(
+                    "band metric {:?} does not match index metric {:?}",
+                    params.band().metric(),
+                    self.metric
+                ),
+            ));
+        }
+        let nprobe = params.validate(self.nlist)?;
+        let mut builder = RangeResultBuilder::new(nq);
+        // Only once all of the above has passed may an empty band 
short-circuit;
+        // doing it earlier would let an empty band mask a wrong dimension, a
+        // non-finite query or a metric mismatch. The index is not touched 
here,
+        // so no list is read.
+        if params.band().is_empty() {
+            return Ok(builder.build());
+        }
+        // Cold start: a freshly opened reader has an empty centroid table and
+        // `loaded == false`. The top-K path calls ensure_loaded() first thing;
+        // omitting it here would panic when the coarse quantizer reads the 
empty
+        // centroid array.
+        self.ensure_loaded()?;
+
+        // Note for metric certification: unlike the top-K path this does not
+        // apply `fvec_normalize` for cosine. That is currently unreachable,
+        // because `params.validate` rejects every non-L2 metric above, but
+        // whoever certifies cosine must add the normalization here as well as
+        // relaxing `ensure_certified_metric`.
+        //
+        // Ranked probe selection, one group per query.
+        //
+        // This relies on `find_topk_batch` selecting the same centroids for a
+        // query whether it runs alone or in a batch: range search promises a
+        // query returns the same rows either way, and the probed lists decide
+        // which rows are reachable at all. The helper takes an SGEMM path for
+        // `nq > 1` and the direct kernel for `nq == 1`, but recomputes the
+        // distance of every selected centroid with the direct kernel wherever
+        // its error bound leaves the ranking ambiguous, which is what makes 
the
+        // two agree. `kmeans` pins that property with a test; if it is ever
+        // relaxed, this caller needs an exact helper of its own again.
+        let (probe_lists, _coarse_distances) = kmeans::find_topk_batch(
+            queries,
+            nq,
+            &self.quantizer_centroids,
+            self.nlist,
+            self.d,
+            nprobe,
+        );
+        for (qi, lists) in probe_lists.iter().enumerate() {
+            builder.record_lists_probed(qi, lists.len());
+        }
+        // The list-to-query fan-out table, which preserves "a unique list is
+        // read once". The inner vectors hold *query* indices in ascending 
query
+        // order, not probe ranks -- a query's own rank is its position within
+        // `probe_lists[qi]`. Scheduling that has to act on rank order is the
+        // later caps work's concern, and it will need to replace this loop.
+        let mut list_to_queries: Vec<Vec<usize>> = vec![Vec::new(); 
self.nlist];
+        let mut unique_lists: Vec<usize> = Vec::new();
+        for (qi, lists) in probe_lists.iter().enumerate() {
+            for &list_id in lists {
+                if list_to_queries[list_id].is_empty() {
+                    unique_lists.push(list_id);
+                }
+                list_to_queries[list_id].push(qi);
+            }
+        }
+
+        // Per-query output buckets. The whole batch of (list, query) results 
is
+        // deliberately **not** materialized and merged afterwards: unlike a
+        // top-K local heap, a range collector has no upper bound, so that 
shape
+        // would hold O(B x Q x hits) resident at once. Each (list, query) task
+        // holds only that list's hits and takes the lock once to merge them.
+        let outputs: Vec<Mutex<Vec<(i64, f32)>>> =

Review Comment:
   Minor: Each `(list, query)` merge acquires `tallies[qi]` and `outputs[qi]` 
separately. Could we keep `(rows, scanned, early_abandoned)` behind one 
per-query mutex instead? This removes one mutex vector and one lock acquisition 
per merge without changing the parallelism or result semantics.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to