JingsongLi commented on code in PR #9315:
URL: https://github.com/apache/paimon/pull/9315#discussion_r3819952494


##########
paimon-python/pypaimon/table/source/vector_search_read.py:
##########
@@ -824,3 +901,115 @@ def _compute_score(query, stored, metric):
     if metric == "inner_product":
         return sum(float(q) * float(s) for q, s in zip(query, stored))
     raise ValueError("Unknown vector search metric: %s" % metric)
+
+
+def _raw_search_vectorized(row_ids, vectors, query_vector, metric, limit,
+                           score_candidates=None):
+    """Vectorized raw search using numpy for batch distance computation."""
+    import numpy as np
+
+    # Filter by score_candidates and null vectors.
+    if score_candidates is not None:
+        candidate_set = set(score_candidates)
+        filtered = [(rid, vec) for rid, vec in zip(row_ids, vectors)
+                    if rid in candidate_set and vec is not None]
+    else:
+        filtered = [(rid, vec) for rid, vec in zip(row_ids, vectors)
+                    if vec is not None]
+
+    if not filtered:
+        return DictBasedScoredIndexResult({})
+
+    filtered_ids, filtered_vecs = zip(*filtered)
+    row_id_array = np.array(filtered_ids, dtype=np.int64)
+    stored_matrix = np.array(
+        [_to_vector_list(v) for v in filtered_vecs], dtype=np.float32)
+    query_np = np.array(
+        _to_vector_list(query_vector) if not isinstance(query_vector, 
np.ndarray)
+        else query_vector, dtype=np.float32)
+
+    return _numpy_topk(row_id_array, stored_matrix, query_np, metric, limit)
+
+
+def _raw_search_from_arrow(arrow_table, vector_column_name, query_vector,
+                           metric, limit, score_candidates=None):
+    """Vectorized raw search directly from Arrow table (avoids Python list 
intermediary)."""
+    import numpy as np
+
+    row_ids_col = arrow_table.column(SpecialFields.ROW_ID.name)
+    vectors_col = arrow_table.column(vector_column_name)
+
+    # Try fast path: fixed-size list → direct numpy reshape.
+    row_id_array = row_ids_col.to_numpy()
+    try:
+        # ChunkedArray has no .values; combine to a single array first.
+        if hasattr(vectors_col, 'combine_chunks'):
+            vectors_arr = vectors_col.combine_chunks()
+        else:
+            vectors_arr = vectors_col
+        flat = vectors_arr.values
+        dim = vectors_arr.type.list_size
+        if dim is not None and flat is not None:
+            stored_matrix = flat.to_numpy(zero_copy_only=False).reshape(-1, 
dim).astype(
+                np.float32)
+        else:
+            stored_matrix = np.array(vectors_col.to_pylist(), dtype=np.float32)
+    except (AttributeError, TypeError, ValueError):
+        stored_matrix = np.array(vectors_col.to_pylist(), dtype=np.float32)

Review Comment:
   Could we filter parent nulls before converting the `ListArray`? 
`ARRAY<FLOAT>` is supported and may be nullable. For a variable-length Arrow 
list such as `[[1.0, 0.0], None]`, this fallback executes `np.array(..., 
dtype=np.float32)` and raises `ValueError` before the null mask below runs, 
whereas the previous loop skipped `stored is None`. Please use the Arrow parent 
validity/offsets (or filter null rows first) before building the matrix.



##########
paimon-python/pypaimon/table/source/vector_search_read.py:
##########
@@ -502,27 +554,40 @@ def _read_indexed(self, splits, query_vector, snapshot):
         index_type = _vector_index_type(splits)
         search_limit = self._indexed_search_limit(index_type)
         pre_filters = self._pre_filters(splits, snapshot)
-        futures = [
-            self._eval(
-                split.row_range_start, split.row_range_end,
-                split.vector_index_files,
-                query_vector,
-                search_limit,
-                None if not pre_filters else pre_filters[i]
-            )
-            for i, split in enumerate(splits)
-        ]
 
-        wait(futures)
-
-        merged_scores = {}
-        for future in futures:
-            split_result = future.result()
-            if split_result is not None:
-                score_getter = split_result.score_getter()
-                for row_id in split_result.results():
-                    if row_id not in merged_scores:
-                        merged_scores[row_id] = score_getter(row_id)
+        max_workers = min(self._index_thread_num, len(splits))
+        if max_workers <= 1:

Review Comment:
   Please validate `global-index.thread-num` as positive before this branch, 
and use `len(splits) == 1` for the single-split fast path. With 
`global-index.thread-num=-1` and multiple splits, `max_workers <= 1` is true, 
so only `splits[0]` is searched and the query silently returns incomplete 
results. Java and Rust fail fast on non-positive values.



-- 
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