From 5f164dc6b7909fc4a8b94fdfa656ed7a8ccdbbdf Mon Sep 17 00:00:00 2001 From: Jim Ferenczi Date: Thu, 9 Jul 2026 14:49:22 +0000 Subject: [PATCH 1/5] Apply sequential read advice during merge for bbq_disk vectors IVFVectorsReader did not override getMergeInstance() or finishMerge(), which Lucene calls during a merge to switch the read advice and restore it afterwards. Because of that, the flat vector reader kept the RANDOM advice used for search throughout the merge, disabling read-ahead on the sequential scans (checkIntegrity and the vector copy) that a merge performs. Override both methods so the merge instance is backed by the flat readers' merge instances, which switch their input to SEQUENTIAL, and finishMerge() restores the search access pattern. The centroid and posting list inputs are opened with the default context and already use read-ahead advice, so they are left unchanged. Signed-off-by: Jim Ferenczi --- docs/changelog/153423.yaml | 5 +++ .../diskbbq/ES920DiskBBQVectorsReader.java | 9 +++++ .../vectors/diskbbq/IVFVectorsReader.java | 36 +++++++++++++++++++ .../es94/ES940DiskBBQVectorsReader.java | 9 +++++ .../es95/ES950DiskBBQVectorsReader.java | 9 +++++ .../next/ESNextDiskBBQVectorsReader.java | 9 +++++ 6 files changed, 77 insertions(+) create mode 100644 docs/changelog/153423.yaml diff --git a/docs/changelog/153423.yaml b/docs/changelog/153423.yaml new file mode 100644 index 0000000000000..4e238d7e903ee --- /dev/null +++ b/docs/changelog/153423.yaml @@ -0,0 +1,5 @@ +pr: 153423 +summary: Apply sequential read advice during merge for `bbq_disk` vectors +area: Vector Search +type: bug +issues: [] diff --git a/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/ES920DiskBBQVectorsReader.java b/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/ES920DiskBBQVectorsReader.java index 2d832d049b03a..36219007b8581 100644 --- a/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/ES920DiskBBQVectorsReader.java +++ b/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/ES920DiskBBQVectorsReader.java @@ -59,6 +59,15 @@ public class ES920DiskBBQVectorsReader extends IVFVectorsReader other, GenericFlatVectorReaders genericReaders) { + this.state = other.state; + this.fieldInfos = other.fieldInfos; + this.fields = other.fields; + this.genericReaders = genericReaders; + this.centroidExtension = other.centroidExtension; + this.clusterExtension = other.clusterExtension; + this.versionDirectIo = other.versionDirectIo; + this.dynamicVisitRatio = other.dynamicVisitRatio; + this.versionMeta = other.versionMeta; + this.ivfCentroids = other.ivfCentroids; + this.ivfClusters = other.ivfClusters; + } + public abstract CentroidIterator getCentroidIterator( FieldInfo fieldInfo, int numCentroids, @@ -284,6 +302,24 @@ public final void checkIntegrity() throws IOException { CodecUtil.checksumEntireFile(ivfClusters); } + @Override + public final KnnVectorsReader getMergeInstance() throws IOException { + // Flat vectors are opened with RANDOM advice for search but read sequentially during a + // merge, so back the merge instance with the flat readers' merge instances (which switch + // their input to SEQUENTIAL). finishMerge() reverts them. + return mergeInstance(genericReaders.getMergeInstance()); + } + + /** Builds a merge instance of this reader backed by the given flat vector merge readers. */ + protected abstract IVFVectorsReader mergeInstance(GenericFlatVectorReaders genericReaders); + + @Override + public final void finishMerge() throws IOException { + for (var reader : genericReaders.allReaders()) { + reader.finishMerge(); + } + } + protected FlatVectorsReader getReaderForField(String field) { FieldInfo info = fieldInfos.fieldInfo(field); if (info == null) throw new IllegalArgumentException("Could not find field [" + field + "]"); diff --git a/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/es94/ES940DiskBBQVectorsReader.java b/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/es94/ES940DiskBBQVectorsReader.java index ec7d1b19c1a64..1a8611e56fa10 100644 --- a/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/es94/ES940DiskBBQVectorsReader.java +++ b/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/es94/ES940DiskBBQVectorsReader.java @@ -73,6 +73,15 @@ public ES940DiskBBQVectorsReader(SegmentReadState state, GenericFlatVectorReader ); } + private ES940DiskBBQVectorsReader(ES940DiskBBQVectorsReader other, GenericFlatVectorReaders genericReaders) { + super(other, genericReaders); + } + + @Override + protected ES940DiskBBQVectorsReader mergeInstance(GenericFlatVectorReaders genericReaders) { + return new ES940DiskBBQVectorsReader(this, genericReaders); + } + CentroidIterator getPostingListPrefetchIterator(CentroidIterator centroidIterator, IndexInput postingListSlice) throws IOException { // TODO we may want to prefetch more than one postings list, however, we will likely want to place a limit // so we don't bother prefetching many lists we won't end up scoring diff --git a/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/es95/ES950DiskBBQVectorsReader.java b/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/es95/ES950DiskBBQVectorsReader.java index 676504dfb593f..e16d7fc5be9cd 100644 --- a/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/es95/ES950DiskBBQVectorsReader.java +++ b/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/es95/ES950DiskBBQVectorsReader.java @@ -76,6 +76,15 @@ public ES950DiskBBQVectorsReader(SegmentReadState state, GenericFlatVectorReader ); } + private ES950DiskBBQVectorsReader(ES950DiskBBQVectorsReader other, GenericFlatVectorReaders genericReaders) { + super(other, genericReaders); + } + + @Override + protected ES950DiskBBQVectorsReader mergeInstance(GenericFlatVectorReaders genericReaders) { + return new ES950DiskBBQVectorsReader(this, genericReaders); + } + CentroidIterator getPostingListPrefetchIterator(CentroidIterator centroidIterator, IndexInput postingListSlice) throws IOException { // TODO we may want to prefetch more than one postings list, however, we will likely want to place a limit // so we don't bother prefetching many lists we won't end up scoring diff --git a/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/next/ESNextDiskBBQVectorsReader.java b/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/next/ESNextDiskBBQVectorsReader.java index 81a076d512f69..4e2c90f6e1d70 100644 --- a/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/next/ESNextDiskBBQVectorsReader.java +++ b/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/next/ESNextDiskBBQVectorsReader.java @@ -77,6 +77,15 @@ public ESNextDiskBBQVectorsReader(SegmentReadState state, GenericFlatVectorReade ); } + private ESNextDiskBBQVectorsReader(ESNextDiskBBQVectorsReader other, GenericFlatVectorReaders genericReaders) { + super(other, genericReaders); + } + + @Override + protected ESNextDiskBBQVectorsReader mergeInstance(GenericFlatVectorReaders genericReaders) { + return new ESNextDiskBBQVectorsReader(this, genericReaders); + } + CentroidIterator getPostingListPrefetchIterator(CentroidIterator centroidIterator, IndexInput postingListSlice) throws IOException { // TODO we may want to prefetch more than one postings list, however, we will likely want to place a limit // so we don't bother prefetching many lists we won't end up scoring From 1fb071e25eeddd75cc3a91c862a07f05151a23ee Mon Sep 17 00:00:00 2001 From: Jim Ferenczi Date: Thu, 9 Jul 2026 16:49:59 +0200 Subject: [PATCH 2/5] Update docs/changelog/153423.yaml --- docs/changelog/153423.yaml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/docs/changelog/153423.yaml b/docs/changelog/153423.yaml index 4e238d7e903ee..fd24e92b4004d 100644 --- a/docs/changelog/153423.yaml +++ b/docs/changelog/153423.yaml @@ -1,5 +1,5 @@ -pr: 153423 -summary: Apply sequential read advice during merge for `bbq_disk` vectors area: Vector Search -type: bug issues: [] +pr: 153423 +summary: Apply sequential read advice during merge for `bbq_disk` vectors +type: enhancement From 2cf9d674eabb136af2f0e9430e76ce91c3a9d346 Mon Sep 17 00:00:00 2001 From: Jim Ferenczi Date: Thu, 6 Aug 2026 16:39:25 +0200 Subject: [PATCH 3/5] Propagate merge lifecycle through vector reader wrappers MergeReaderWrapper returned its merge reader without calling getMergeInstance() on it, and never implemented finishMerge(). Both calls therefore stopped at the wrapper, so a reader that is only reachable behind it never got the chance to prepare for, or recover from, a merge. This is the path taken whenever direct I/O is enabled for the raw vectors. ES818BinaryQuantizedVectorsReader had the other half of the problem: it built a merge instance from the raw reader but never told that reader the merge had finished, leaving whatever getMergeInstance() changed in place for subsequent searches. Delegate both calls in both readers. --- .../codec/vectors/MergeReaderWrapper.java | 9 +- .../ES818BinaryQuantizedVectorsReader.java | 5 + .../vectors/MergeReaderWrapperTests.java | 126 ++++++++++++++++++ 3 files changed, 138 insertions(+), 2 deletions(-) create mode 100644 server/src/test/java/org/elasticsearch/index/codec/vectors/MergeReaderWrapperTests.java diff --git a/server/src/main/java/org/elasticsearch/index/codec/vectors/MergeReaderWrapper.java b/server/src/main/java/org/elasticsearch/index/codec/vectors/MergeReaderWrapper.java index ebbadddecef30..cd5782cc0d435 100644 --- a/server/src/main/java/org/elasticsearch/index/codec/vectors/MergeReaderWrapper.java +++ b/server/src/main/java/org/elasticsearch/index/codec/vectors/MergeReaderWrapper.java @@ -75,8 +75,13 @@ public void search(String field, byte[] target, KnnCollector knnCollector, Accep } @Override - public FlatVectorsReader getMergeInstance() { - return mergeReader; + public FlatVectorsReader getMergeInstance() throws IOException { + return mergeReader.getMergeInstance(); + } + + @Override + public void finishMerge() throws IOException { + mergeReader.finishMerge(); } @Override diff --git a/server/src/main/java/org/elasticsearch/index/codec/vectors/es818/ES818BinaryQuantizedVectorsReader.java b/server/src/main/java/org/elasticsearch/index/codec/vectors/es818/ES818BinaryQuantizedVectorsReader.java index 591bd6dcc9403..0cb7100020017 100644 --- a/server/src/main/java/org/elasticsearch/index/codec/vectors/es818/ES818BinaryQuantizedVectorsReader.java +++ b/server/src/main/java/org/elasticsearch/index/codec/vectors/es818/ES818BinaryQuantizedVectorsReader.java @@ -146,6 +146,11 @@ public FlatVectorsReader getMergeInstance() throws IOException { return new ES818BinaryQuantizedVectorsReader(this, rawVectorsReader.getMergeInstance()); } + @Override + public void finishMerge() throws IOException { + rawVectorsReader.finishMerge(); + } + private void readFields(ChecksumIndexInput meta, FieldInfos infos) throws IOException { for (int fieldNumber = meta.readInt(); fieldNumber != -1; fieldNumber = meta.readInt()) { FieldInfo info = infos.fieldInfo(fieldNumber); diff --git a/server/src/test/java/org/elasticsearch/index/codec/vectors/MergeReaderWrapperTests.java b/server/src/test/java/org/elasticsearch/index/codec/vectors/MergeReaderWrapperTests.java new file mode 100644 index 0000000000000..be19a5e43cedd --- /dev/null +++ b/server/src/test/java/org/elasticsearch/index/codec/vectors/MergeReaderWrapperTests.java @@ -0,0 +1,126 @@ +/* + * Copyright Elasticsearch B.V. and/or licensed to Elasticsearch B.V. under one + * or more contributor license agreements. Licensed under the "Elastic License + * 2.0", the "GNU Affero General Public License v3.0 only", and the "Server Side + * Public License v 1"; you may not use this file except in compliance with, at + * your election, the "Elastic License 2.0", the "GNU Affero General Public + * License v3.0 only", or the "Server Side Public License, v 1". + */ + +package org.elasticsearch.index.codec.vectors; + +import org.apache.lucene.codecs.hnsw.FlatVectorsReader; +import org.apache.lucene.codecs.hnsw.FlatVectorsScorer; +import org.apache.lucene.index.ByteVectorValues; +import org.apache.lucene.index.FieldInfo; +import org.apache.lucene.index.FloatVectorValues; +import org.apache.lucene.search.AcceptDocs; +import org.apache.lucene.search.KnnCollector; +import org.apache.lucene.util.hnsw.RandomVectorScorer; +import org.elasticsearch.test.ESTestCase; + +import java.io.IOException; +import java.util.Map; + +/** + * {@link MergeReaderWrapper} serves searches from one reader and merges from another. Lucene brackets + * a merge with {@code getMergeInstance()} and {@code finishMerge()} on the reader it is handed, so + * both calls have to reach the reader that actually performs the merge. + */ +public class MergeReaderWrapperTests extends ESTestCase { + + /** Minimal reader that only records the merge lifecycle calls it receives. */ + private static class RecordingReader extends FlatVectorsReader { + int getMergeInstanceCalls; + int finishMergeCalls; + final FlatVectorsReader mergeInstance; + + RecordingReader(FlatVectorsReader mergeInstance) { + this.mergeInstance = mergeInstance == null ? this : mergeInstance; + } + + @Override + public FlatVectorsReader getMergeInstance() { + getMergeInstanceCalls++; + return mergeInstance; + } + + @Override + public void finishMerge() { + finishMergeCalls++; + } + + @Override + public FlatVectorsScorer getFlatVectorScorer(String field) { + throw new UnsupportedOperationException(); + } + + @Override + public RandomVectorScorer getRandomVectorScorer(String field, float[] target) { + throw new UnsupportedOperationException(); + } + + @Override + public RandomVectorScorer getRandomVectorScorer(String field, byte[] target) { + throw new UnsupportedOperationException(); + } + + @Override + public void checkIntegrity() {} + + @Override + public FloatVectorValues getFloatVectorValues(String field) { + return null; + } + + @Override + public ByteVectorValues getByteVectorValues(String field) { + return null; + } + + @Override + public void search(String field, float[] target, KnnCollector knnCollector, AcceptDocs acceptDocs) {} + + @Override + public void search(String field, byte[] target, KnnCollector knnCollector, AcceptDocs acceptDocs) {} + + @Override + public long ramBytesUsed() { + return 0; + } + + @Override + public Map getOffHeapByteSize(FieldInfo fieldInfo) { + return Map.of(); + } + + @Override + public void close() {} + } + + public void testGetMergeInstanceIsDelegatedToTheMergeReader() throws IOException { + RecordingReader mergeInstance = new RecordingReader(null); + RecordingReader mergeReader = new RecordingReader(mergeInstance); + RecordingReader mainReader = new RecordingReader(null); + + try (MergeReaderWrapper wrapper = new MergeReaderWrapper(mainReader, mergeReader)) { + assertSame(mergeInstance, wrapper.getMergeInstance()); + } + + assertEquals(1, mergeReader.getMergeInstanceCalls); + assertEquals("the search reader must not be asked for a merge instance", 0, mainReader.getMergeInstanceCalls); + } + + public void testFinishMergeIsDelegatedToTheMergeReader() throws IOException { + RecordingReader mergeReader = new RecordingReader(null); + RecordingReader mainReader = new RecordingReader(null); + + try (MergeReaderWrapper wrapper = new MergeReaderWrapper(mainReader, mergeReader)) { + wrapper.getMergeInstance(); + wrapper.finishMerge(); + } + + assertEquals(1, mergeReader.finishMergeCalls); + assertEquals("the search reader takes no part in the merge", 0, mainReader.finishMergeCalls); + } +} From 05e5bec485f3ac43cc2970528e63c3962ca3310a Mon Sep 17 00:00:00 2001 From: Jim Ferenczi Date: Fri, 7 Aug 2026 08:17:01 +0200 Subject: [PATCH 4/5] Broaden changelog to cover all merge lifecycle fixes --- docs/changelog/153423.yaml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/changelog/153423.yaml b/docs/changelog/153423.yaml index fd24e92b4004d..e2934c11bb8ea 100644 --- a/docs/changelog/153423.yaml +++ b/docs/changelog/153423.yaml @@ -1,5 +1,5 @@ area: Vector Search issues: [] pr: 153423 -summary: Apply sequential read advice during merge for `bbq_disk` vectors -type: enhancement +summary: Apply sequential read advice during vector merges +type: bug From 24644a12998c18f4859cecc1878ea5e694e98f66 Mon Sep 17 00:00:00 2001 From: Jim Ferenczi Date: Fri, 7 Aug 2026 12:06:03 +0200 Subject: [PATCH 5/5] Drop comments describing wrapped reader behaviour --- .../index/codec/vectors/diskbbq/IVFVectorsReader.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/IVFVectorsReader.java b/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/IVFVectorsReader.java index 5c545fd7fe620..6b771f07266e0 100644 --- a/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/IVFVectorsReader.java +++ b/server/src/main/java/org/elasticsearch/index/codec/vectors/diskbbq/IVFVectorsReader.java @@ -163,7 +163,7 @@ protected IVFVectorsReader( /** * Copy constructor used to build a merge instance: shares everything with {@code other} but uses - * the provided flat vector readers (their merge instances). + * the provided flat vector readers. */ protected IVFVectorsReader(IVFVectorsReader other, GenericFlatVectorReaders genericReaders) { this.state = other.state; @@ -337,9 +337,6 @@ public final void checkIntegrity() throws IOException { @Override public final KnnVectorsReader getMergeInstance() throws IOException { - // Flat vectors are opened with RANDOM advice for search but read sequentially during a - // merge, so back the merge instance with the flat readers' merge instances (which switch - // their input to SEQUENTIAL). finishMerge() reverts them. return mergeInstance(genericReaders.getMergeInstance()); }