jimczi commented on code in PR #16705:
URL: https://github.com/apache/lucene/pull/16705#discussion_r4093562796


##########
lucene/core/src/java/org/apache/lucene/index/KnnVectorValues.java:
##########
@@ -50,12 +50,47 @@ public int ordToDoc(int ord) {
   }
 
   /**
-   * Prefetches the provided ordinals
+   * Prefetches {@code count} consecutive vectors starting at the given 
ordinal, so that later calls
+   * to read them are more likely to hit memory. Implementations should start 
the reads and return
+   * without waiting for them, and are free to prefetch fewer vectors than 
asked for, including none
+   * at all. {@code count} is clamped to the number of vectors remaining after 
{@code ord}. The
+   * default implementation is a no-op.
+   *
+   * @param ord the ordinal of the first vector to prefetch
+   * @param count how many consecutive vectors to prefetch, starting at {@code 
ord}
+   * @return the number of vectors a prefetch was actually issued for, {@code 
0} if none, in which
+   *     case the caller gains nothing by deferring the reads
+   */
+  public int prefetch(int ord, int count) throws IOException {

Review Comment:
   Can this return a boolean? `IndexInput#prefetch` returns true if it 
prefetched something, and `StoredFields`/`TermVectors` return void. 
`RescoreTopNQuery` only needs to know whether anything was issued, and count is 
clamped anyway so the number doesn't say much.



##########
lucene/core/src/java/org/apache/lucene/search/RescoreTopNQuery.java:
##########
@@ -71,24 +95,82 @@ public Query rewrite(IndexSearcher indexSearcher) throws 
IOException {
       DoubleValues rescores = rewrittenValueSource.getValues(leaf, 
getDoubleValues(innerScorer));
       DocIdSetIterator iterator = innerScorer.iterator();
       while (iterator.nextDoc() != DocIdSetIterator.NO_MORE_DOCS) {
-        int docId = iterator.docID();
-        if (rescores.advanceExact(docId)) {
-          double v = rescores.doubleValue();
-          queue.insertWithOverflow(new ScoreDoc(leaf.docBase + docId, (float) 
v));
-        } else {
-          queue.insertWithOverflow(new ScoreDoc(leaf.docBase + docId, 0f));
-        }
+        rescoreInto(queue, rescores, leaf.docBase, iterator.docID());
         originalCount++;
       }
     }
-    int i = 0;
-    ScoreDoc[] scoreDocs = new ScoreDoc[queue.size()];
-    for (ScoreDoc topDoc : queue) {
-      scoreDocs[i++] = topDoc;
+    return originalCount;
+  }
+
+  /**
+   * Starts the loads for every candidate before scoring any of them, so that 
more than one read is
+   * in flight when the values live on slow storage. Each segment's prefetches 
are issued on the
+   * searcher's executor: prefetching costs real CPU per candidate, so issuing 
a whole shortlist
+   * from a single thread caps how many reads can be outstanding. Only doc ids 
are buffered, never
+   * values.
+   */
+  private int rescoreWithPrefetch(
+      IndexSearcher indexSearcher,
+      IndexReader reader,
+      Weight weight,
+      DoubleValuesSource rewrittenValueSource,
+      HitQueue queue)
+      throws IOException {
+    final List<LeafReaderContext> leaves = reader.leaves();
+    final DoubleValues[] leafValues = new DoubleValues[leaves.size()];
+    final int[][] leafDocs = new int[leaves.size()][];
+    final List<Callable<Void>> tasks = new ArrayList<>(leaves.size());
+    for (int i = 0; i < leaves.size(); i++) {
+      final int idx = i;
+      final LeafReaderContext leaf = leaves.get(i);
+      tasks.add(
+          () -> {
+            Scorer innerScorer = weight.scorer(leaf);
+            if (innerScorer == null) {
+              return null;
+            }
+            DoubleValues rescores = rewrittenValueSource.getValues(leaf, null);
+            DocIdSetIterator iterator = innerScorer.iterator();
+            int[] docs = new int[16];
+            int count = 0;
+            while (iterator.nextDoc() != DocIdSetIterator.NO_MORE_DOCS) {
+              int docId = iterator.docID();
+              rescores.prefetch(docId);
+              if (count == docs.length) {
+                docs = ArrayUtil.grow(docs, count + 1);
+              }
+              docs[count++] = docId;
+            }
+            leafValues[idx] = rescores;
+            leafDocs[idx] = ArrayUtil.copyOfSubArray(docs, 0, count);
+            return null;
+          });
+    }
+    indexSearcher.getTaskExecutor().invokeAll(tasks);

Review Comment:
   The description says the task executor part is orthogonal and comes as a 
separate PR, but it's still here. I'd drop it and keep this single threaded for 
now.
   
   One task per leaf isn't a fair split. The topN candidates land in segments 
in proportion to their size, so one task gets most of them and the others 
return immediately. To do it properly you'd flatten the (leaf, doc) pairs and 
cut them into equal ranges, and then you need a `DoubleValues` per task per 
leaf since they're not thread safe. That's worth its own PR.
   
   Your own numbers point the same way, on 6.18 the single threaded path 
already reaches aqu-sz 10.6, and the 1.6x was on 6.1 where it only got 2.91.



-- 
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