This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 3d07e28e17 [python] Stream training samples into native vector index 
trainers (#9758)
3d07e28e17 is described below

commit 3d07e28e17b9802eaae372f24e3ac78c6c2d5c68
Author: chaoyang <[email protected]>
AuthorDate: Mon Sep 14 10:36:33 2026 +0800

    [python] Stream training samples into native vector index trainers (#9758)
---
 paimon-python/README.md                            |  10 ++
 .../vindex/vindex_vector_index_writer.py           |  53 ++++++----
 .../pypaimon/tests/global_index_build_test.py      |  30 +++++-
 .../pypaimon/tests/vindex_training_test.py         | 115 +++++++++++++++++++++
 4 files changed, 185 insertions(+), 23 deletions(-)

diff --git a/paimon-python/README.md b/paimon-python/README.md
index c2e5b41b91..7c2708bb54 100644
--- a/paimon-python/README.md
+++ b/paimon-python/README.md
@@ -304,6 +304,16 @@ budgets. This option controls index I/O, not shard search 
or native compute
 threads.
 
 
+# Native vector index training
+
+The native vector index writer submits training vectors in bounded batches.
+`<index-type>.train.sample-ratio` (or its field-level override) still selects
+the same evenly spaced non-null vectors in the same order. Native training
+receives the final corpus size for automatic IVF sizing. This bounds Python
+training buffers; native training and index construction have their own
+memory requirements.
+
+
 # Vector fallback scoring and refinement
 
 Raw vector fallback and refinement score regular FLOAT vectors in bounded
diff --git 
a/paimon-python/pypaimon/globalindex/vindex/vindex_vector_index_writer.py 
b/paimon-python/pypaimon/globalindex/vindex/vindex_vector_index_writer.py
index 87c334fa1b..96443cca34 100644
--- a/paimon-python/pypaimon/globalindex/vindex/vindex_vector_index_writer.py
+++ b/paimon-python/pypaimon/globalindex/vindex/vindex_vector_index_writer.py
@@ -152,18 +152,8 @@ class VindexVectorIndexWriter:
 
             self._close_temp_files()
             self._file_io.check_or_mkdirs(self._index_path)
-            vectors = np.fromfile(
-                self._vector_temp_path,
-                dtype=np.float32,
-                count=self._vector_count * self._dimension,
-            ).reshape(self._vector_count, self._dimension)
-            training_vectors = _sample_training_vectors(
-                np, vectors, self._train_sample_ratio)
-            training = VectorIndexTrainer.train(
-                self._training_options(), training_vectors)
+            training = self._train(np, VectorIndexTrainer)
             try:
-                del training_vectors
-                del vectors
                 with VectorIndexWriter(training) as writer:
                     self._add_vectors_in_batches(np, writer)
                     with self._file_io.new_output_stream(file_path) as 
output_stream:
@@ -178,6 +168,16 @@ class VindexVectorIndexWriter:
 
         return [ResultEntry(self.file_name, self._row_count, b"{}")]
 
+    def _train(self, np, trainer_type):
+        with open(self._vector_temp_path, "rb") as vector_file:
+            with trainer_type.create(self._training_options()) as trainer:
+                for batch in _iter_training_batches(
+                    np, vector_file, self._vector_count, self._dimension,
+                    self._train_sample_ratio, batch_size=ADD_BATCH_SIZE,
+                ):
+                    trainer.add_training_vectors(batch)
+                return trainer.finish_training()
+
     def _file_path(self) -> str:
         return "%s/%s" % (self._index_path, self.file_name)
 
@@ -389,16 +389,31 @@ def _is_float_type(data_type: DataType) -> bool:
     )
 
 
-def _sample_training_vectors(np, vectors, sample_ratio: float):
-    vector_count = vectors.shape[0]
+def _iter_training_batches(
+    np, vector_file, vector_count: int, dimension: int, sample_ratio: float,
+    batch_size: int = ADD_BATCH_SIZE,
+):
+    """Yield the existing evenly spaced sample using bounded reads and 
buffers."""
     train_count = max(1, min(vector_count, int(math.ceil(
         vector_count * sample_ratio))))
-    if train_count == vector_count:
-        return vectors
-    indexes = (
-        np.arange(train_count, dtype=np.int64) * vector_count // train_count
-    )
-    return np.ascontiguousarray(vectors[indexes])
+    position = 0
+    item_size = np.dtype(np.float32).itemsize
+    while position < train_count:
+        start = position * vector_count // train_count
+        end = min(start + batch_size, vector_count)
+        # First sample position whose source row is at or beyond this block.
+        next_position = min(train_count, (end * train_count + vector_count - 
1) // vector_count)
+        vector_file.seek(start * dimension * item_size)
+        vectors = np.fromfile(
+            vector_file, dtype=np.float32, count=(end - start) * dimension,
+        ).reshape(end - start, dimension)
+        if train_count == vector_count:
+            yield vectors
+        else:
+            indexes = np.arange(position, next_position, dtype=np.int64)
+            indexes = indexes * vector_count // train_count - start
+            yield np.ascontiguousarray(vectors[indexes])
+        position = next_position
 
 
 def _float32_batch_values(np, pa, vectors, dimension):
diff --git a/paimon-python/pypaimon/tests/global_index_build_test.py 
b/paimon-python/pypaimon/tests/global_index_build_test.py
index 1633b5270f..9beac060bb 100644
--- a/paimon-python/pypaimon/tests/global_index_build_test.py
+++ b/paimon-python/pypaimon/tests/global_index_build_test.py
@@ -22,6 +22,7 @@ import os
 import struct
 import sys
 import types
+import tempfile
 from unittest.mock import Mock, patch
 
 import pyarrow as pa
@@ -42,7 +43,7 @@ from 
pypaimon.globalindex.full_text.native_full_text_index_writer import (
 )
 from pypaimon.globalindex.vindex.vindex_vector_index_writer import (
     VindexVectorIndexWriter,
-    _sample_training_vectors,
+    _iter_training_batches,
     native_options,
     train_sample_ratio,
 )
@@ -94,9 +95,26 @@ class _FakeVectorIndexTraining:
 
 class _FakeVectorIndexTrainer:
 
+    def __init__(self, options):
+        self.options = options
+        self.batches = []
+
     @classmethod
-    def train(cls, options, data):
-        return _FakeVectorIndexTraining(options, data)
+    def create(cls, options):
+        return cls(options)
+
+    def add_training_vectors(self, data):
+        self.batches.append(data.copy())
+
+    def finish_training(self):
+        import numpy as np
+        return _FakeVectorIndexTraining(self.options, 
np.concatenate(self.batches))
+
+    def __enter__(self):
+        return self
+
+    def __exit__(self, *args):
+        pass
 
 
 class _FakeVectorIndexWriter:
@@ -1075,7 +1093,11 @@ class GlobalIndexBuildTest(
 
         import numpy as np
         vectors = np.arange(20, dtype=np.float32).reshape(10, 2)
-        sampled = _sample_training_vectors(np, vectors, 0.4)
+        with tempfile.TemporaryFile() as vector_file:
+            vectors.tofile(vector_file)
+            vector_file.flush()
+            sampled = np.concatenate(list(_iter_training_batches(
+                np, vector_file, 10, 2, 0.4, batch_size=3)))
         self.assertEqual(
             [[0.0, 1.0], [4.0, 5.0], [10.0, 11.0], [14.0, 15.0]],
             sampled.tolist(),
diff --git a/paimon-python/pypaimon/tests/vindex_training_test.py 
b/paimon-python/pypaimon/tests/vindex_training_test.py
new file mode 100644
index 0000000000..a238e30296
--- /dev/null
+++ b/paimon-python/pypaimon/tests/vindex_training_test.py
@@ -0,0 +1,115 @@
+# 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 math
+import os
+import tempfile
+import unittest
+from unittest import mock
+
+import numpy as np
+
+from pypaimon.filesystem.local_file_io import LocalFileIO
+from pypaimon.globalindex.vindex.vindex_vector_index_writer import (
+    VindexVectorIndexWriter, _iter_training_batches,
+)
+from pypaimon.schema.data_types import ArrayType, AtomicType
+
+
+class VindexTrainingTest(unittest.TestCase):
+
+    def test_batches_preserve_sample_positions_and_bound_reads(self):
+        vectors = np.arange(10003 * 3, dtype=np.float32).reshape(-1, 3)
+        with tempfile.TemporaryFile() as stream:
+            vectors.tofile(stream)
+            stream.flush()
+            for ratio in (1.0, 0.999, 0.37, 0.01, 1e-8):
+                for batch_size in (1, 17, 1000):
+                    with self.subTest(ratio=ratio, batch_size=batch_size):
+                        with mock.patch.object(np, "fromfile", 
wraps=np.fromfile) as read:
+                            batches = list(_iter_training_batches(
+                                np, stream, len(vectors), 3, ratio, 
batch_size))
+                        count = max(1, math.ceil(len(vectors) * ratio))
+                        indexes = np.arange(count) * len(vectors) // count
+                        np.testing.assert_array_equal(vectors[indexes], 
np.concatenate(batches))
+                        self.assertTrue(all(b.flags.c_contiguous for b in 
batches))
+                        self.assertTrue(all(len(b) <= batch_size for b in 
batches))
+                        self.assertTrue(all(c[1]["count"] <= batch_size * 3
+                                            for c in read.call_args_list))
+
+    def test_training_failure_closes_trainer_and_removes_temp_files(self):
+        with tempfile.TemporaryDirectory() as directory:
+            writer = self._writer(directory, {})
+            writer.write([1.0] * 8, 0)
+            paths = [writer._vector_temp_path, writer._row_id_temp_path]
+            for phase in ("add_training_vectors", "finish_training"):
+                trainer = mock.MagicMock()
+                trainer.__enter__.return_value = trainer
+                getattr(trainer, phase).side_effect = RuntimeError("training 
failed")
+                module = mock.Mock()
+                module.VectorIndexTrainer.create.return_value = trainer
+                with mock.patch.dict("sys.modules", {"paimon_vindex": module}):
+                    with self.assertRaisesRegex(RuntimeError, "training 
failed"):
+                        writer.finish()
+                trainer.__exit__.assert_called_once()
+                self.assertTrue(all(not os.path.exists(path) for path in 
paths))
+                self.assertFalse(os.path.exists(writer._file_path()))
+                writer = self._writer(directory, {})
+                writer.write([1.0] * 8, 0)
+                paths = [writer._vector_temp_path, writer._row_id_temp_path]
+            writer.close()
+
+    def test_native_streamed_build_matches_one_shot(self):
+        try:
+            from paimon_vindex import VectorIndexTrainer, VectorIndexWriter
+        except ImportError:
+            self.skipTest("paimon-vindex is not installed")
+        vectors = np.random.default_rng(42).standard_normal((2049, 
8)).astype(np.float32)
+        with tempfile.TemporaryDirectory() as directory:
+            for ratio in (1.0, 0.37):
+                for index_type in ("ivf-flat", "ivf-pq", "ivf-sq", "ivf-rq", 
"diskann"):
+                    with self.subTest(ratio=ratio, index_type=index_type):
+                        options = {index_type + ".train.sample-ratio": 
str(ratio)}
+                        if index_type != "diskann":
+                            options[index_type + ".nlist"] = "16"
+                        writer = self._writer(directory, options, index_type)
+                        writer.write(None, 0)
+                        for i, vector in enumerate(vectors):
+                            writer.write(vector, i + 1)
+                        reference_path = os.path.join(directory, "reference")
+                        count = math.ceil(len(vectors) * ratio)
+                        sample = vectors[np.arange(count) * len(vectors) // 
count]
+                        with 
VectorIndexTrainer.train(writer._training_options(), sample) as training:
+                            with VectorIndexWriter(training) as native:
+                                native.add_vectors(np.arange(1, len(vectors) + 
1), vectors)
+                                with open(reference_path, "wb") as output:
+                                    native.write(output)
+                        with mock.patch(
+                            
"pypaimon.globalindex.vindex.vindex_vector_index_writer.ADD_BATCH_SIZE", 127
+                        ):
+                            result = writer.finish()
+                        self.assertEqual(1, len(result))
+                        with open(reference_path, "rb") as reference, 
open(writer._file_path(), "rb") as actual:
+                            self.assertEqual(reference.read(), actual.read())
+
+    @staticmethod
+    def _writer(directory, options, index_type="ivf-flat"):
+        options = dict(options)
+        options[index_type + ".dimension"] = "8"
+        return VindexVectorIndexWriter(
+            LocalFileIO(), directory, ArrayType(True, AtomicType("FLOAT")),
+            index_type, options, "embedding")

Reply via email to