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