diff --git a/paimon-python/README.md b/paimon-python/README.md index ed96878ce459..c2e5b41b91aa 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -302,3 +302,16 @@ only support `seek` and `read` remain serialized. Workers are created lazily and released when the index reader closes; separate readers have separate budgets. This option controls index I/O, not shard search or native compute threads. + + +# Vector fallback scoring and refinement + +Raw vector fallback and refinement score regular FLOAT vectors in bounded +blocks using NumPy. List, large-list and fixed-size-list Arrow arrays are +supported, including slices and multiple chunks. Null or unsupported blocks +use the scalar path. Candidate filters are applied before scoring. + +L2 and cosine retain scalar accumulation order. Inner product retains Python +`sum` semantics, including its behavior on newer Python versions. Existing +Top-K tie-breaking rules are preserved. The same scoring path is used for raw and +refined primary-key vector results. diff --git a/paimon-python/pypaimon/table/source/primary_key_vector_read.py b/paimon-python/pypaimon/table/source/primary_key_vector_read.py index 51159f32167f..e379aaa802db 100644 --- a/paimon-python/pypaimon/table/source/primary_key_vector_read.py +++ b/paimon-python/pypaimon/table/source/primary_key_vector_read.py @@ -25,7 +25,7 @@ from pypaimon.table.source.primary_key_vector_scan import PrimaryKeyVectorScanPlan from pypaimon.table.source.vector_search_read import DataEvolutionVectorRead from pypaimon.table.source.vector_search_read import ( - _check_vector_dimension, _compute_score, _to_vector_list) + _iter_arrow_scores) from pypaimon.read.split import DataSplit from pypaimon.globalindex.indexed_split import IndexedSplit from pypaimon.deletionvectors.deletion_vector import DeletionVector @@ -131,7 +131,7 @@ def reranked_iter(): if key[:3] == (partition, data_split.bucket, data_file_name)) position_iter = iter(positions) for batch in reader.to_arrow([split]).to_batches(): - for stored in batch.column(0).to_pylist(): + for score in _iter_arrow_scores(batch.column(0), self._query_vector, metric): try: row_position = next(position_iter) except StopIteration: @@ -144,14 +144,11 @@ def reranked_iter(): raise ValueError( "Primary-key vector rerank read unexpected position %s." % (key,)) - if stored is None: + if score is None: raise ValueError( "Primary-key vector candidate %s contains a null vector." % (key,)) - stored = _to_vector_list(stored) - _check_vector_dimension(self._query_vector, stored) - yield candidate.with_score(_compute_score( - self._query_vector, stored, metric)) + yield candidate.with_score(score) try: next(position_iter) raise ValueError( @@ -196,21 +193,18 @@ def _raw_candidates(self, plan): if _allowed(split, data_file.file_name, position)) position_iter = iter(positions) for batch in reader.to_arrow([read_split]).to_batches(): - for stored in batch.column(0).to_pylist(): + for score in _iter_arrow_scores(batch.column(0), self._query_vector, metric): try: row_position = next(position_iter) except StopIteration: raise ValueError( "Raw vector read returned an unexpected row.") - if stored is None: + if score is None: continue - stored = _to_vector_list(stored) - _check_vector_dimension(self._query_vector, stored) yield PrimaryKeySearchPosition( _partition_bytes(split.data_split.partition), split.data_split.bucket, data_file.file_name, - row_position, _compute_score( - self._query_vector, stored, metric)) + row_position, score) try: next(position_iter) raise ValueError( diff --git a/paimon-python/pypaimon/table/source/vector_search_read.py b/paimon-python/pypaimon/table/source/vector_search_read.py index 5f5dc6ec7038..b05f6bb8b4fe 100644 --- a/paimon-python/pypaimon/table/source/vector_search_read.py +++ b/paimon-python/pypaimon/table/source/vector_search_read.py @@ -306,21 +306,23 @@ def _read_raw_search(self, raw_row_ranges, pre_filter, query_vector, top_k_heap = [] metric = self._search_metric(index_type) - row_ids = table.column(SpecialFields.ROW_ID.name).to_pylist() - vectors = table.column(self._vector_column.name).to_pylist() - for row_id, stored in zip(row_ids, vectors): - if score_candidates is not None and row_id not in score_candidates: - continue - if stored is None: - continue - stored_vector = _to_vector_list(stored) - _check_vector_dimension(query_vector, stored_vector) - _offer_score( - top_k_heap, - self._limit, - row_id, - _compute_score(query_vector, stored_vector, metric), - ) + block_size = _score_block_size(query_vector) + for start in range(0, table.num_rows, block_size): + block = table.slice(start, block_size) + row_ids = block.column(SpecialFields.ROW_ID.name).to_pylist() + vectors = block.column(self._vector_column.name) + if score_candidates is not None: + positions = [i for i, row_id in enumerate(row_ids) + if row_id in score_candidates] + if not positions: + continue + vectors = vectors.take(positions) + row_ids = [row_ids[i] for i in positions] + for row_id, score in zip( + row_ids, _iter_arrow_scores(vectors, query_vector, metric) + ): + if score is not None: + _offer_score(top_k_heap, self._limit, row_id, score) return _scored_result(top_k_heap) def _read_raw_vectors(self, candidates, include_filter=True, snapshot=None): @@ -368,17 +370,24 @@ def _plan_raw_read(self, raw_row_ranges, include_filter, snapshot=None): def _score_raw_vectors(self, candidates, raw_vectors, query_vector, metric, top_k): top_k_heap = [] + row_ids, vectors = [], [] + + def offer_block(): + for row_id, score in zip(row_ids, _score_rows(vectors, query_vector, metric)): + _offer_score(top_k_heap, top_k, row_id, score) + + block_size = _score_block_size(query_vector) for row_id in candidates: stored_vector = raw_vectors.get(row_id) if stored_vector is None: continue - _check_vector_dimension(query_vector, stored_vector) - _offer_score( - top_k_heap, - top_k, - row_id, - _compute_score(query_vector, stored_vector, metric), - ) + row_ids.append(row_id) + vectors.append(stored_vector) + if len(vectors) == block_size: + offer_block() + row_ids, vectors = [], [] + if vectors: + offer_block() return _scored_result(top_k_heap) def _read_raw_refine_search(self, candidates, query_vector, index_type=None, @@ -646,8 +655,7 @@ def _read_raw_batch_search(self, raw_row_ranges, pre_filter, return [_scored_result(heap) for heap in heaps] table_read, splits = self._plan_raw_read(raw_row_ranges, True, snapshot) - metric = _raw_search_metric( - self._table, self._vector_column, self._options, index_type) + metric = self._search_metric(index_type) workers = min(len(splits), table_read._resolve_parallelism(None, len(splits))) if workers <= 1: return self._score_raw_splits(table_read, splits, metric) @@ -884,6 +892,116 @@ def _normalize_metric(metric): return str(metric).lower().replace("-", "_") +def _score_block_size(query): + # Target 8 MiB per float64 matrix, allowing at least one vector. + return max(1, min(1024, (1 << 20) // max(1, len(query)))) + + +def _iter_arrow_scores(vectors, query, metric): + import numpy as np + import pyarrow as pa + + block_size = _score_block_size(query) + for start in range(0, len(vectors), block_size): + block = vectors.slice(start, block_size) + if isinstance(block, pa.ChunkedArray): + block = block.combine_chunks() + dtype = block.type + scores = None + if (pa.types.is_list(dtype) or pa.types.is_large_list(dtype) or + pa.types.is_fixed_size_list(dtype)) and not block.null_count: + flat = block.flatten() + if pa.types.is_float32(flat.type) and not flat.null_count: + if pa.types.is_fixed_size_list(dtype): + regular = dtype.list_size == len(query) + else: + offsets = block.offsets.to_numpy(zero_copy_only=False) + regular = bool(np.all(np.diff(offsets) == len(query))) + if regular and len(query): + matrix = flat.to_numpy(zero_copy_only=False).astype(np.float64).reshape( + len(block), len(query)) + scores = _compute_scores(query, matrix, metric) + if scores is None: + for vector in block.to_pylist(): + if vector is None: + yield None + else: + vector = _to_vector_list(vector) + _check_vector_dimension(query, vector) + yield _compute_score(query, vector, metric) + else: + yield from scores + + +def _score_rows(vectors, query, metric): + import numpy as np + + scores = None + if vectors and all(vector is not None for vector in vectors): + try: + matrix = np.array(vectors, dtype=np.float64) + except (ValueError, TypeError, OverflowError): + matrix = None + if matrix is not None: + scores = _compute_scores(query, matrix, metric) + if scores is not None: + return scores + result = [] + for vector in vectors: + if vector is None: + result.append(None) + else: + vector = _to_vector_list(vector) + _check_vector_dimension(query, vector) + result.append(_compute_score(query, vector, metric)) + return result + + +def _compute_scores(query, matrix, metric): + """Score an owned float64 matrix while preserving scalar reduction semantics. + + L2/cosine use left-to-right accumulation, as in the scalar loops. Inner + product uses Python's sum on precomputed products, retaining its behavior + across Python versions (including compensated summation on Python 3.12+). + Returning None selects the scalar fallback. + """ + import numpy as np + + try: + query = np.asarray(query, dtype=np.float64) + except (ValueError, TypeError, OverflowError): + return None + if (matrix.ndim != 2 or query.ndim != 1 or not len(query) or + matrix.shape[1] != len(query) or not np.isfinite(matrix).all() or + not np.isfinite(query).all()): + return None + if metric not in ("l2", "cosine", "inner_product"): + return None + with np.errstate(over="ignore", invalid="ignore", divide="ignore"): + if metric == "l2": + np.subtract(query, matrix, out=matrix) + np.square(matrix, out=matrix) + np.add.accumulate(matrix, axis=1, out=matrix) + return (1.0 / (1.0 + matrix[:, -1])).tolist() + if metric == "cosine": + norms = np.square(matrix) + np.add.accumulate(norms, axis=1, out=norms) + stored_norms = norms[:, -1].copy() + del norms + query_norm = 0.0 + for value in query: + value = float(value) + query_norm += value * value + np.multiply(matrix, query, out=matrix) + if metric == "inner_product": + return [sum(row) for row in matrix.tolist()] + matrix[:, 0] += 0.0 # Match a scalar accumulator initialized to +0.0. + np.add.accumulate(matrix, axis=1, out=matrix) + denominators = [query_norm ** 0.5 * float(norm) ** 0.5 for norm in stored_norms] + return [0.0 if denominator == 0 else float(dot) / denominator + for dot, denominator in zip(matrix[:, -1], denominators)] + + def _compute_score(query, stored, metric): if metric == "l2": sum_sq = 0.0 diff --git a/paimon-python/pypaimon/tests/primary_key_global_index_golden_test.py b/paimon-python/pypaimon/tests/primary_key_global_index_golden_test.py index afa55e15ed3d..4981b80128ef 100644 --- a/paimon-python/pypaimon/tests/primary_key_global_index_golden_test.py +++ b/paimon-python/pypaimon/tests/primary_key_global_index_golden_test.py @@ -148,3 +148,50 @@ def test_java_primary_key_raw_only_uses_column_metric(catalog, metric, expected_ assert rows.column("id").to_pylist() == [expected_id] assert result.positions[0].score == _compute_score( query, rows.column("embedding")[0].as_py(), metric) + + +def test_java_primary_key_vector_refinement_matches_scalar(catalog): + from unittest import mock + from pypaimon.table.source import vector_search_read as scoring + + _require_native("paimon_vindex") + table = catalog.get_table("default.test_pk_vector_golden") + + def search(): + return (table.new_vector_search_builder() + .with_vector_column("embedding") + .with_query_vector([1.0, 0.0, 0.0, 0.0]) + .with_option("ivf.refine_factor", "2") + .with_limit(2).execute_local()) + + with mock.patch.object(scoring, "_compute_scores", lambda *args: None): + expected = search() + with mock.patch.object(scoring, "_compute_scores", wraps=scoring._compute_scores) as fast: + actual = search() + assert actual.positions == expected.positions + assert actual.positions + assert fast.call_count > 0 + + +def test_java_primary_key_raw_vectors_match_scalar(catalog): + from dataclasses import replace + from unittest import mock + from pypaimon.table.source import vector_search_read as scoring + from pypaimon.table.source.primary_key_vector_scan import PrimaryKeyVectorScanPlan + + table = catalog.get_table("default.test_pk_vector_golden") + builder = (table.new_vector_search_builder() + .with_vector_column("embedding") + .with_query_vector([1.0, 0.0, 0.0, 0.0]).with_limit(2)) + plan = builder.new_vector_search_scan().scan() + raw_plan = PrimaryKeyVectorScanPlan(plan.snapshot_id, [ + replace(split, payloads=(), uncovered_data_files=tuple( + file.file_name for file in split.data_split.files)) for split in plan.splits()]) + reader = builder.new_vector_search_read() + with mock.patch.object(scoring, "_compute_scores", lambda *args: None): + expected = list(reader._raw_candidates(raw_plan)) + with mock.patch.object(scoring, "_compute_scores", wraps=scoring._compute_scores) as fast: + actual = list(reader._raw_candidates(raw_plan)) + assert actual == expected + assert actual + assert fast.call_count > 0 diff --git a/paimon-python/pypaimon/tests/vector_scoring_test.py b/paimon-python/pypaimon/tests/vector_scoring_test.py new file mode 100644 index 000000000000..b900efd2a5fe --- /dev/null +++ b/paimon-python/pypaimon/tests/vector_scoring_test.py @@ -0,0 +1,121 @@ +# 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. + +import unittest +from unittest import mock + +import numpy as np +import pyarrow as pa + +from pypaimon.table.source.vector_search_read import ( + DataEvolutionVectorRead, _compute_score, _compute_scores, _iter_arrow_scores, + _score_block_size, _score_rows, +) +from pypaimon.table.special_fields import SpecialFields +from pypaimon.tests.vector_search_filter_test import _StubTable, _field +from pypaimon.utils.range import Range +from pypaimon.utils.roaring_bitmap import RoaringBitmap64 + + +class VectorScoringTest(unittest.TestCase): + + def test_scores_exactly_match_scalar_accumulation(self): + rng = np.random.default_rng(42) + for dimension in (1, 3, 128, 384): + vectors = rng.standard_normal((37, dimension)).astype(np.float32) + query = rng.standard_normal(dimension).tolist() + for metric in ("l2", "cosine", "inner_product"): + with self.subTest(dimension=dimension, metric=metric): + expected = [_compute_score(query, row, metric) for row in vectors.tolist()] + actual = _compute_scores(query, vectors.astype(np.float64), metric) + self.assertEqual(expected, actual) + self.assertEqual(np.array(expected).tobytes(), np.array(actual).tobytes()) + + def test_cancellation_zero_norms_and_close_scores(self): + cases = [ + ([1e16, 1.0, -1e16], [[1, 1, 1], [1, 0, 1], [0, 0, 0]]), + ([0, 0, 0], [[1, 2, 3], [0, 0, 0], [-0.0, -0.0, -0.0]]), + ([1, 1e-8, 1e-8], [[1, 1e-8, 0], [1, 0, 1e-8], [1, 0, 0]]), + ([1e150, 1e-150, -1e150], [[1e30, 1e-30, 1e30], [0, 1e-30, 0]]), + ] + for query, values in cases: + vectors = np.array(values, dtype=np.float32) + for metric in ("l2", "cosine", "inner_product"): + expected = [_compute_score(query, row, metric) for row in vectors.tolist()] + self.assertEqual(expected, _compute_scores(query, vectors.astype(np.float64), metric)) + + def test_arrow_slices_chunks_and_list_layouts(self): + values = np.random.default_rng(7).standard_normal((1031, 7)).astype(np.float32).tolist() + query = [1.0] * 7 + for dtype in (pa.list_(pa.float32()), pa.large_list(pa.float32()), pa.list_(pa.float32(), 7)): + array = pa.array(values, type=dtype) + sliced = array.slice(2, 1027) + chunked = pa.chunked_array([sliced.slice(0, 511), sliced.slice(511)]) + for data in (sliced, chunked): + for metric in ("l2", "cosine", "inner_product"): + expected = [_compute_score(query, row, metric) for row in values[2:1029]] + self.assertEqual(expected, list(_iter_arrow_scores(data, query, metric))) + + def test_null_and_unsupported_data_preserve_scalar_behavior(self): + for dtype in (pa.list_(pa.float32()), pa.list_(pa.float64()), pa.list_(pa.int64())): + array = pa.array([[1, 2], None, [0, 0]], type=dtype) + self.assertEqual([1.0, None, 1.0 / 6.0], list(_iter_arrow_scores(array, [1, 2], "l2"))) + scores = _iter_arrow_scores(pa.array([None, [1.0]], type=pa.list_(pa.float32())), [1, 2], "l2") + self.assertIsNone(next(scores)) + with self.assertRaisesRegex(ValueError, "dimension mismatch"): + next(scores) + with self.assertRaisesRegex(ValueError, "dimension mismatch"): + list(_iter_arrow_scores(pa.array([[1.0], [1.0, 2.0]]), [1, 2], "l2")) + with self.assertRaises(TypeError): + list(_iter_arrow_scores(pa.array([[1.0, None]], type=pa.list_(pa.float32())), [1, 2], "l2")) + values = [[float("nan"), 0], [1, 0]] + scores = list(_iter_arrow_scores(pa.array(values, type=pa.list_(pa.float32())), [1, 0], "l2")) + self.assertTrue(np.isnan(scores[0])) + self.assertEqual(1.0, scores[1]) + self.assertEqual([1.0], list(_iter_arrow_scores(pa.array([[]], type=pa.list_(pa.float32())), [], "l2"))) + + def test_raw_search_filters_before_scoring_and_preserves_ties(self): + column = _field(1, "embedding", "FLOAT") + reader = DataEvolutionVectorRead(_StubTable([column], []), 2, column, [1, 0]) + table = pa.table({SpecialFields.ROW_ID.name: [9, 3, 1, 2], + "embedding": pa.array([[1], [1, 0], [1, 0], None], + type=pa.list_(pa.float32()))}) + candidates = RoaringBitmap64() + for row_id in (1, 2, 3): + candidates.add(row_id) + with mock.patch.object(reader, "_read_raw_arrow", return_value=table): + result = reader._read_raw_search([Range(0, 9)], None, [1, 0], score_candidates=candidates) + self.assertEqual([1, 3], result.results().to_list()) + reader._limit = 1 + with mock.patch.object(reader, "_read_raw_arrow", return_value=table): + result = reader._read_raw_search([Range(0, 9)], None, [1, 0], score_candidates=candidates) + self.assertEqual([1], result.results().to_list()) + + def test_refinement_and_block_memory_bound(self): + column = _field(1, "embedding", "FLOAT") + reader = DataEvolutionVectorRead(_StubTable([column], []), 10, column, [1, 0]) + values = {i: [float(i % 7), 0] for i in range(2051)} + original = {row_id: vector[:] for row_id, vector in values.items()} + candidates = RoaringBitmap64() + for row_id in values: + candidates.add(row_id) + result = reader._score_raw_vectors(candidates, values, [1, 0], "inner_product", 10) + expected = sorted(values, key=lambda i: (-values[i][0], i))[:10] + self.assertEqual(sorted(expected), result.results().to_list()) + self.assertEqual(original, values) + self.assertLessEqual(_score_block_size([0] * 4096) * 4096, 1 << 20) + self.assertEqual([0.0], _score_rows([[0, 0]], [1, 0], "cosine"))