Visitar URL original
[python] Evaluate Ray vector prefilters on workers by TheR1sing3un · Pull Request #10446 · apache/paimon · GitHub
Skip to content

[python] Evaluate Ray vector prefilters on workers - #10446

Merged
JingsongLi merged 1 commit into
apache:masterfrom
TheR1sing3un:contribute/ray-worker-prefilter
Oct 9, 2026
Merged

JingsongLi merged 1 commit into
apache:masterfrom
TheR1sing3un:contribute/ray-worker-prefilter

Conversation

@TheR1sing3un

Copy link
Copy Markdown
Member

Purpose

Ray vector search currently builds scalar and live-row filter bitmaps on the driver before dispatching indexed searches. Candidate-only scalar predicates can therefore read and verify filter columns for all indexed candidates on the driver, limiting the benefit of additional Ray workers.

Move prefilter execution for single and batch vector searches into Ray tasks. Each search task loads deletion vectors and verifies candidates only within its index shard's row-ID range, using the planned snapshot. Tasks with identical scalar-index inputs share one scalar evaluation through a top-level Ray ObjectRef dependency, avoiding repeated scans of a shared index without materializing its bitmap on the driver. Candidate verification does not rediscover scalar indexes, and batch tasks reuse their final filter across query blocks. Local search keeps its existing execution path.

Performance was measured with real Parquet data, BTree/Bitmap indexes, IVF_FLAT vector searches and final row lookup on a local Ray cluster: Python 3.11.15, Ray 2.54.0, PyArrow 19.0.1, paimon-vindex 0.5.0; 15 logical CPUs, 48 GiB RAM. The dataset has 1,000,000 rows, eight vector shards, a shared BTree index, 16 float32 dimensions, L2, nlist=1, k=10, and a filter matching 10% of rows. Both variants return identical IDs. The baseline reproduces the original driver-prefilter dispatch. Workers and filesystem cache were warmed and timed runs alternated variants.

  • With name LIKE '%needle%', global-index.filter.refine-from-data=true, and the default BTree fallback scan threshold, four workers improve median end-to-end latency from 2.598 s to 1.705 s (1.52x) over three timed runs per variant.
  • With btree-index.fallback-scan-max-size=0 b to exercise filter-column verification without scanning the candidate-only BTree, five timed runs per variant give:
Workers Driver prefilter Worker prefilter Speedup
1 1.290 s 1.245 s 1.04x
2 1.277 s 0.661 s 1.93x
4 1.229 s 0.379 s 3.24x
8 1.213 s 0.219 s 5.53x

There is a scheduling tradeoff for cheap filters: at four workers, the exact Bitmap-filter case changes from 34.0 ms to 41.2 ms; unfiltered search is 36.4 ms versus 34.8 ms. On 1,000 rows (seven timed runs), exact filtering changes from 23.5 ms to 28.9 ms and candidate verification from 30.7 ms to 32.1 ms. These are single-machine, warm-filesystem results, not a claim of universal or multi-node speedup. Partially overlapping scalar-index input sets can still repeat index reads.

Tests

  • 286 targeted tests passed with Ray 2.54.0, covering single/batch search, global refinement and duplicate precedence, scalar exactness, full-text regression, and scoped deletion-vector reads.
  • 19 targeted tests passed with Ray 2.59.0.
  • New/updated tests verify actual worker-side candidate filtering, filter-column-only reads without repeated scalar discovery, snapshot consistency across a concurrent delete, shared scalar evaluation occurring once outside the driver, and dependency scheduling when each task requests all cluster CPUs.
  • Flake8 with the repository configuration and git diff --check passed.

Performance scripts and data were kept outside the repository; no benchmark code is included.

@JingsongLi

Copy link
Copy Markdown
Contributor

+1

@JingsongLi
JingsongLi merged commit e93ab38 into apache:master Oct 9, 2026
12 of 14 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants