abernardi597 opened a new issue, #16729: URL: https://github.com/apache/lucene/issues/16729
### Description ###### Findings by humans, report and reproductions were AI-assisted (Opus 5). Spinoff from [this thread](https://github.com/mikemccand/luceneutil/pull/507#discussion_r3045850830) on `luceneutil`. ### Summary `HnswConcurrentMergeBuilder` hands out work through one shared batch counter, so any single worker can drain the whole merge. Its workers are all submitted before any of them runs, and the intra-merge executor runs work on the calling thread when it has no spare thread. Those two facts combine badly: one such inline execution consumes the entire graph merge, and the workers submitted afterwards find no work left. The merge then runs single-threaded for its whole duration, even when threads free up moments later. ### Mechanism `ConcurrentMergeWorker#run` claims work by advancing a shared `workProgress` counter one batch at a time. A worker keeps claiming batches until the counter passes `maxOrd`. So a worker that runs alone does not do `1/N` of the merge. It does all of it. `TaskExecutor#invokeAll` submits `count - 1` tasks in a loop, then runs whatever is left on the calling thread. That loop is sequential and blocking. `ConcurrentMergeScheduler.CachedExecutor#execute` runs the command on the calling thread when no thread is available, deciding availability from `maxThreadCount - mergeThreads.size() - 1`. Put together: the first submission runs inline, that worker drains the counter, and `invokeAll` does not reach the remaining submissions until it returns. Those later submissions do get threads. They find the counter already past `maxOrd` and return immediately. This is not a defect in `TaskExecutor`. Running a command on the calling thread is a supported input to it, and `IndexSearcher` and `PostCollectionFaceting` both pass `Runnable::run` deliberately. The combination only bites a caller whose tasks share a work source. ### When it manifests Three conditions: - The vectors format is built with `numMergeWorkers > 1` and a null `mergeExec`, so the codec uses `MergeState#intraMergeTaskExecutor`. Lucene's default codec passes `DEFAULT_NUM_MERGE_WORKER == 1` and never constructs a `ConcurrentHnswMerger`, so a default `IndexWriter` is unaffected. - The merge is big enough to be given real threads. `ConcurrentMergeScheduler#getIntraMergeExecutor` returns a `SameThreadExecutorService` when `merge.estimatedMergeBytes < MIN_BIG_MERGE_MB * 1024 * 1024`, and `MIN_BIG_MERGE_MB` is 50. Smaller merges are serial by design, so the lost-parallelism bug does not apply to them. Note they are not routed away from the concurrent builder though; see below. - `ConcurrentMergeScheduler` has no spare thread at the moment the workers are submitted. - Threads free up before the graph merge finishes. The third condition is what makes this a bug rather than a capacity limit. If the executor stays saturated for the whole merge then single-threaded execution is the correct outcome. The problem is capacity that arrives after submission and is never taken up. Concurrent large merges make this ordinary rather than exotic, and no unusual configuration is needed to hit it. Each running merge occupies a merge thread, `CachedExecutor` subtracts those from `maxThreadCount`, and merges finish at different times. A merge that starts while its siblings are running gets no intra-merge threads, and cannot pick any up when a sibling completes. Three merges of equal size starting together is enough. So is one large tier merge beginning while a batch of smaller merges holds the merge threads. `InfoStream` already exposes it. `HnswConcurrentMergeBuilder#build` logs effective concurrency, so an affected merge reports `1.00x` while reporting N workers on the line above. ### Scope `HnswConcurrentMergeBuilder` is the only `TaskExecutor#invokeAll` caller whose tasks share a work source, and the only consumer of `MergeState#intraMergeTaskExecutor`. The total degradation described here is therefore specific to HNSW graph merging. A weaker form is general: achieved concurrency is whatever the executor can offer during the submit loop, and it is never revisited. Callers with pre-partitioned tasks degrade proportionally rather than totally, because each task carries its own share of the work. `BPIndexReorderer` and `BpVectorReorderer` are a useful contrast. They split recursively and call `invokeAll(leftTask, rightTask)` at each node, so later nodes re-offer work as the recursion proceeds. Capacity appearing mid-run does get used there. ### Reproduction Two exhibits. The first shows stock `ConcurrentMergeScheduler` producing the bug on its own. The second removes thread capacity as a possible explanation for it. #### 1. Organic: stock ConcurrentMergeScheduler, unmodified Lucene Twelve segments of 15,000 docs with 512-dimension `FLOAT32` vectors, then `maybeMerge()`. Stock `ConcurrentMergeScheduler` with `setMaxMergesAndThreads(12, 6)`, `TieredMergePolicy` with `setSegmentsPerTier(2)`, and the vectors format built with `numMergeWorkers = 8` and a null `mergeExec`. Nothing is overridden: CMS decides for itself whether a thread is available, from its own `maxThreadCount - mergeThreads.size() - 1`. Three merges start together. All three are 58.6MB, so all three clear `MIN_BIG_MERGE_MB` and get the real intra-merge executor. ``` merge vectors size window(ms) requested reported threads that worked 1 30000 58.6M 0->4673 8 2.00x {Lucene Merge Thread #2=6} 2 30000 58.6M 0->4760 8 1.98x {Lucene Merge Thread #0=6} 3 30000 58.6M 0->7971 8 1.00x {Lucene Merge Thread #1=15} merge 3 reported 1.00x and spent its last 3211ms as the only merge in the system, so CMS had 4 intra-merge threads free and unused ``` With three merges running, the budget is `6 - 3 - 1 = 2` helper threads for three merges. Merges 1 and 2 each took one and reached about 2x. Merge 3 got none, so its first worker ran on the merge thread and claimed all 15 of its 15 batches. Merges 1 and 2 finish at about 4.7s, after which merge 3 is alone and the budget is back to 4. It stays single-threaded for another 3.2s and takes 8.0s against the 4.7s its equally sized siblings took. Everything above comes from upstream `InfoStream`. `HnswConcurrentMergeBuilder` emits its messages from the worker thread, so an `InfoStream` that records `Thread.currentThread().getName()` attributes work to threads with no changes to Lucene. One caveat on that column: when merges overlap, a helper thread's message cannot be tied to a specific merge from `InfoStream` alone, so the counts for merges 1 and 2 are incomplete. Merge 3's is complete, since 15 batches is its entire workload. The arrival pattern is arranged: the segments are built under `NoMergePolicy` first so that three equal merges start together. The merge behaviour itself is entirely stock. An index that accumulated segments with merging paused, then resumed, reaches the same state. #### 2. Controlled: capacity ruled out The organic run leaves one objection open, that merge 3 was simply short of threads. This removes that. Same index and codec, but both runs override `getIntraMergeExecutor` to return an **unbounded** `Executors.newCachedThreadPool()`, so neither run is ever short of a thread. The only difference is that in the second run exactly one submission, the first, is served on the calling thread rather than handed to that pool. That is what `CachedExecutor` does when the budget is exhausted. Three trials of each, alternating: ``` forceMerge of 60000 vectors, numMergeWorkers=8, intra-merge pool = unbounded in both runs intra-merge executor wall(ms) reported threads batches per thread unbounded, all submits forked 3449 7.28x 8 mergeThread=4 pool-1=3 ... pool-7=4 unbounded, 1st submit inline 21634 1.00x 1 mergeThread=30 unbounded, all submits forked 3210 7.18x 8 mergeThread=2 pool-1=6 ... pool-7=2 unbounded, 1st submit inline 21704 1.00x 1 mergeThread=30 unbounded, all submits forked 3245 7.20x 8 mergeThread=2 pool-1=9 ... pool-7=2 unbounded, 1st submit inline 21505 1.00x 1 mergeThread=30 ``` Serving one submission on the calling thread costs 6.7x, 21.6s against 3.2s, on identical input with threads freely available throughout. The degraded runs are identical across trials: always exactly one thread, always all 30 batches. This is a controlled experiment rather than a reproduction. It does not show CMS deciding to inline anything; exhibit 1 does that. What it shows is that once a submission is served on the calling thread, the parallelism is lost for reasons that have nothing to do with how many threads exist. ### Properties a fix has to respect I am not prescribing a fix, since the shape of one depends on how much machinery Lucene wants to own here. These three properties of the workload seem worth stating, because each one rules out an otherwise reasonable approach. **Reads during the concurrent phase are random, so the counter's ordered pass is worth little.** The counter walks ordinals in order, which looks like a sequential pass over the vector file. Almost all of the reads are not that pass. Inserting a node reads its own vector once, then beam-searches the graph and scores the vectors of every node visited, which are scattered across the ordinal space by construction. Lucene already classifies this access as random: the reader that the graph build scores against is opened with `DataAccessHint.RANDOM`, and `Lucene99FlatVectorsReader#getMergeInstance` flips the hint to `SEQUENTIAL` only for the separate source-side merge read. That sequential read also happens outside the concurrent phase, since `mergeOneField` merges flat vectors eagerly and defers only the graph build. **Work is not uniformly priced across ordinals, so contiguous partitioning is exposed.** `ConcurrentMergeWorker#addGraphNode` returns immediately for any node in `initializedNodes`, and `IncrementalHnswGraphMerger#getNewOrdMapping` sets those bits at the ordinal in merged doc order. For an unsorted index a reused segment's nodes therefore occupy one contiguous run of ordinals. A 28,483 vector merge in a separate session showed this directly, 43% of its ordinals being a contiguous prefix of no-ops and 6 of its 14 batches free: ``` addVectors [0 2048) 0.13 ms <- already in the reused graph addVectors [2048 4096) 0.00 ms addVectors [4096 6144) 0.00 ms addVectors [6144 8192) 0.00 ms addVectors [8192 10240) 0.05 ms addVectors [10240 12288) 0.00 ms addVectors [12288 14336) 56.17 ms <- reused run ends mid-batch addVectors [14336 16384) 1142.85 ms <- real insertions from here on addVectors [16384 18432) 1093.85 ms ... ``` Split that merge into eight contiguous ranges and the first three workers would have nothing to do. Note this is a static property of the ordinal rather than a function of elapsed time, so concurrency does not average it out. **Static assignment gives up tail balance.** The counter lets a worker that finishes early absorb straggler batches. Any fixed division loses that, so per-vector cost variance shows up as tail latency on `invokeAll`. ### Note on the current contract One observation that may matter whichever direction a fix takes. `CachedExecutor` already computes available capacity under the `ConcurrentMergeScheduler` lock, but does not expose it. A caller submitting work cannot distinguish "no thread was available" from "the work ran", so it has no basis for asking again later. ### Version and environment details _No response_ -- 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] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
