TheR1sing3un opened a new pull request, #9751:
URL: https://github.com/apache/paimon/pull/9751
### Purpose
Generic global-index builds currently read an entire index split into an
Arrow table, convert the whole indexed column and row IDs to Python lists, and
retain the Arrow table while the writer finishes. Large vector shards therefore
require substantial memory before training starts.
Consume Arrow batches incrementally and write each batch before reading the
next. Create the writer lazily for nonempty input, preserve shard ranges,
relative row IDs and null handling, and close both the Arrow reader and its
Python iterator before finishing the index. Explicit iterator cleanup also
closes source readers when writing fails.
This shared path covers vindex and native full-text indexes. Sorted index
builders are unchanged. The change removes whole-shard materialization at the
builder boundary; decoder buffers, allocator caches, training samples and
native index allocations still consume memory.
This PR is based directly on master and is independent of #9750. Both sides
of the benchmarks below use the existing training-sampling implementation.
### Validation
- `python -m pytest pypaimon/tests/global_index_build_test.py -q`: **30
passed**.
- Changed files pass `flake8 --config dev/cfg.ini` and `git diff --check`.
- Tests cover incremental consumption, empty input, null vectors, null row
IDs, range filtering, relative IDs across shard boundaries, native full-text
builds, and cleanup on writer creation, reading, writing and finish failures. A
real table write-failure test verifies temporary-file removal and no snapshot
commit.
- Benchmark inputs produce identical vector and row-ID SHA-256 hashes across
all configurations. All 65,536 vectors are ingested.
- Native IVF-Flat builds produce identical top-10 IDs and distances for 16
fixed queries at both nprobe=4 and nprobe=16, across three baseline and three
streaming runs.
### Benchmark and ablation
A reproducible script is included at
`paimon-python/dev/benchmark_global_index_streaming.py`.
Environment: macOS 26.4.1 arm64, Python 3.9.6, NumPy 2.0.2, PyArrow 19.0.1,
paimon-vindex 0.4.0; CPU execution with OMP_NUM_THREADS=1 and
OPENBLAS_NUM_THREADS=1. Three fresh, sequential processes per configuration;
shuffled order within ingestion and native suites. Dataset preparation runs
separately. Filesystem cache is not controlled. Arrow reported sandbox
restrictions on some sysctl CPU probes in every mode.
Data: a real Paimon data-evolution Parquet table with row tracking, one
index shard, 65,536 vectors of 256 float32 elements (64 MiB raw vectors).
Default read batch size is 1,024.
Ingestion-only ablation uses the real reader and vector writer's temporary
files, replacing native finish with flushing and input verification. Values are
medians of three runs; RSS is absolute process peak, including roughly 150 MiB
of import/runtime overhead. Hashing time is excluded.
| Arrow lifetime | Python row lifetime | Mode | Peak RSS (MiB) | Time (s) |
| --- | --- | --- | ---: | ---: |
| Whole shard | Whole shard | baseline | 942.05 | 3.392 |
| Streaming | Whole shard | arrow-only | 923.39 | 3.320 |
| Whole shard | Per batch | python-only | 299.95 | 3.184 |
| Streaming | Per batch | stream | 232.62 | 3.219 |
The combined change reduces ingestion peak RSS by about **75.3%**. Avoiding
whole-shard Python objects accounts for most of the gain; streaming Arrow
provides an additional reduction.
Scaling: increasing rows from 16,384 to 65,536 changes baseline peak RSS
from 366.28 to 942.05 MiB, and streaming RSS from 188.14 to 232.62 MiB. At
65,536 rows, streaming batch sizes of 256 / 1,024 / 4,096 use 217.84 / 232.62 /
284.48 MiB.
Including native IVF-Flat training and index serialization (nlist=16, 25%
training sample, L2) changes peak RSS from **942.16 to 433.58 MiB** (about
54.0% lower), with median build times of 3.404 and 3.301 seconds. Timing
excludes input hashing, snapshot commit and subsequent query validation. The
small timing difference is not evidence of a stable speedup.
Example reproduction from `paimon-python`:
```sh
PYTHONPATH=. python dev/benchmark_global_index_streaming.py prepare \
--warehouse /tmp/paimon-stream-bench --rows 65536 --dimension 256
OMP_NUM_THREADS=1 OPENBLAS_NUM_THREADS=1 PYTHONPATH=. \
python dev/benchmark_global_index_streaming.py run \
--warehouse /tmp/paimon-stream-bench --mode stream --batch-size 1024 \
--native --output /tmp/stream-native.json
```
For the ingestion ablation, omit `--native` and run each of `baseline`,
`arrow-only`, `python-only` and `stream` in a fresh process. Repeat each
configuration three times.
--
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]