vibex-wang commented on code in PR #9315:
URL: https://github.com/apache/paimon/pull/9315#discussion_r3977680426


##########
paimon-python/pypaimon/table/source/vector_search_read.py:
##########
@@ -824,3 +923,246 @@ 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
+    import pyarrow.compute as pc
+
+    row_ids_col = arrow_table.column(SpecialFields.ROW_ID.name)
+    vectors_col = arrow_table.column(vector_column_name)
+
+    # Filter out null vectors at the Arrow level before conversion.
+    valid_mask = pc.is_valid(vectors_col)
+    if not pc.all(valid_mask).as_py():
+        arrow_table = arrow_table.filter(valid_mask)
+        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)
+
+    query_np = np.asarray(query_vector, dtype=np.float32)
+
+    if stored_matrix.ndim < 2 or stored_matrix.shape[0] == 0:
+        return DictBasedScoredIndexResult({})
+
+    if stored_matrix.shape[1] != query_np.shape[0]:
+        raise ValueError(
+            "Query vector dimension mismatch: expected %d, got %d"
+            % (stored_matrix.shape[1], query_np.shape[0]))
+
+    # Handle null vectors and score_candidates filtering.
+    if score_candidates is not None:
+        candidate_set = set(score_candidates)
+        mask = np.array([rid in candidate_set for rid in row_id_array], 
dtype=bool)
+        # Also mask null vectors (check for any NaN row).
+        null_mask = ~np.isnan(stored_matrix).any(axis=1)
+        mask = mask & null_mask
+        row_id_array = row_id_array[mask]
+        stored_matrix = stored_matrix[mask]
+    else:
+        null_mask = ~np.isnan(stored_matrix).any(axis=1)
+        if not null_mask.all():
+            row_id_array = row_id_array[null_mask]
+            stored_matrix = stored_matrix[null_mask]
+
+    if len(row_id_array) == 0:
+        return DictBasedScoredIndexResult({})
+
+    return _numpy_topk(row_id_array, stored_matrix, query_np, metric, limit)
+
+
+def _numpy_topk(row_id_array, stored_matrix, query_np, metric, limit):
+    """Core numpy distance computation + topK selection."""
+    import numpy as np
+
+    if metric == "l2":
+        diffs = stored_matrix - query_np

Review Comment:
   Fixed.



##########
paimon-python/pypaimon/tests/benchmark_vector_search_standalone.py:
##########
@@ -0,0 +1,247 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+"""
+Standalone benchmark for Phase 0 vector search optimizations.
+No pypaimon imports needed — tests the algorithm directly.
+
+Usage:
+    python3 benchmark_vector_search_standalone.py
+    python3 benchmark_vector_search_standalone.py --num-rows 100000 --dim 768
+"""
+
+import argparse
+import time
+
+import numpy as np
+
+
+# ============================================================
+# Original pure-Python implementation (copied from vector_search_read.py)
+# ============================================================
+
+def _compute_score_python(query, stored, metric):
+    if metric == "l2":
+        sum_sq = 0.0
+        for q, s in zip(query, stored):
+            diff = float(q) - float(s)
+            sum_sq += diff * diff
+        return 1.0 / (1.0 + sum_sq)
+    if metric == "cosine":
+        dot = 0.0
+        norm_a = 0.0
+        norm_b = 0.0
+        for q, s in zip(query, stored):
+            q = float(q)
+            s = float(s)
+            dot += q * s
+            norm_a += q * q
+            norm_b += s * s
+        denominator = (norm_a ** 0.5) * (norm_b ** 0.5)
+        return 0.0 if denominator == 0 else dot / denominator
+    if metric == "inner_product":
+        return sum(float(q) * float(s) for q, s in zip(query, stored))
+    raise ValueError("Unknown metric: %s" % metric)
+
+
+def raw_search_python(row_ids, vectors, query_vector, metric, limit):
+    """Original pure-Python raw search with heap."""
+    import heapq
+    top_k_heap = []
+    for row_id, stored in zip(row_ids, vectors):
+        if stored is None:
+            continue
+        score = _compute_score_python(query_vector, stored, metric)
+        entry = (score, -row_id, row_id)
+        if len(top_k_heap) < limit:
+            heapq.heappush(top_k_heap, entry)
+        elif entry[:2] > top_k_heap[0][:2]:
+            heapq.heapreplace(top_k_heap, entry)
+    return {row_id: score for score, _, row_id in top_k_heap}
+
+
+# ============================================================
+# New numpy-vectorized implementation
+# ============================================================
+
+def raw_search_numpy(row_ids_list, vectors_list, query_vector, metric, limit):
+    """Numpy-vectorized raw search."""
+    # Filter nulls.
+    filtered = [(rid, vec) for rid, vec in zip(row_ids_list, vectors_list)
+                if vec is not None]
+    if not filtered:
+        return {}
+
+    filtered_ids, filtered_vecs = zip(*filtered)
+    row_id_array = np.array(filtered_ids, dtype=np.int64)
+    stored_matrix = np.array(filtered_vecs, dtype=np.float32)
+    query_np = np.asarray(query_vector, dtype=np.float32)
+
+    return _numpy_distance_topk(row_id_array, stored_matrix, query_np, metric, 
limit)
+
+
+def raw_search_numpy_fast(row_id_array, stored_matrix, query_np, metric, 
limit):
+    """Numpy fast path: data already in numpy arrays (simulates Arrow 
zero-copy)."""
+    return _numpy_distance_topk(row_id_array, stored_matrix, query_np, metric, 
limit)
+
+
+def _numpy_distance_topk(row_id_array, stored_matrix, query_np, metric, limit):

Review Comment:
   Fixed.



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