From 0346cc65ec3f203a78085e7b2f24ef90b8545e2d Mon Sep 17 00:00:00 2001 From: Manasvi Goyal Date: Mon, 28 Sep 2026 11:54:07 -0700 Subject: [PATCH 1/3] changes to send fp16 bytes directly to gpu Signed-off-by: Manasvi Goyal --- .../NativeIndexBuildStrategyFactory.java | 9 +- .../DefaultVectorRepositoryAccessor.java | 7 +- .../remote/RemoteIndexBuildMetrics.java | 8 +- .../remote/VectorValuesInputStream.java | 17 ++-- .../NativeIndexBuildStrategyFactoryTests.java | 53 ++++++++++++ .../DefaultVectorRepositoryAccessorTests.java | 27 ++++++ .../KnnVectorValuesInputStreamTests.java | 82 ++++++++++++++++--- .../remote/RemoteIndexBuildMetricsTests.java | 31 +++++++ .../remote/RemoteIndexBuildStrategyTests.java | 60 ++++++++++++++ 9 files changed, 253 insertions(+), 41 deletions(-) diff --git a/src/main/java/org/opensearch/knn/index/codec/nativeindex/NativeIndexBuildStrategyFactory.java b/src/main/java/org/opensearch/knn/index/codec/nativeindex/NativeIndexBuildStrategyFactory.java index 2fb992c9c7..d6abc95333 100644 --- a/src/main/java/org/opensearch/knn/index/codec/nativeindex/NativeIndexBuildStrategyFactory.java +++ b/src/main/java/org/opensearch/knn/index/codec/nativeindex/NativeIndexBuildStrategyFactory.java @@ -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; @@ -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() diff --git a/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/DefaultVectorRepositoryAccessor.java b/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/DefaultVectorRepositoryAccessor.java index cc318a90d2..d021946ad4 100644 --- a/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/DefaultVectorRepositoryAccessor.java +++ b/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/DefaultVectorRepositoryAccessor.java @@ -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 diff --git a/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildMetrics.java b/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildMetrics.java index 99a7b28211..90d6f40e8c 100644 --- a/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildMetrics.java +++ b/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildMetrics.java @@ -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; @@ -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(); diff --git a/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/VectorValuesInputStream.java b/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/VectorValuesInputStream.java index 6b2d808cbd..c3f23d3ae1 100644 --- a/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/VectorValuesInputStream.java +++ b/src/main/java/org/opensearch/knn/index/codec/nativeindex/remote/VectorValuesInputStream.java @@ -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; @@ -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); @@ -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 @@ -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); diff --git a/src/test/java/org/opensearch/knn/index/codec/nativeindex/NativeIndexBuildStrategyFactoryTests.java b/src/test/java/org/opensearch/knn/index/codec/nativeindex/NativeIndexBuildStrategyFactoryTests.java index eee550ea3a..a0caaa808a 100644 --- a/src/test/java/org/opensearch/knn/index/codec/nativeindex/NativeIndexBuildStrategyFactoryTests.java +++ b/src/test/java/org/opensearch/knn/index/codec/nativeindex/NativeIndexBuildStrategyFactoryTests.java @@ -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; @@ -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 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 mockedExtractor = Mockito.mockStatic(FieldInfoExtractor.class); + MockedStatic mockedSettings = Mockito.mockStatic(KNNSettings.class); + MockedStatic 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 attributes = new HashMap<>(); diff --git a/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/DefaultVectorRepositoryAccessorTests.java b/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/DefaultVectorRepositoryAccessorTests.java index b8177daf81..d5f941990a 100644 --- a/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/DefaultVectorRepositoryAccessorTests.java +++ b/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/DefaultVectorRepositoryAccessorTests.java @@ -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> 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. */ diff --git a/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/KnnVectorValuesInputStreamTests.java b/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/KnnVectorValuesInputStreamTests.java index 120f9dacfb..df116d4252 100644 --- a/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/KnnVectorValuesInputStreamTests.java +++ b/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/KnnVectorValuesInputStreamTests.java @@ -8,6 +8,7 @@ import org.apache.lucene.search.DocIdSetIterator; import org.opensearch.knn.KNNTestCase; import org.opensearch.knn.index.VectorDataType; +import org.opensearch.knn.index.codec.util.KNNVectorAsCollectionOfHalfFloatsSerializer; import org.opensearch.knn.index.vectorvalues.KNNVectorValues; import org.opensearch.knn.index.vectorvalues.KNNVectorValuesFactory; import org.opensearch.knn.index.vectorvalues.TestVectorValues; @@ -101,9 +102,8 @@ public void testFloatVectorValuesInputStream() throws IOException { } /** - * Tests that reading half_float vectors out of a VectorValuesInputStream yields raw fp32 bytes, 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. + * Tests that reading half_float vectors out of a VectorValuesInputStream yields fp16 bytes (2 bytes per + * dimension), encoded on the data node. The remote build service consumes fp16 directly. */ public void testHalfFloatVectorValuesInputStream() throws IOException { int NUM_DOCS = randomIntBetween(1, 1000); @@ -122,20 +122,18 @@ public void testHalfFloatVectorValuesInputStream() throws IOException { // 1. Read all input stream bytes byte[] vectorStreamBytes = vectorValuesInputStream.readAllBytes(); - FloatBuffer vectorStreamFloats = ByteBuffer.wrap(vectorStreamBytes).order(ByteOrder.LITTLE_ENDIAN).asFloatBuffer(); - // 2. Bytes should be dim*4 per vector (fp32), not dim*2 (fp16) - assertEquals((long) NUM_DOCS * NUM_DIMENSION * Float.BYTES, vectorStreamBytes.length); + // 2. Bytes should be dim*2 per vector (fp16), not dim*4 (fp32) + assertEquals((long) NUM_DOCS * NUM_DIMENSION * Short.BYTES, vectorStreamBytes.length); - // 3. Content should match the source vectors exactly, unconverted. - FloatBuffer expectedBuffer = ByteBuffer.allocate(NUM_DOCS * NUM_DIMENSION * Float.BYTES) - .order(ByteOrder.LITTLE_ENDIAN) - .asFloatBuffer(); + // 3. Content should be the fp16 encoding of each source vector, in order + ByteBuffer expectedBuffer = ByteBuffer.allocate(NUM_DOCS * NUM_DIMENSION * Short.BYTES); + byte[] encoded = new byte[NUM_DIMENSION * Short.BYTES]; for (float[] vector : vectorValues) { - expectedBuffer.put(vector); + KNNVectorAsCollectionOfHalfFloatsSerializer.INSTANCE.floatToByteArray(vector, encoded, NUM_DIMENSION); + expectedBuffer.put(encoded); } - expectedBuffer.position(0); - assertEquals(expectedBuffer, vectorStreamFloats); + assertArrayEquals(expectedBuffer.array(), vectorStreamBytes); } public void testByteVectorValuesInputStream() throws IOException { @@ -227,6 +225,64 @@ public void testMultiPartVectorValueInputStream() throws IOException { assertArrayEquals(expectedStream.readAllBytes(), testBuffer.array()); } + public void testMultiPartHalfFloatVectorValueInputStream() throws IOException { + final int NUM_DOCS = randomIntBetween(100, 1000); + final int NUM_DIMENSION = randomIntBetween(1, 1000); + final int NUM_PARTS = randomIntBetween(1, NUM_DOCS / 10); + + List vectorValues = getRandomFloatVectors(NUM_DOCS, NUM_DIMENSION); + final Supplier> knnVectorValuesSupplier = () -> KNNVectorValuesFactory.getVectorValues( + VectorDataType.HALF_FLOAT, + new TestVectorValues.PreDefinedFloatVectorValues(vectorValues) + ); + + final KNNVectorValues knnVectorValues = knnVectorValuesSupplier.get(); + initializeVectorValues(knnVectorValues); + final int vectorBlobLength = knnVectorValues.bytesPerVector() * NUM_DOCS; + assertEquals(NUM_DIMENSION * Short.BYTES * NUM_DOCS, vectorBlobLength); + final int PART_SIZE = vectorBlobLength / NUM_PARTS; + final int LAST_PART_SIZE = (vectorBlobLength % PART_SIZE) != 0 ? vectorBlobLength % PART_SIZE : PART_SIZE; + + final List streamList = new ArrayList<>(NUM_PARTS); + for (int partNumber = 0; partNumber < NUM_PARTS; partNumber++) { + streamList.add( + new VectorValuesInputStream( + knnVectorValuesSupplier.get(), + VectorDataType.HALF_FLOAT, + (long) partNumber * PART_SIZE, + PART_SIZE + ) + ); + } + if (LAST_PART_SIZE != PART_SIZE) { + streamList.add( + new VectorValuesInputStream( + knnVectorValuesSupplier.get(), + VectorDataType.HALF_FLOAT, + vectorBlobLength - LAST_PART_SIZE, + LAST_PART_SIZE + ) + ); + } + + ByteBuffer testBuffer = ByteBuffer.allocate(vectorBlobLength); + for (VectorValuesInputStream stream : streamList) { + testBuffer.put(stream.readAllBytes()); + } + assertEquals(vectorBlobLength, testBuffer.position()); + + // Expected: sequential stream, and independently the serializer output vector by vector + VectorValuesInputStream expectedStream = new VectorValuesInputStream(knnVectorValuesSupplier.get(), VectorDataType.HALF_FLOAT); + assertArrayEquals(expectedStream.readAllBytes(), testBuffer.array()); + ByteBuffer serialized = ByteBuffer.allocate(vectorBlobLength); + byte[] encoded = new byte[NUM_DIMENSION * Short.BYTES]; + for (float[] vector : vectorValues) { + KNNVectorAsCollectionOfHalfFloatsSerializer.INSTANCE.floatToByteArray(vector, encoded, NUM_DIMENSION); + serialized.put(encoded); + } + assertArrayEquals(serialized.array(), testBuffer.array()); + } + /** * Tests that invoking {@link VectorValuesInputStream#read()} N times yields the same results as {@link VectorValuesInputStream#read(byte[], 0, N)} */ diff --git a/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildMetricsTests.java b/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildMetricsTests.java index 203300e305..9e9ac65dc3 100644 --- a/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildMetricsTests.java +++ b/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildMetricsTests.java @@ -8,6 +8,7 @@ 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.KNNVectorValuesFactory; import org.opensearch.knn.plugin.stats.KNNRemoteIndexBuildValue; import java.io.IOException; @@ -93,6 +94,36 @@ public void testEndMetricsMergeUpdatesMergeTimerAndGauges() throws IOException { assertEquals(0L, (long) REMOTE_INDEX_BUILD_FLUSH_TIME.getValue()); } + /** + * The size gauge reflects the upload layout: fp32 for FLOAT, fp16 (2 bytes per dimension) for HALF_FLOAT. + */ + public void testStartMetricsSizeReflectsVectorDataType() throws IOException { + int docs = buildIndexParams.getTotalLiveDocs(); + int dimension = 2; + + RemoteIndexBuildMetrics floatMetrics = new RemoteIndexBuildMetrics(); + floatMetrics.startRemoteIndexBuildMetrics(buildParamsWithFlush(true)); + assertEquals((long) docs * dimension * Float.BYTES, (long) REMOTE_INDEX_BUILD_CURRENT_FLUSH_SIZE.getValue()); + floatMetrics.endRemoteIndexBuildMetrics(BuildResult.SUCCESS); + + BuildIndexParams halfFloatParams = BuildIndexParams.builder() + .indexOutputWithBuffer(indexOutputWithBuffer) + .knnEngine(buildIndexParams.getKnnEngine()) + .field(buildIndexParams.getField()) + .vectorDataType(VectorDataType.HALF_FLOAT) + .indexParameters(buildIndexParams.getIndexParameters()) + .knnVectorValuesSupplier(KNNVectorValuesFactory.getVectorValuesSupplier(VectorDataType.HALF_FLOAT, randomVectorValues)) + .totalLiveDocs(docs) + .segmentWriteState(segmentWriteState) + .isFlush(true) + .build(); + RemoteIndexBuildMetrics halfFloatMetrics = new RemoteIndexBuildMetrics(); + halfFloatMetrics.startRemoteIndexBuildMetrics(halfFloatParams); + assertEquals((long) docs * dimension * Short.BYTES, (long) REMOTE_INDEX_BUILD_CURRENT_FLUSH_SIZE.getValue()); + halfFloatMetrics.endRemoteIndexBuildMetrics(BuildResult.SUCCESS); + assertEquals(0L, (long) REMOTE_INDEX_BUILD_CURRENT_FLUSH_SIZE.getValue()); + } + /** * Runs a full start/end metrics cycle for the given {@link BuildResult} and asserts that only {@code expectedCounter} * among the four outcome counters was incremented. diff --git a/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildStrategyTests.java b/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildStrategyTests.java index 2bc21aa646..e705ac14c0 100644 --- a/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildStrategyTests.java +++ b/src/test/java/org/opensearch/knn/index/codec/nativeindex/remote/RemoteIndexBuildStrategyTests.java @@ -30,6 +30,7 @@ import org.opensearch.knn.index.engine.ResolvedIndexSpec; import org.opensearch.knn.index.mapper.CompressionLevel; import org.opensearch.knn.index.store.IndexOutputWithBuffer; +import org.opensearch.knn.index.vectorvalues.KNNVectorValuesFactory; import org.opensearch.knn.plugin.stats.KNNRemoteIndexBuildValue; import org.opensearch.remoteindexbuild.model.RemoteBuildRequest; import org.opensearch.repositories.RepositoriesService; @@ -551,6 +552,65 @@ public void testBuildRequestFP16() throws IOException { assertFalse(request.isSkipStoredVectors()); } + public void testBuildRequestHalfFloatFlat() throws IOException { + ResolvedIndexSpec resolvedSpec = ResolvedIndexSpec.builder() + .engine(KNNEngine.FAISS) + .methodName("hnsw") + .encoderType(Encoder.EncoderType.FLAT) + .compressionLevel(CompressionLevel.NOT_CONFIGURED) + .vectorDataType(VectorDataType.HALF_FLOAT) + .dimension(2) + .build(); + RemoteBuildRequest request = RemoteIndexBuildStrategy.buildRemoteBuildRequest( + createTestIndexSettings(), + halfFloatBuildIndexParams(), + createTestRepositoryMetadata(), + MOCK_FULL_PATH, + getMockParameterMap(), + resolvedSpec + ); + assertEquals(VectorDataType.HALF_FLOAT.getValue(), request.getVectorDataType()); + assertEquals(2, request.getDimension()); + assertEquals(3, request.getDocCount()); + assertFalse(request.isSkipStoredVectors()); + } + + public void testBuildRequestHalfFloatSQOneBit() throws IOException { + ResolvedIndexSpec resolvedSpec = ResolvedIndexSpec.builder() + .engine(KNNEngine.FAISS) + .methodName("hnsw") + .encoderType(Encoder.EncoderType.SQ) + .quantizationBits(Encoder.QuantizationBits.ONE) + .compressionLevel(CompressionLevel.x16) + .vectorDataType(VectorDataType.HALF_FLOAT) + .dimension(2) + .build(); + RemoteBuildRequest request = RemoteIndexBuildStrategy.buildRemoteBuildRequest( + createTestIndexSettings(), + halfFloatBuildIndexParams(), + createTestRepositoryMetadata(), + MOCK_FULL_PATH, + getMockSQOneBitParameterMap(), + resolvedSpec + ); + assertEquals(VectorDataType.HALF_FLOAT.getValue(), request.getVectorDataType()); + assertTrue(request.isSkipStoredVectors()); + } + + private BuildIndexParams halfFloatBuildIndexParams() { + return BuildIndexParams.builder() + .indexOutputWithBuffer(indexOutputWithBuffer) + .knnEngine(KNNEngine.FAISS) + .field(buildIndexParams.getField()) + .vectorDataType(VectorDataType.HALF_FLOAT) + .indexParameters(buildIndexParams.getIndexParameters()) + .knnVectorValuesSupplier(KNNVectorValuesFactory.getVectorValuesSupplier(VectorDataType.HALF_FLOAT, randomVectorValues)) + .totalLiveDocs(buildIndexParams.getTotalLiveDocs()) + .segmentWriteState(segmentWriteState) + .isFlush(randomBoolean()) + .build(); + } + public Map getMockParameterMap() { Map encoderSq = Map.of(ENCODER_SQ, Map.of()); Map encoderMap = Map.of(METHOD_ENCODER_PARAMETER, encoderSq); From 8738e4e6999a7ac041d81410e4b58f3bf386bf23 Mon Sep 17 00:00:00 2001 From: Manasvi Goyal Date: Mon, 28 Sep 2026 14:11:04 -0700 Subject: [PATCH 2/3] add changelog Signed-off-by: Manasvi Goyal --- CHANGELOG.md | 1 + 1 file changed, 1 insertion(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 074ae07651..92d26f430e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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) From ea431df9f4ea7ce3e073bd668365c666db8b1885 Mon Sep 17 00:00:00 2001 From: Manasvi Goyal Date: Thu, 1 Oct 2026 14:20:51 -0700 Subject: [PATCH 3/3] use api-snapshot remote index builder image in remote build CI Signed-off-by: Manasvi Goyal --- .github/workflows/remote_index_build.yml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/.github/workflows/remote_index_build.yml b/.github/workflows/remote_index_build.yml index 7dd32d7f19..85002e6072 100644 --- a/.github/workflows/remote_index_build.yml +++ b/.github/workflows/remote_index_build.yml @@ -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: | @@ -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