Repository navigation
[python] Build global index shards concurrently - #9798
Conversation
42134f1 to
da87cce
Compare
JingsongLi
left a comment
There was a problem hiding this comment.
Requirement fit: SUPPORTED. Implementation: FINDINGS.
Reviewed da87cce08afa. Building independent global-index shards concurrently has a real caller-visible throughput path with a conservative default. Normal worker failure cleanup and deterministic manifest ordering are covered, but output ownership is lost when task submission itself fails.
Validation: 46 tests and 32 subtests passed. A separate reproduction using the production scheduler, a real ThreadPoolExecutor, temporary output files and an injected second-worker startup failure left both uncommitted shard outputs behind. Current CI is green. Native throughput was not measured locally.
The parallel build assigned `futures` only after the whole comprehension succeeded, so a submission that raised (for example `RuntimeError: can't start new thread`) left it as the empty list. Shards accepted before the failure still ran during executor shutdown, but rollback saw no futures and left their index files behind. Appending each future as it is returned is not enough: `submit()` puts the work item on the queue before it starts an extra worker, so even the submission that raises can run its shard on an already running worker. Track output ownership in the worker itself so rollback no longer depends on the future list, and cancel still-queued shards when submission fails. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
JingsongLi
left a comment
There was a problem hiding this comment.
Reviewed the output rollback ownership changes. Retaining rollback responsibility when a later executor submission fails prevents already-created outputs from being orphaned. The failure-path coverage looks good to me.
Purpose
PyPaimon's generic global-index builder currently builds vector and full-text shards one at a time. On multi-shard builds this leaves storage bandwidth and native index capacity idle.
Add
global-index.build.parallelism(default1) and use a bounded thread pool when the value is greater than one. Each worker keeps the existing streaming Arrow reader and one writer per shard, manifest messages are returned in shard-plan order, and a failed build closes active resources and removes every uncommitted index file. The documentation warns that each shard may also start native worker threads.Benchmark
A temporary local benchmark exercised the production shard scheduler and streaming batch path with 4,000 approximately 2 KiB text rows per shard. It injected 20 ms of read latency plus 40 ms of native-build/upload latency per shard to isolate the scheduling change. Results are medians of three runs on macOS arm64; RSS is peak resident memory.
The ablation shows the throughput/memory tradeoff behind the conservative default. The benchmark uses simulated native and I/O latency, so end-to-end gains will depend on the index implementation, storage, shard size, and native thread settings. The temporary benchmark harness is not included in this change.
Tests
python -m pytest -q pypaimon/tests/global_index_build_test.py(36 passed)python -m flake8 --config=dev/cfg.ini pypaimon/common/options/core_options.py pypaimon/globalindex/create_global_index.py pypaimon/tests/global_index_build_test.py