diff --git a/lucene/CHANGES.txt b/lucene/CHANGES.txt index ef52cbb7d42c..54ee152f0298 100644 --- a/lucene/CHANGES.txt +++ b/lucene/CHANGES.txt @@ -130,6 +130,10 @@ Bug Fixes GlobalOrdinalsWithScoreCollector, fixing an intermittent failure caused by non-associative float addition that #16378 missed. (Luca Cavanna) +* GITHUB#16419: Restore KnnVectorsReader#finishMerge() calls after merges, including + exceptional exits, so vector inputs switch back from sequential read advice. + (Simon Cooper, John Maurice) + Other --------------------- * GITHUB#16266: Remove deprecated search(Query, Collector) calls in QueryUtils by replacing diff --git a/lucene/core/src/java/org/apache/lucene/index/IndexWriter.java b/lucene/core/src/java/org/apache/lucene/index/IndexWriter.java index 80ff607a4efc..9a0a00b5eeb9 100644 --- a/lucene/core/src/java/org/apache/lucene/index/IndexWriter.java +++ b/lucene/core/src/java/org/apache/lucene/index/IndexWriter.java @@ -3565,22 +3565,26 @@ public void addIndexesReaderMerge(MergePolicy.OneMerge merge) throws IOException intraMergeExecutor, merge); - if (!merger.shouldMerge()) { - return; - } - - merge.checkAborted(); - synchronized (this) { - runningAddIndexesMerges.add(merger); - } - merge.mergeStartNS = System.nanoTime(); try { - merger.merge(); // merge 'em - } finally { + if (!merger.shouldMerge()) { + return; + } + + merge.checkAborted(); synchronized (this) { - runningAddIndexesMerges.remove(merger); - notifyAll(); + runningAddIndexesMerges.add(merger); } + merge.mergeStartNS = System.nanoTime(); + try { + merger.merge(); // merge 'em + } finally { + synchronized (this) { + runningAddIndexesMerges.remove(merger); + notifyAll(); + } + } + } finally { + merger.cleanupMerge(); } merge.setMergeInfo( @@ -5372,32 +5376,36 @@ public int length() { context, intraMergeExecutor, merge); - merge.info.setSoftDelCount(Math.toIntExact(softDeleteCount.get())); - merge.checkAborted(); - MergeState mergeState = merger.mergeState; MergeState.DocMap[] docMaps; - if (reorderDocMaps == null) { - docMaps = mergeState.docMaps; - } else { - // Since the reader was reordered, we passed a merged view to MergeState and from its - // perspective there is a single input segment to the merge and the - // SlowCompositeCodecReaderWrapper is effectively doing the merge. - assert mergeState.docMaps.length == 1 - : "Got " + mergeState.docMaps.length + " docMaps, but expected 1"; - MergeState.DocMap compactionDocMap = mergeState.docMaps[0]; - docMaps = new MergeState.DocMap[reorderDocMaps.length]; - for (int i = 0; i < docMaps.length; ++i) { - MergeState.DocMap reorderDocMap = reorderDocMaps[i]; - docMaps[i] = docID -> compactionDocMap.get(reorderDocMap.get(docID)); + try { + merge.info.setSoftDelCount(Math.toIntExact(softDeleteCount.get())); + merge.checkAborted(); + + if (reorderDocMaps == null) { + docMaps = mergeState.docMaps; + } else { + // Since the reader was reordered, we passed a merged view to MergeState and from its + // perspective there is a single input segment to the merge and the + // SlowCompositeCodecReaderWrapper is effectively doing the merge. + assert mergeState.docMaps.length == 1 + : "Got " + mergeState.docMaps.length + " docMaps, but expected 1"; + MergeState.DocMap compactionDocMap = mergeState.docMaps[0]; + docMaps = new MergeState.DocMap[reorderDocMaps.length]; + for (int i = 0; i < docMaps.length; ++i) { + MergeState.DocMap reorderDocMap = reorderDocMaps[i]; + docMaps[i] = docID -> compactionDocMap.get(reorderDocMap.get(docID)); + } } - } - merge.mergeStartNS = System.nanoTime(); + merge.mergeStartNS = System.nanoTime(); - // This is where all the work happens: - if (merger.shouldMerge()) { - merger.merge(); + // This is where all the work happens: + if (merger.shouldMerge()) { + merger.merge(); + } + } finally { + merger.cleanupMerge(); } assert mergeState.segmentInfo == merge.info.info; diff --git a/lucene/core/src/java/org/apache/lucene/index/SegmentMerger.java b/lucene/core/src/java/org/apache/lucene/index/SegmentMerger.java index b83b5d51c032..321780406ad7 100644 --- a/lucene/core/src/java/org/apache/lucene/index/SegmentMerger.java +++ b/lucene/core/src/java/org/apache/lucene/index/SegmentMerger.java @@ -23,6 +23,7 @@ import org.apache.lucene.codecs.Codec; import org.apache.lucene.codecs.DocValuesConsumer; import org.apache.lucene.codecs.FieldsConsumer; +import org.apache.lucene.codecs.KnnVectorsReader; import org.apache.lucene.codecs.KnnVectorsWriter; import org.apache.lucene.codecs.NormsConsumer; import org.apache.lucene.codecs.NormsProducer; @@ -227,7 +228,7 @@ private void mergeTerms(SegmentWriteState segmentWriteState, SegmentReadState se } } - public void mergeFieldInfos() { + private void mergeFieldInfos() { for (FieldInfos readerFieldInfos : mergeState.fieldInfos) { for (FieldInfo fi : readerFieldInfos) { fieldInfosBuilder.add(fi); @@ -326,4 +327,12 @@ private void mergeWithLogging( + " docs]"); } } + + void cleanupMerge() throws IOException { + for (KnnVectorsReader reader : mergeState.knnVectorsReaders) { + if (reader != null) { + reader.finishMerge(); + } + } + } } diff --git a/lucene/core/src/test/org/apache/lucene/codecs/lucene99/TestMergeReadAdviceRevert.java b/lucene/core/src/test/org/apache/lucene/codecs/lucene99/TestMergeReadAdviceRevert.java new file mode 100644 index 000000000000..3036d40d65ee --- /dev/null +++ b/lucene/core/src/test/org/apache/lucene/codecs/lucene99/TestMergeReadAdviceRevert.java @@ -0,0 +1,205 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.lucene.codecs.lucene99; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import org.apache.lucene.document.Document; +import org.apache.lucene.document.KnnFloatVectorField; +import org.apache.lucene.index.DirectoryReader; +import org.apache.lucene.index.IndexWriter; +import org.apache.lucene.index.IndexWriterConfig; +import org.apache.lucene.index.TieredMergePolicy; +import org.apache.lucene.index.VectorSimilarityFunction; +import org.apache.lucene.store.DataAccessHint; +import org.apache.lucene.store.Directory; +import org.apache.lucene.store.FilterDirectory; +import org.apache.lucene.store.FilterIndexInput; +import org.apache.lucene.store.IOContext; +import org.apache.lucene.store.IndexInput; +import org.apache.lucene.store.MMapDirectory; +import org.apache.lucene.tests.index.BaseKnnVectorsFormatTestCase; +import org.apache.lucene.tests.util.LuceneTestCase; +import org.apache.lucene.tests.util.TestUtil; + +public class TestMergeReadAdviceRevert extends LuceneTestCase { + + private static final int DIM = 16; + + public void testSequentialAdviceIsRevertedAfterMerge() throws Exception { + Recorder recorder = new Recorder(); + try (Directory raw = new MMapDirectory(createTempDir()); + Directory dir = new RecordingDirectory(raw, recorder)) { + + IndexWriterConfig iwc = new IndexWriterConfig(); + iwc.setCodec(TestUtil.alwaysKnnVectorsFormat(new Lucene99HnswVectorsFormat())); + // Expose .vec files to RecordingDirectory. + iwc.setUseCompoundFile(false); + iwc.setMergePolicy(new TieredMergePolicy()); + + try (IndexWriter w = new IndexWriter(dir, iwc)) { + for (int seg = 0; seg < 2; seg++) { + for (int i = 0; i < 64; i++) { + Document doc = new Document(); + doc.add( + new KnnFloatVectorField( + "field", + BaseKnnVectorsFormatTestCase.randomNormalizedVector(DIM), + VectorSimilarityFunction.DOT_PRODUCT)); + w.addDocument(doc); + } + w.commit(); + } + + // Keep the source SegmentReaders open so the merge reuses their vector inputs. + try (DirectoryReader nrt = DirectoryReader.open(w)) { + assertEquals(2, nrt.leaves().size()); + recorder.mark("--- forceMerge(1) start ---"); + w.forceMerge(1); + recorder.mark("--- forceMerge(1) done ---"); + + // Check before the source SegmentReaders are closed. + List offenders = new ArrayList<>(); + List sawSequential = new ArrayList<>(); + for (Map.Entry> e : recorder.events().entrySet()) { + String file = e.getKey(); + if (file.endsWith("." + Lucene99FlatVectorsFormat.VECTOR_DATA_EXTENSION) == false) { + continue; + } + List events = e.getValue(); + String lastAdvice = null; + for (String ev : events) { + if (ev.startsWith("ADVICE:")) { + lastAdvice = ev.substring("ADVICE:".length()); + } + } + if (events.contains("ADVICE:SEQUENTIAL")) { + sawSequential.add(file); + if ("SEQUENTIAL".equals(lastAdvice)) { + offenders.add(file + " " + events); + } + } + } + + assertFalse( + "no .vec input received SEQUENTIAL advice:\n" + recorder.dump(), + sawSequential.isEmpty()); + + assertTrue( + ".vec inputs still using SEQUENTIAL advice:\n " + + String.join("\n ", offenders) + + "\n\nEvents:\n" + + recorder.dump(), + offenders.isEmpty()); + } + } + } + } + + static final class Recorder { + private final Map> events = new LinkedHashMap<>(); + private final List timeline = new ArrayList<>(); + + synchronized void record(String file, String event) { + events.computeIfAbsent(file, k -> new ArrayList<>()).add(event); + timeline.add(file + " -> " + event); + } + + synchronized void mark(String note) { + timeline.add(note); + } + + synchronized Map> events() { + return events; + } + + synchronized String dump() { + return String.join("\n", timeline); + } + } + + static final class RecordingDirectory extends FilterDirectory { + private final Recorder recorder; + + RecordingDirectory(Directory in, Recorder recorder) { + super(in); + this.recorder = recorder; + } + + @Override + public IndexInput openInput(String name, IOContext context) throws IOException { + return new RecordingIndexInput( + "Recording(" + name + ")", super.openInput(name, context), name, recorder, true); + } + } + + static final class RecordingIndexInput extends FilterIndexInput { + private final String name; + private final Recorder recorder; + private final boolean top; + + RecordingIndexInput(String desc, IndexInput in, String name, Recorder recorder, boolean top) { + super(desc, in); + this.name = name; + this.recorder = recorder; + this.top = top; + } + + private static String hintOf(IOContext ctx) { + return ctx.hints(DataAccessHint.class).findFirst().map(DataAccessHint::name).orElse("NONE"); + } + + @Override + public void updateIOContext(IOContext context) throws IOException { + recorder.record(name, "ADVICE:" + hintOf(context) + (top ? "" : "@clone")); + in.updateIOContext(context); + } + + @Override + public IndexInput clone() { + return new RecordingIndexInput(toString(), in.clone(), name, recorder, false); + } + + @Override + public IndexInput slice(String sliceDescription, long offset, long length) throws IOException { + return new RecordingIndexInput( + sliceDescription, in.slice(sliceDescription, offset, length), name, recorder, false); + } + + @Override + public IndexInput slice(String sliceDescription, long offset, long length, IOContext context) + throws IOException { + return new RecordingIndexInput( + sliceDescription, + in.slice(sliceDescription, offset, length, context), + name, + recorder, + false); + } + + @Override + public void close() throws IOException { + if (top) { + recorder.record(name, "CLOSE"); + } + super.close(); + } + } +} diff --git a/lucene/core/src/test/org/apache/lucene/index/TestDoc.java b/lucene/core/src/test/org/apache/lucene/index/TestDoc.java index 0027e7409f71..97315d435b9f 100644 --- a/lucene/core/src/test/org/apache/lucene/index/TestDoc.java +++ b/lucene/core/src/test/org/apache/lucene/index/TestDoc.java @@ -243,6 +243,7 @@ private SegmentCommitInfo merge( null); merger.merge(); + merger.cleanupMerge(); r1.close(); r2.close(); si.setFiles(new HashSet<>(trackingDir.getCreatedFiles())); diff --git a/lucene/core/src/test/org/apache/lucene/index/TestSegmentMerger.java b/lucene/core/src/test/org/apache/lucene/index/TestSegmentMerger.java index fb290d0a467d..56dd6ed2e2ef 100644 --- a/lucene/core/src/test/org/apache/lucene/index/TestSegmentMerger.java +++ b/lucene/core/src/test/org/apache/lucene/index/TestSegmentMerger.java @@ -110,6 +110,7 @@ public void testMerge() throws IOException { new SameThreadExecutorService(), null); MergeState mergeState = merger.merge(); + merger.cleanupMerge(); int docsMerged = mergeState.segmentInfo.maxDoc(); assertTrue(docsMerged == 2); // Should be able to open a new SegmentReader against the new directory diff --git a/lucene/test-framework/src/java/org/apache/lucene/tests/codecs/asserting/AssertingKnnVectorsFormat.java b/lucene/test-framework/src/java/org/apache/lucene/tests/codecs/asserting/AssertingKnnVectorsFormat.java index a9f57990191f..dda8c26830c2 100644 --- a/lucene/test-framework/src/java/org/apache/lucene/tests/codecs/asserting/AssertingKnnVectorsFormat.java +++ b/lucene/test-framework/src/java/org/apache/lucene/tests/codecs/asserting/AssertingKnnVectorsFormat.java @@ -234,7 +234,7 @@ public Map getOffHeapByteSize(FieldInfo fieldInfo) { public void close() throws IOException { delegate.close(); delegate.close(); // impls should be able to handle multiple closes - assert finishMergeCount.get() <= 0 || mergeInstanceCount.get() == finishMergeCount.get(); + assert mergeInstanceCount.get() == finishMergeCount.get(); } @Override