Skip to content
Open
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
4 changes: 2 additions & 2 deletions .github/workflows/remote_index_build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ jobs:

- name: Pull Remote Index Build Docker Image from Docker Hub
run: |
docker pull opensearchstaging/remote-vector-index-builder:api-latest
docker pull opensearchstaging/remote-vector-index-builder:api-snapshot

- name: Pull LocalStack Docker image
run: |
Expand All @@ -108,7 +108,7 @@ jobs:

- name: Run Docker container
run: |
docker run --rm -d --name remote-index-builder-container --gpus all -p 80:1025 -e S3_ENDPOINT_URL=http://172.17.0.1:4566 -e AWS_ACCESS_KEY_ID=${{ env.AWS_ACCESS_KEY_ID }} -e AWS_SECRET_ACCESS_KEY=${{ env.AWS_SECRET_ACCESS_KEY }} -e AWS_DEFAULT_REGION=${{env.AWS_DEFAULT_REGION}} -e AWS_SESSION_TOKEN=${{env.AWS_SESSION_TOKEN}} opensearchstaging/remote-vector-index-builder:api-latest
docker run --rm -d --name remote-index-builder-container --gpus all -p 80:1025 -e S3_ENDPOINT_URL=http://172.17.0.1:4566 -e AWS_ACCESS_KEY_ID=${{ env.AWS_ACCESS_KEY_ID }} -e AWS_SECRET_ACCESS_KEY=${{ env.AWS_SECRET_ACCESS_KEY }} -e AWS_DEFAULT_REGION=${{env.AWS_DEFAULT_REGION}} -e AWS_SESSION_TOKEN=${{env.AWS_SESSION_TOKEN}} opensearchstaging/remote-vector-index-builder:api-snapshot
sleep 5

- name: Run tests
Expand Down
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -59,3 +59,4 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
* Skip warmup for warm-tier indices to avoid unnecessary graph loading from remote store [#3565](https://github.com/opensearch-project/k-NN/pull/3565)
* Avoid FP16 -> FP32 -> FP16 round trip when merging `half_float` segments [#3610](https://github.com/opensearch-project/k-NN/pull/3610)
* Report `exact_search` timings in the Profile API for nested k-NN queries with `expand_nested_docs` on the Lucene engine [#3579](https://github.com/opensearch-project/k-NN/pull/3579)
* Send `half_float` vectors to the remote vector index build service as fp16 instead of fp32 [#3608](https://github.com/opensearch-project/k-NN/pull/3608)
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,6 @@
import org.apache.lucene.index.FieldInfo;
import org.opensearch.index.IndexSettings;
import org.opensearch.knn.common.FieldInfoExtractor;
import org.opensearch.knn.index.VectorDataType;
import org.opensearch.knn.index.codec.nativeindex.remote.RemoteIndexBuildStrategy;
import org.opensearch.knn.index.engine.KNNEngine;
import org.opensearch.knn.index.engine.KNNLibraryIndexingContext;
Expand Down Expand Up @@ -75,13 +74,7 @@ public NativeIndexBuildStrategy getBuildStrategy(
}

initializeVectorValues(knnVectorValues);
// HALF_FLOAT is uploaded as raw fp32 for remote build (see VectorValuesInputStream#reloadBuffer) -
// bytesPerVector() reports the 2-byte on-disk size, which understates the actual upload size here.
final VectorDataType vectorDataType = FieldInfoExtractor.extractVectorDataType(fieldInfo);
long bytesPerVectorForUpload = vectorDataType == VectorDataType.HALF_FLOAT
? (long) knnVectorValues.dimension() * Float.BYTES
: knnVectorValues.bytesPerVector();
long vectorBlobLength = bytesPerVectorForUpload * totalLiveDocs;
long vectorBlobLength = (long) knnVectorValues.bytesPerVector() * totalLiveDocs;

if (totalLiveDocs > MIN_DOCS_FOR_REMOTE_INDEX_BUILD
&& isKNNRemoteVectorBuildEnabled()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,12 +67,7 @@ public void writeToRepository(
assert blobContainer != null;
KNNVectorValues<?> knnVectorValues = knnVectorValuesSupplier.get();
initializeVectorValues(knnVectorValues);
// HALF_FLOAT is uploaded as raw fp32 (see VectorValuesInputStream#reloadBuffer) - bytesPerVector()
// reports the 2-byte on-disk size, which does not match the actual upload size here.
long bytesPerVectorForUpload = vectorDataType == VectorDataType.HALF_FLOAT
? (long) knnVectorValues.dimension() * Float.BYTES
: knnVectorValues.bytesPerVector();
long vectorBlobLength = bytesPerVectorForUpload * totalLiveDocs;
long vectorBlobLength = (long) knnVectorValues.bytesPerVector() * totalLiveDocs;

// TODO : Once Lucene patch https://github.com/apache/lucene/issues/14992 is merged, remove vector data type check in condition.
// That issue is specific to MergedByteVectorValues (BYTE/BINARY); HALF_FLOAT rides on Lucene's
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@

import lombok.extern.log4j.Log4j2;
import org.opensearch.common.StopWatch;
import org.opensearch.knn.index.VectorDataType;
import org.opensearch.knn.index.codec.nativeindex.model.BuildIndexParams;
import org.opensearch.knn.index.codec.nativeindex.remote.RemoteIndexBuildStrategy.BuildResult;
import org.opensearch.knn.index.vectorvalues.KNNVectorValues;
Expand Down Expand Up @@ -64,12 +63,7 @@ public RemoteIndexBuildMetrics() {
public void startRemoteIndexBuildMetrics(BuildIndexParams indexInfo) throws IOException {
KNNVectorValues<?> knnVectorValues = indexInfo.getKnnVectorValuesSupplier().get();
initializeVectorValues(knnVectorValues);
// HALF_FLOAT is uploaded as raw fp32 for remote build (see VectorValuesInputStream#reloadBuffer) -
// bytesPerVector() reports the 2-byte on-disk size, which understates the actual upload size here.
long bytesPerVectorForUpload = indexInfo.getVectorDataType() == VectorDataType.HALF_FLOAT
? (long) knnVectorValues.dimension() * Float.BYTES
: knnVectorValues.bytesPerVector();
this.size = (long) indexInfo.getTotalLiveDocs() * bytesPerVectorForUpload;
this.size = (long) indexInfo.getTotalLiveDocs() * knnVectorValues.bytesPerVector();
this.isFlush = indexInfo.isFlush();
this.fieldName = indexInfo.getField();
overallStopWatch.start();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
import lombok.extern.log4j.Log4j2;
import org.apache.lucene.search.DocIdSetIterator;
import org.opensearch.knn.index.VectorDataType;
import org.opensearch.knn.index.codec.util.KNNVectorAsCollectionOfHalfFloatsSerializer;
import org.opensearch.knn.index.vectorvalues.KNNBinaryVectorValues;
import org.opensearch.knn.index.vectorvalues.KNNByteVectorValues;
import org.opensearch.knn.index.vectorvalues.KNNFloatVectorValues;
Expand Down Expand Up @@ -38,6 +39,8 @@ class VectorValuesInputStream extends InputStream {
// will be filled 1 vector at a time.
private ByteBuffer currentBuffer;
private final int bytesPerVector;
// Reused per vector: fp16 encoding of the current vector, only allocated for HALF_FLOAT
private final byte[] halfFloatVectorBytes;
private long bytesRemaining;
private final VectorDataType vectorDataType;
private final AtomicBoolean closed = new AtomicBoolean(false);
Expand Down Expand Up @@ -67,9 +70,8 @@ public VectorValuesInputStream(KNNVectorValues<?> knnVectorValues, VectorDataTyp
this.knnVectorValues = knnVectorValues;
this.vectorDataType = vectorDataType;
initializeVectorValues(this.knnVectorValues);
this.bytesPerVector = vectorDataType == HALF_FLOAT
? this.knnVectorValues.dimension() * Float.BYTES
: this.knnVectorValues.bytesPerVector();
this.bytesPerVector = this.knnVectorValues.bytesPerVector();
this.halfFloatVectorBytes = vectorDataType == HALF_FLOAT ? new byte[bytesPerVector] : null;
// We use currentBuffer == null to indicate that there are no more vectors to be read
this.currentBuffer = ByteBuffer.allocate(bytesPerVector).order(ByteOrder.LITTLE_ENDIAN);
// Position the InputStream at the specific byte within the specific vector that startPosition references
Expand Down Expand Up @@ -207,11 +209,12 @@ private void reloadBuffer() throws IOException {
float[] floatVector = ((KNNFloatVectorValues) knnVectorValues).getVector();
currentBuffer.asFloatBuffer().put(floatVector);
} else if (vectorDataType == HALF_FLOAT) {
// Uploaded as raw fp32, same as FLOAT - the remote build service converts fp32 -> fp16 itself
// while streaming (FP32ToFP16ConvertingBytesIO), the same way it already does for the existing
// FLOAT+sq,bits:16 case. Do not encode to fp16 bytes here - the service does not expect that.
// Encoded to fp16 (2 bytes per dimension) on the data node; the remote build service consumes fp16
// directly. This is the HALF_FLOAT field type only - a FLOAT field with the fp16 SQ encoder is
// still uploaded as raw fp32 and converted by the service.
float[] floatVector = ((KNNHalfFloatVectorValues) knnVectorValues).getVector();
currentBuffer.asFloatBuffer().put(floatVector);
KNNVectorAsCollectionOfHalfFloatsSerializer.INSTANCE.floatToByteArray(floatVector, halfFloatVectorBytes, floatVector.length);
currentBuffer.put(halfFloatVectorBytes);
} else if (vectorDataType == BYTE) {
byte[] byteVector = ((KNNByteVectorValues) knnVectorValues).getVector();
currentBuffer.put(byteVector);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,15 +22,20 @@
import org.opensearch.knn.index.engine.KNNLibraryIndexingContext;
import org.opensearch.knn.index.engine.ResolvedIndexSpec;
import org.opensearch.knn.index.vectorvalues.KNNVectorValues;
import org.opensearch.knn.index.vectorvalues.KNNVectorValuesFactory;
import org.opensearch.knn.index.vectorvalues.TestVectorValues;
import org.opensearch.repositories.RepositoriesService;

import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.function.Supplier;

import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
import static org.opensearch.knn.common.KNNConstants.METHOD_HNSW;
Expand Down Expand Up @@ -131,6 +136,54 @@ public void testGetBuildStrategy_nonFaissEngine_returnsDefault() {
}
}

@SneakyThrows
public void testGetBuildStrategy_halfFloat_thenGateSeesFp16UploadSize() {
when(fieldInfo.attributes()).thenReturn(new HashMap<>());
final int dimension = 128;
final int totalLiveDocs = 10;
final List<float[]> vectors = new ArrayList<>();
for (int i = 0; i < totalLiveDocs; i++) {
vectors.add(TestVectorValues.getRandomVector(dimension));
}
// Real half_float vector values, so bytesPerVector() is the codec's own 2 bytes per dimension
final KNNVectorValues<?> halfFloatValues = KNNVectorValuesFactory.getVectorValues(
VectorDataType.HALF_FLOAT,
new TestVectorValues.PreDefinedFloatVectorValues(vectors)
);

try (
MockedStatic<FieldInfoExtractor> mockedExtractor = Mockito.mockStatic(FieldInfoExtractor.class);
MockedStatic<KNNSettings> mockedSettings = Mockito.mockStatic(KNNSettings.class);
MockedStatic<RemoteIndexBuildStrategy> mockedRemote = Mockito.mockStatic(RemoteIndexBuildStrategy.class)
) {
mockedExtractor.when(() -> FieldInfoExtractor.extractKNNEngine(fieldInfo)).thenReturn(KNNEngine.FAISS);
mockedSettings.when(KNNSettings::isKNNRemoteVectorBuildEnabled).thenReturn(true);
mockedRemote.when(() -> RemoteIndexBuildStrategy.shouldBuildIndexRemotely(any(IndexSettings.class), anyLong(), anyInt()))
.thenReturn(true);

when(knnLibraryIndexingContext.getResolvedSpec()).thenReturn(
ResolvedIndexSpec.builder()
.engine(KNNEngine.FAISS)
.methodName(METHOD_HNSW)
.encoderType(Encoder.EncoderType.FLAT)
.vectorDataType(VectorDataType.HALF_FLOAT)
.dimension(dimension)
.indexVersionCreated(Version.CURRENT)
.build()
);
NativeIndexBuildStrategyFactory factory = new NativeIndexBuildStrategyFactory(repositoriesServiceSupplier, indexSettings);
factory.setKnnLibraryIndexingContext(knnLibraryIndexingContext);

NativeIndexBuildStrategy strategy = factory.getBuildStrategy(fieldInfo, totalLiveDocs, halfFloatValues);

assertTrue(strategy instanceof RemoteIndexBuildStrategy);
final long fp16UploadBytes = (long) totalLiveDocs * dimension * Short.BYTES;
mockedRemote.verify(
() -> RemoteIndexBuildStrategy.shouldBuildIndexRemotely(any(IndexSettings.class), eq(fp16UploadBytes), eq(dimension))
);
}
}

@SneakyThrows
public void testGetBuildStrategy_remoteConditionsMet_returnsRemoteStrategy() {
Map<String, String> attributes = new HashMap<>();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,33 @@ public void testRepositoryInteractionWithBlobContainer_halfFloat() throws IOExce
verify(testContainer).writeBlob(eq(BLOB_NAME + DOC_ID_FILE_EXTENSION), any(), eq((long) NUM_DOCS * Integer.BYTES), eq(true));
}

/**
* half_float vectors are uploaded as fp16, so the vector blob length must be 2 bytes per dimension.
*/
public void testRepositoryInteraction_halfFloatBlobLengthIsFp16() throws IOException, InterruptedException {
BlobPath testBasePath = new BlobPath().add("testBasePath");
BlobContainer testContainer = Mockito.spy(new TestBlobContainer(mock(FsBlobStore.class), testBasePath, mock(Path.class)));
VectorRepositoryAccessor objectUnderTest = new DefaultVectorRepositoryAccessor(testContainer);

String BLOB_NAME = "test_blob";
int NUM_DOCS = 100;
Supplier<KNNVectorValues<?>> supplier = KNNVectorValuesFactory.getVectorValuesSupplier(
VectorDataType.HALF_FLOAT,
randomVectorValues
);
objectUnderTest.writeToRepository(BLOB_NAME, NUM_DOCS, VectorDataType.HALF_FLOAT, supplier);

KNNVectorValues<?> knnVectorValues = supplier.get();
initializeVectorValues(knnVectorValues);
verify(testContainer).writeBlob(
eq(BLOB_NAME + VECTOR_BLOB_FILE_EXTENSION),
any(),
eq((long) NUM_DOCS * knnVectorValues.dimension() * Short.BYTES),
eq(true)
);
verify(testContainer).writeBlob(eq(BLOB_NAME + DOC_ID_FILE_EXTENSION), any(), eq((long) NUM_DOCS * Integer.BYTES), eq(true));
}

/**
* Test that when an exception is thrown during asyncBlobUpload, the exception is rethrown.
*/
Expand Down
Loading
Loading