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]

Reply via email to