Skip to content

Fix HNSW graph merge concurrency in case of intra-merge executor (#16729) - #16765

Open
CH-Abhinav wants to merge 1 commit into
apache:mainfrom
CH-Abhinav:hnsw-concurrent-merge
Open

CH-Abhinav wants to merge 1 commit into
apache:mainfrom
CH-Abhinav:hnsw-concurrent-merge

Conversation

@CH-Abhinav

Copy link
Copy Markdown
Contributor

Description

Fixes #16729.

Problem

In HnswConcurrentMergeBuilder, $N$ workers were previously constructed as greedy tasks where each worker executed an internal loop draining batches from a shared atomic counter (workProgress) until all vectors were added. These $N$ tasks were submitted in a single call to TaskExecutor#invokeAll.

When ConcurrentMergeScheduler.CachedExecutor is briefly saturated at the time of submission (i.e. maxThreadCount - mergeThreads.size() - 1 <= 0 due to concurrent merges), it falls back to executing the task inline on the calling thread. Because worker 0 was an unbounded loop draining all available batches, it consumed 100% of the merge's work inline before returning control to the submit loop. Even if sibling merges completed milliseconds later and freed up intra-merge thread capacity, helper threads were never enlisted, locking the entire merge to 1.00x concurrency (running single-threaded for its entire duration).

Solution

  1. Chunked Task Submission: Instead of submitting $N$ long-lived greedy tasks, the ordinal space [0, maxOrd) is partitioned into discrete batch tasks of batchSize (default 2048) and submitted to TaskExecutor#invokeAll.
  2. Worker Pool: A BlockingQueue<ConcurrentMergeWorker> workerPool = new ArrayBlockingQueue<>(workers.length) holds the pre-initialized workers (each with its own RandomVectorScorerSupplier and scratch buffers).
  3. Mid-Merge Concurrency Recovery: When a task runs (either on an intra-merge pool thread or inline on the calling thread), it leases a worker from workerPool, executes a single slice [start, end), and releases the worker back in a finally block. If the calling thread must run a task inline, it only processes one batch (~50ms) rather than draining the whole merge, allowing subsequent batches to be scheduled onto helper threads as soon as sibling merges finish.
  4. Cleanup: Removed unused workProgress, run(int maxOrd), and getStartPos(int maxOrd) from ConcurrentMergeWorker.

Tests

  • Added TestHnswConcurrentMergeSubmission:
    • Reproduces organic thread contention using stock ConcurrentMergeScheduler, TieredMergePolicy, and vector merges $&gt; 50\text{MB}$ (clearing MIN_BIG_MERGE_MB).
    • Verifies that a merge outliving sibling merges recovers intra-merge threads mid-merge (batchesByThread.size() > 1).
    • Annotated with @LuceneTestCase.Monster and @Repeat(iterations = 3) given segment size and indexing requirements.
  • Verified ./gradlew check -x test passes all precommit checks (licenses, ecj linter for main and test, forbidden APIs, and tidy formatting).
  • Verified existing tests (TestConcurrentMergeScheduler, TestHnswFloatVectorGraph) pass without regressions.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

A briefly saturated intra-merge executor can serialize an entire HNSW graph merge

1 participant