Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions docs/changelog/153423.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
area: Vector Search
issues: []
pr: 153423
summary: Apply sequential read advice during vector merges
type: bug
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,15 @@ public class ES920DiskBBQVectorsReader extends IVFVectorsReader<IVFVectorsReader
);
}

private ES920DiskBBQVectorsReader(ES920DiskBBQVectorsReader other, GenericFlatVectorReaders genericReaders) {
super(other, genericReaders);
}

@Override
protected ES920DiskBBQVectorsReader mergeInstance(GenericFlatVectorReaders genericReaders) {
return new ES920DiskBBQVectorsReader(this, genericReaders);
}

public CentroidIterator getPostingListPrefetchIterator(CentroidIterator centroidIterator, IndexInput postingListSlice)
throws IOException {
return new PrefetchingCentroidIterator(centroidIterator, postingListSlice);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,24 @@ protected IVFVectorsReader(
}
}

/**
* Copy constructor used to build a merge instance: shares everything with {@code other} but uses
* the provided flat vector readers.
*/
protected IVFVectorsReader(IVFVectorsReader<E> 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,
Expand Down Expand Up @@ -317,6 +335,21 @@ public final void checkIntegrity() throws IOException {
CodecUtil.checksumEntireFile(ivfClusters);
}

@Override
public final KnnVectorsReader getMergeInstance() throws IOException {
return mergeInstance(genericReaders.getMergeInstance());
}

/** Builds a merge instance of this reader backed by the given flat vector merge readers. */
protected abstract IVFVectorsReader<E> 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 + "]");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
@@ -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<String, Long> 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);
}
}
Loading