From 52294efdb0a2625caeed8880622a76dbcc59f788 Mon Sep 17 00:00:00 2001 From: nvzm123 Date: Fri, 18 Sep 2026 10:40:06 +0000 Subject: [PATCH 1/2] Parallelize bounded HNSW graph post-processing --- .../cuvs/lucene/AcceleratedHNSWParams.java | 8 +- .../cuvs/lucene/AcceleratedHNSWUtils.java | 225 ++++++++++++++---- .../cuvs/lucene/BoundedParallelExecutor.java | 203 ++++++++++++++++ .../nvidia/cuvs/lucene/GPUBuiltHnswGraph.java | 201 +++++++++++++--- .../Lucene99AcceleratedHNSWVectorsWriter.java | 6 +- ...ratedHNSWBinaryQuantizedVectorsWriter.java | 6 +- ...ratedHNSWScalarQuantizedVectorsWriter.java | 6 +- .../cuvs/lucene/IntGraphTestMatrix.java | 134 +++++++++++ .../lucene/TestBoundedParallelExecutor.java | 219 +++++++++++++++++ ...TestWriterThreadsGraphMaterialization.java | 190 +++++++++++++++ .../TestWriterThreadsGraphSerialization.java | 214 +++++++++++++++++ .../TestWriterThreadsPersistedIndex.java | 119 +++++++++ 12 files changed, 1451 insertions(+), 80 deletions(-) create mode 100644 java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/BoundedParallelExecutor.java create mode 100644 java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/IntGraphTestMatrix.java create mode 100644 java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestBoundedParallelExecutor.java create mode 100644 java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsGraphMaterialization.java create mode 100644 java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsGraphSerialization.java create mode 100644 java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsPersistedIndex.java diff --git a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWParams.java b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWParams.java index a5f164b70b..1e467b773a 100644 --- a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWParams.java +++ b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWParams.java @@ -95,7 +95,7 @@ public static enum Strategy { /** * Constructs an instance of {@link AcceleratedHNSWParams} with specific parameter values. * - * @param writerThreads Number of cuVS writer threads to use. + * @param writerThreads Maximum threads to use for cuVS writes and post-build graph processing. * @param intermediateGraphDegree The intermediate graph degree while building the CAGRA index. * @param graphdegree The graph degree to use while building the CAGRA index. * @param hnswLayers The number of HNSW layers to build in the HNSW index. @@ -143,9 +143,9 @@ private AcceleratedHNSWParams( } /** - * Get the cuVS writer threads parameter + * Get the maximum thread count for cuVS writes and post-build graph processing. * - * @return cuVS writer threads parameter + * @return maximum writer and graph-processing thread count */ public int getWriterThreads() { return writerThreads; @@ -326,7 +326,7 @@ public static class Builder { private HnswHeuristicType hnswHeuristicType = DEFAULT_HNSW_HEURISTIC_TYPE; /** - * Set the number of cuVS writer threads while building the index + * Set the maximum number of threads used for cuVS writes and post-build graph processing. * Valid range - Minimum: {@value MIN_WRITER_THREADS}, Maximum: {@value MAX_WRITER_THREADS} * Default value - {@value DEFAULT_WRITER_THREADS} * diff --git a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java index 9c49c07fe0..6ce7dda331 100644 --- a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java +++ b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java @@ -21,6 +21,8 @@ import java.util.TreeSet; import org.apache.lucene.index.FieldInfo; import org.apache.lucene.index.VectorSimilarityFunction; +import org.apache.lucene.store.ByteBuffersDataOutput; +import org.apache.lucene.store.DataOutput; import org.apache.lucene.store.IndexOutput; import org.apache.lucene.util.InfoStream; import org.apache.lucene.util.hnsw.HnswGraph; @@ -89,6 +91,29 @@ public static GPUBuiltHnswGraph createMultiLayerHnswGraph( CagraIndexParams params, QuantizationType quantization) throws Throwable { + return createMultiLayerHnswGraph( + fieldInfo, + size, + dimensions, + adjacencyListMatrix, + vectors, + hnswLayers, + params, + quantization, + 1); + } + + static GPUBuiltHnswGraph createMultiLayerHnswGraph( + FieldInfo fieldInfo, + int size, + int dimensions, + CuVSMatrix adjacencyListMatrix, + List vectors, + int hnswLayers, + CagraIndexParams params, + QuantizationType quantization, + int requestedWorkers) + throws Throwable { int M = Math.ceilDiv((int) adjacencyListMatrix.columns(), 2); @@ -166,7 +191,8 @@ public static GPUBuiltHnswGraph createMultiLayerHnswGraph( } // Create the multi-layer graph with all layers - return new GPUBuiltHnswGraph(size, dimensions, layerNodes, layerAdjacencies); + return GPUBuiltHnswGraph.create( + size, dimensions, layerNodes, layerAdjacencies, requestedWorkers); } /** @@ -237,54 +263,169 @@ private static CuVSMatrix buildCagraGraphForSubset( */ public static int[][] writeGraph(GPUBuiltHnswGraph graph, IndexOutput vectorIndex) throws IOException { - // write vectors' neighbors on each level into the vectorIndex file + return writeGraph(graph, vectorIndex, 1); + } + + static int[][] writeGraph(GPUBuiltHnswGraph graph, IndexOutput vectorIndex, int requestedWorkers) + throws IOException { + return writeGraph( + graph, vectorIndex, requestedWorkers, Runtime.getRuntime().availableProcessors()); + } + + static int[][] writeGraph( + GPUBuiltHnswGraph graph, + IndexOutput vectorIndex, + int requestedWorkers, + int availableProcessors) + throws IOException { int countOnLevel0 = graph.size(); - int[][] offsets = new int[graph.numLevels()][]; - int[] scratch = new int[graph.maxConn() * 2]; - for (int level = 0; level < graph.numLevels(); level++) { + int numLevels = graph.numLevels(); + int[][] offsets = new int[numLevels][]; + int maxConn = graph.maxConn(); + + int[] level0Nodes = NodesIterator.getSortedNodes(graph.getNodesOnLevel(0)); + offsets[0] = new int[level0Nodes.length]; + if (requestedWorkers > 1 && level0Nodes.length >= PARALLEL_MIN_NODES) { + writeLevel0Parallel( + graph, + vectorIndex, + level0Nodes, + offsets[0], + countOnLevel0, + maxConn, + requestedWorkers, + availableProcessors); + } else { + writeLevelSerial(graph, vectorIndex, 0, level0Nodes, offsets[0], countOnLevel0, maxConn); + } + + for (int level = 1; level < numLevels; level++) { int[] sortedNodes = NodesIterator.getSortedNodes(graph.getNodesOnLevel(level)); offsets[level] = new int[sortedNodes.length]; - int nodeOffsetId = 0; - - for (int node : sortedNodes) { - // Get node neighbors - NeighborArray neighbors = graph.getNeighbors(level, node); - // Get the size of the neighbor array - int size = neighbors.size(); - // Write size in VInt as the neighbors list is typically small - long offsetStart = vectorIndex.getFilePointer(); - // Get neighbors - int[] nnodes = neighbors.nodes(); - // Sort them - Arrays.sort(nnodes, 0, size); - // Now that we have sorted, do delta encoding to minimize the required bits to store the - // information - int actualSize = 0; - if (size > 0) { - scratch[0] = nnodes[0]; - actualSize = 1; - } - // De-duplication - for (int i = 1; i < size; i++) { - assert nnodes[i] < countOnLevel0 : "node too large: " + nnodes[i] + ">=" + countOnLevel0; - // Sorting step helps here - if (nnodes[i - 1] == nnodes[i]) { - continue; - } - scratch[actualSize++] = nnodes[i] - nnodes[i - 1]; + writeLevelSerial( + graph, vectorIndex, level, sortedNodes, offsets[level], countOnLevel0, maxConn); + } + return offsets; + } + + /** Node count below which parallel level-zero serialization costs more than it saves. */ + static final int PARALLEL_MIN_NODES = 1 << 16; + + /** Maximum encoded payload retained by one parallel wave before it is written to the index. */ + static final long MAX_PARALLEL_ENCODE_BYTES = 64L << 20; + + private static final int MAX_VINT_BYTES = 5; + + private static void writeLevelSerial( + GPUBuiltHnswGraph graph, + IndexOutput output, + int level, + int[] sortedNodes, + int[] offsets, + int countOnLevel0, + int maxConn) + throws IOException { + int[] scratch = new int[maxConn * 2]; + for (int index = 0; index < sortedNodes.length; index++) { + long start = output.getFilePointer(); + encodeNode( + requireNeighbors(graph, level, sortedNodes[index]), scratch, output, countOnLevel0); + offsets[index] = Math.toIntExact(output.getFilePointer() - start); + } + } + + private static void writeLevel0Parallel( + GPUBuiltHnswGraph graph, + IndexOutput output, + int[] nodes, + int[] offsets, + int countOnLevel0, + int maxConn, + int requestedWorkers, + int availableProcessors) + throws IOException { + try (BoundedParallelExecutor executor = + BoundedParallelExecutor.create(requestedWorkers, nodes.length, availableProcessors)) { + if (!executor.isParallel()) { + writeLevelSerial(graph, output, 0, nodes, offsets, countOnLevel0, maxConn); + return; + } + + int waveSize = nodesPerSerializationWave(maxConn); + for (int waveStart = 0; waveStart < nodes.length; ) { + int waveEnd = (int) Math.min(nodes.length, (long) waveStart + waveSize); + int nodesInWave = waveEnd - waveStart; + int currentWaveStart = waveStart; + ByteBuffersDataOutput[] buffers = + new ByteBuffersDataOutput[executor.workerCountFor(nodesInWave)]; + + executor.invokeRanges( + nodesInWave, + (taskIndex, start, end) -> { + ByteBuffersDataOutput buffer = new ByteBuffersDataOutput(); + int[] scratch = new int[maxConn * 2]; + for (int relativeIndex = start; relativeIndex < end; relativeIndex++) { + if ((relativeIndex & 0x3ff) == 0 && Thread.currentThread().isInterrupted()) { + throw new InterruptedException("graph serialization interrupted"); + } + int nodeIndex = currentWaveStart + relativeIndex; + long before = buffer.size(); + encodeNode( + requireNeighbors(graph, 0, nodes[nodeIndex]), scratch, buffer, countOnLevel0); + offsets[nodeIndex] = Math.toIntExact(buffer.size() - before); + } + buffers[taskIndex] = buffer; + }); + + for (ByteBuffersDataOutput buffer : buffers) { + buffer.copyTo(output); } - // Write the size after duplicates are removed - vectorIndex.writeVInt(actualSize); - // Write de-duplicated neighbors - for (int i = 0; i < actualSize; i++) { - vectorIndex.writeVInt(scratch[i]); + waveStart = waveEnd; + } + } + } + + static int nodesPerSerializationWave(int maxConn) { + long maxBytesPerNode = (Math.max(0L, maxConn) + 1L) * MAX_VINT_BYTES; + return Math.toIntExact( + Math.max(1L, Math.min(Integer.MAX_VALUE, MAX_PARALLEL_ENCODE_BYTES / maxBytesPerNode))); + } + + private static void encodeNode( + NeighborArray neighbors, int[] scratch, DataOutput output, int countOnLevel0) + throws IOException { + int size = neighbors.size(); + if (size > scratch.length) { + throw new IllegalArgumentException( + "Neighbor count " + size + " exceeds scratch capacity " + scratch.length); + } + int actualSize = 0; + if (size > 0) { + int[] nodes = neighbors.nodes(); + Arrays.sort(nodes, 0, size); + scratch[0] = nodes[0]; + actualSize = 1; + for (int index = 1; index < size; index++) { + assert nodes[index] < countOnLevel0 + : "node too large: " + nodes[index] + ">=" + countOnLevel0; + if (nodes[index - 1] != nodes[index]) { + scratch[actualSize++] = nodes[index] - nodes[index - 1]; } - offsets[level][nodeOffsetId++] = - Math.toIntExact(vectorIndex.getFilePointer() - offsetStart); } } - // Return offsets (information written while writing the meta info) - return offsets; + output.writeVInt(actualSize); + for (int index = 0; index < actualSize; index++) { + output.writeVInt(scratch[index]); + } + } + + private static NeighborArray requireNeighbors(GPUBuiltHnswGraph graph, int level, int node) + throws IOException { + NeighborArray neighbors = graph.getNeighbors(level, node); + if (neighbors == null) { + throw new IOException("Missing neighbors for node " + node + " on level " + level); + } + return neighbors; } /** diff --git a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/BoundedParallelExecutor.java b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/BoundedParallelExecutor.java new file mode 100644 index 0000000000..d873df4694 --- /dev/null +++ b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/BoundedParallelExecutor.java @@ -0,0 +1,203 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ +package com.nvidia.cuvs.lucene; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.Callable; +import java.util.concurrent.CancellationException; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; + +/** Runs contiguous work ranges with a bounded number of short-lived platform threads. */ +final class BoundedParallelExecutor implements AutoCloseable { + + @FunctionalInterface + interface RangeTask { + void run(int taskIndex, int startInclusive, int endExclusive) throws Exception; + } + + private final int workers; + private final ExecutorService executor; + private boolean closed; + + static BoundedParallelExecutor create(int requestedWorkers, int workItems) { + return create(requestedWorkers, workItems, Runtime.getRuntime().availableProcessors()); + } + + static BoundedParallelExecutor create( + int requestedWorkers, int workItems, int availableProcessors) { + int workers = effectiveWorkerCount(requestedWorkers, workItems, availableProcessors); + return new BoundedParallelExecutor(workers); + } + + static int effectiveWorkerCount(int requestedWorkers, int workItems, int availableProcessors) { + if (requestedWorkers <= 1 || workItems <= 1 || availableProcessors <= 1) { + return 1; + } + return Math.min(requestedWorkers, Math.min(workItems, availableProcessors)); + } + + private BoundedParallelExecutor(int workers) { + this(workers, workers == 1 ? null : Executors.newFixedThreadPool(workers)); + } + + BoundedParallelExecutor(int workers, ExecutorService executor) { + if (workers < 1 || (workers == 1) != (executor == null)) { + throw new IllegalArgumentException( + "executor must be present exactly when workers exceed one"); + } + this.workers = workers; + this.executor = executor; + } + + int workerCountFor(int workItems) { + return Math.min(workers, Math.max(1, workItems)); + } + + boolean isParallel() { + return executor != null; + } + + void invokeRanges(int workItems, RangeTask task) throws IOException { + if (closed) { + throw new IllegalStateException("executor is closed"); + } + if (workItems < 0) { + throw new IllegalArgumentException("workItems must not be negative"); + } + if (workItems == 0) { + return; + } + + int taskCount = workerCountFor(workItems); + if (taskCount == 1) { + runSerial(task, workItems); + return; + } + + List> tasks = new ArrayList<>(taskCount); + for (int taskIndex = 0; taskIndex < taskCount; taskIndex++) { + int rangeStart = Math.toIntExact((long) taskIndex * workItems / taskCount); + int rangeEnd = Math.toIntExact((long) (taskIndex + 1) * workItems / taskCount); + int rangeIndex = taskIndex; + tasks.add( + () -> { + task.run(rangeIndex, rangeStart, rangeEnd); + return null; + }); + } + + List> futures; + try { + futures = executor.invokeAll(tasks); + } catch (InterruptedException interrupted) { + stopAfterInterruption(); + Thread.currentThread().interrupt(); + throw new IOException("Interrupted while waiting for parallel graph work", interrupted); + } catch (Throwable submissionFailure) { + stopAndAwait(); + rethrow(submissionFailure); + throw new AssertionError("unreachable"); + } + + Throwable failure = null; + for (Future future : futures) { + try { + future.get(); + } catch (InterruptedException interrupted) { + stopAfterInterruption(); + Thread.currentThread().interrupt(); + throw new IOException("Interrupted while collecting parallel graph work", interrupted); + } catch (ExecutionException execution) { + failure = addFailure(failure, execution.getCause()); + } catch (CancellationException cancelled) { + failure = addFailure(failure, cancelled); + } + } + rethrow(failure); + } + + private static void runSerial(RangeTask task, int workItems) throws IOException { + try { + task.run(0, 0, workItems); + } catch (Throwable failure) { + rethrow(failure); + } + } + + private static Throwable addFailure(Throwable primary, Throwable secondary) { + if (primary == null) { + return secondary; + } + if (primary != secondary) { + primary.addSuppressed(secondary); + } + return primary; + } + + private static void rethrow(Throwable failure) throws IOException { + if (failure == null) { + return; + } + if (failure instanceof IOException ioException) { + throw ioException; + } + if (failure instanceof RuntimeException runtimeException) { + throw runtimeException; + } + if (failure instanceof Error error) { + throw error; + } + throw new IOException("Parallel graph work failed", failure); + } + + private void stopAfterInterruption() { + stopAndAwait(); + closed = true; + } + + private void stopAndAwait() { + if (executor != null) { + executor.shutdownNow(); + awaitTerminationPreservingInterrupt(); + } + } + + @Override + public void close() throws IOException { + if (closed) { + return; + } + closed = true; + if (executor == null) { + return; + } + executor.shutdown(); + awaitTerminationPreservingInterrupt(); + } + + private void awaitTerminationPreservingInterrupt() { + boolean interrupted = Thread.interrupted(); + try { + while (!executor.isTerminated()) { + try { + executor.awaitTermination(1, TimeUnit.SECONDS); + } catch (InterruptedException ignored) { + interrupted = true; + executor.shutdownNow(); + } + } + } finally { + if (interrupted) { + Thread.currentThread().interrupt(); + } + } + } +} diff --git a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java index 7e9f888e32..89ef4bb7d0 100644 --- a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java +++ b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java @@ -6,10 +6,14 @@ import static org.apache.lucene.search.DocIdSetIterator.NO_MORE_DOCS; +import com.nvidia.cuvs.CuVSDeviceMatrix; +import com.nvidia.cuvs.CuVSHostMatrix; import com.nvidia.cuvs.CuVSMatrix; import com.nvidia.cuvs.RowView; +import java.io.IOException; import java.util.ArrayList; import java.util.List; +import java.util.function.Supplier; import org.apache.lucene.util.hnsw.HnswGraph; import org.apache.lucene.util.hnsw.NeighborArray; @@ -31,6 +35,12 @@ public class GPUBuiltHnswGraph extends HnswGraph { // Layer 0 is special - it contains all nodes private final NeighborArray[] layer0Neighbors; + private record MaterializedGraph( + int numLevels, + List layerNodes, + NeighborArray[] layer0Neighbors, + List layerNeighbors) {} + /** * Multi-layer constructor that supports arbitrary number of layers. * @@ -41,45 +51,180 @@ public class GPUBuiltHnswGraph extends HnswGraph { */ public GPUBuiltHnswGraph( int size, int dimensions, List layerNodes, List layerAdjacencies) { + this(size, dimensions, materializeSerial(size, layerNodes, layerAdjacencies)); + } + private GPUBuiltHnswGraph(int size, int dimensions, MaterializedGraph graph) { this.size = size; this.dimensions = dimensions; - this.numLevels = layerAdjacencies.size(); - this.layerNodes = new ArrayList<>(); - this.layerNeighbors = new ArrayList<>(); + this.numLevels = graph.numLevels(); + this.layerNodes = graph.layerNodes(); + this.layerNeighbors = graph.layerNeighbors(); + this.layer0Neighbors = graph.layer0Neighbors(); + } + + /** Builds a graph while bounding adjacency materialization to the requested writer threads. */ + static GPUBuiltHnswGraph create( + int size, + int dimensions, + List layerNodes, + List layerAdjacencies, + int requestedWorkers) + throws IOException { + return create( + size, + dimensions, + layerNodes, + layerAdjacencies, + requestedWorkers, + Runtime.getRuntime().availableProcessors()); + } + + static GPUBuiltHnswGraph create( + int size, + int dimensions, + List layerNodes, + List layerAdjacencies, + int requestedWorkers, + int availableProcessors) + throws IOException { + if (requestedWorkers <= 1 || size < PARALLEL_MIN_NODES) { + return new GPUBuiltHnswGraph(size, dimensions, layerNodes, layerAdjacencies); + } + try (BoundedParallelExecutor executor = + BoundedParallelExecutor.create(requestedWorkers, size, availableProcessors)) { + if (!executor.isParallel()) { + return new GPUBuiltHnswGraph(size, dimensions, layerNodes, layerAdjacencies); + } + return new GPUBuiltHnswGraph( + size, dimensions, materializeParallel(size, layerNodes, layerAdjacencies, executor)); + } + } + + /** Node count below which parallel materialization costs more than it saves. */ + static final int PARALLEL_MIN_NODES = 1 << 16; + + /** Maximum temporary host copy used to make a device adjacency safe for concurrent reads. */ + static final long MAX_PARALLEL_GRAPH_COPY_BYTES = 4L << 30; - // Process Layer 0 (base layer with all nodes) - CuVSMatrix layer0Adjacency = layerAdjacencies.get(0); - this.layer0Neighbors = fillNeighborArray(layer0Adjacency, size); + private static MaterializedGraph materializeSerial( + int size, List layerNodes, List layerAdjacencies) { + List upperLayerNodes = new ArrayList<>(); + List upperLayerNeighbors = new ArrayList<>(); + NeighborArray[] baseLayerNeighbors = fillNeighborArraySerial(layerAdjacencies.get(0), size); - // Process higher layers (1 to numLevels-1) - for (int level = 1; level < numLevels; level++) { + for (int level = 1; level < layerAdjacencies.size(); level++) { int[] nodes = layerNodes.get(level); - CuVSMatrix adjacency = layerAdjacencies.get(level); - this.layerNodes.add(nodes); - this.layerNeighbors.add(fillNeighborArray(adjacency, nodes.length)); + upperLayerNodes.add(nodes); + upperLayerNeighbors.add(fillNeighborArraySerial(layerAdjacencies.get(level), nodes.length)); } + return new MaterializedGraph( + layerAdjacencies.size(), upperLayerNodes, baseLayerNeighbors, upperLayerNeighbors); } - /** - * Fills the neighbor array using the adjacency matrix. - * - * @param adjacency instance of adjacency CuVSMatrix - * @param size the number of nodes - * @return the NeighborArray - */ - private NeighborArray[] fillNeighborArray(CuVSMatrix adjacency, int size) { - NeighborArray[] neighbors = new NeighborArray[size]; - for (int i = 0; i < size; i++) { - RowView rv = adjacency.getRow(i); - if (rv != null && rv.size() > 0) { - neighbors[i] = new NeighborArray((int) rv.size(), true); - for (int j = 0; j < rv.size(); j++) { - neighbors[i].addInOrder(rv.getAsInt(j), 1.0f - (j * 0.001f)); + private static MaterializedGraph materializeParallel( + int size, + List layerNodes, + List layerAdjacencies, + BoundedParallelExecutor executor) + throws IOException { + List upperLayerNodes = new ArrayList<>(); + List upperLayerNeighbors = new ArrayList<>(); + NeighborArray[] baseLayerNeighbors = fillNeighborArray(layerAdjacencies.get(0), size, executor); + + for (int level = 1; level < layerAdjacencies.size(); level++) { + int[] nodes = layerNodes.get(level); + upperLayerNodes.add(nodes); + upperLayerNeighbors.add( + fillNeighborArray(layerAdjacencies.get(level), nodes.length, executor)); + } + return new MaterializedGraph( + layerAdjacencies.size(), upperLayerNodes, baseLayerNeighbors, upperLayerNeighbors); + } + + private static NeighborArray[] fillNeighborArray( + CuVSMatrix adjacency, int size, BoundedParallelExecutor executor) throws IOException { + if (size < PARALLEL_MIN_NODES + || !executor.isParallel() + || (adjacency instanceof CuVSDeviceMatrix + && !fitsParallelGraphCopyBudget(adjacency.size(), adjacency.columns()))) { + return fillNeighborArraySerial(adjacency, size); + } + + if (adjacency instanceof CuVSDeviceMatrix deviceAdjacency) { + try (CuVSHostMatrix hostCopy = copyToHost(deviceAdjacency)) { + return fillNeighborArrayParallel(hostCopy, size, executor); + } + } + return fillNeighborArrayParallel(adjacency, size, executor); + } + + /** Returns whether an INT32 graph can be copied without exceeding the host-memory budget. */ + static boolean fitsParallelGraphCopyBudget(long rows, long columns) { + if (rows < 0 || columns < 0) { + return false; + } + if (rows == 0 || columns == 0) { + return true; + } + return rows <= MAX_PARALLEL_GRAPH_COPY_BYTES / Integer.BYTES / columns; + } + + private static CuVSHostMatrix copyToHost(CuVSDeviceMatrix source) { + return copyToHost( + source, + () -> CuVSMatrix.hostBuilder(source.size(), source.columns(), source.dataType()).build()); + } + + static CuVSHostMatrix copyToHost( + CuVSDeviceMatrix source, Supplier hostCopyFactory) { + CuVSHostMatrix hostCopy = hostCopyFactory.get(); + try { + source.toHost(hostCopy); + return hostCopy; + } catch (RuntimeException | Error failure) { + try { + hostCopy.close(); + } catch (RuntimeException | Error closeFailure) { + if (failure != closeFailure) { + failure.addSuppressed(closeFailure); } - } else { - neighbors[i] = new NeighborArray(0, true); } + throw failure; + } + } + + private static NeighborArray[] fillNeighborArraySerial(CuVSMatrix adjacency, int size) { + NeighborArray[] neighbors = new NeighborArray[size]; + for (int i = 0; i < size; i++) { + neighbors[i] = materializeRow(adjacency.getRow(i)); + } + return neighbors; + } + + private static NeighborArray[] fillNeighborArrayParallel( + CuVSMatrix adjacency, int size, BoundedParallelExecutor executor) throws IOException { + NeighborArray[] neighbors = new NeighborArray[size]; + executor.invokeRanges( + size, + (taskIndex, start, end) -> { + for (int row = start; row < end; row++) { + if ((row & 0x3ff) == 0 && Thread.currentThread().isInterrupted()) { + throw new InterruptedException("graph materialization interrupted"); + } + neighbors[row] = materializeRow(adjacency.getRow(row)); + } + }); + return neighbors; + } + + private static NeighborArray materializeRow(RowView row) { + if (row == null || row.size() == 0) { + return new NeighborArray(0, true); + } + NeighborArray neighbors = new NeighborArray(Math.toIntExact(row.size()), true); + for (int index = 0; index < row.size(); index++) { + neighbors.addInOrder(row.getAsInt(index), 1.0f - (index * 0.001f)); } return neighbors; } diff --git a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java index 13edc64975..ab57dfed64 100644 --- a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java +++ b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java @@ -179,9 +179,11 @@ private void writeFieldInternal(FieldInfo fieldInfo, List vectors) thro vectors, acceleratedHNSWParams.getHnswLayers(), params, - QuantizationType.NONE); + QuantizationType.NONE, + acceleratedHNSWParams.getWriterThreads()); long vectorIndexOffset = hnswVectorIndex.getFilePointer(); - int[][] graphLevelNodeOffsets = writeGraph(hnswGraph, hnswVectorIndex); + int[][] graphLevelNodeOffsets = + writeGraph(hnswGraph, hnswVectorIndex, acceleratedHNSWParams.getWriterThreads()); long vectorIndexLength = hnswVectorIndex.getFilePointer() - vectorIndexOffset; writeMeta( hnswVectorIndex, diff --git a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java index 87907d2cbb..c23c89cc6e 100644 --- a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java +++ b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java @@ -185,11 +185,13 @@ private void writeFieldInternal(FieldInfo fieldInfo, List vectors) throw vectors, acceleratedHNSWParams.getHnswLayers(), params, - QuantizationType.BINARY); + QuantizationType.BINARY, + acceleratedHNSWParams.getWriterThreads()); long vectorIndexOffset = hnswVectorIndex.getFilePointer(); // Write the graph to the vector index - int[][] graphLevelNodeOffsets = writeGraph(hnswGraph, hnswVectorIndex); + int[][] graphLevelNodeOffsets = + writeGraph(hnswGraph, hnswVectorIndex, acceleratedHNSWParams.getWriterThreads()); long vectorIndexLength = hnswVectorIndex.getFilePointer() - vectorIndexOffset; // Write metadata diff --git a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java index 7141af56ee..670665dd2f 100644 --- a/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java +++ b/java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java @@ -210,12 +210,14 @@ private void writeFieldInternal(FieldInfo fieldInfo, List vectors) throws IOE unsignedVectors, acceleratedHNSWParams.getHnswLayers(), params, - QuantizationType.SCALAR); + QuantizationType.SCALAR, + acceleratedHNSWParams.getWriterThreads()); long vectorIndexOffset = hnswVectorIndex.getFilePointer(); // Write the graph to the vector index - int[][] graphLevelNodeOffsets = writeGraph(hnswGraph, hnswVectorIndex); + int[][] graphLevelNodeOffsets = + writeGraph(hnswGraph, hnswVectorIndex, acceleratedHNSWParams.getWriterThreads()); long vectorIndexLength = hnswVectorIndex.getFilePointer() - vectorIndexOffset; diff --git a/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/IntGraphTestMatrix.java b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/IntGraphTestMatrix.java new file mode 100644 index 0000000000..57750d039e --- /dev/null +++ b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/IntGraphTestMatrix.java @@ -0,0 +1,134 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ +package com.nvidia.cuvs.lucene; + +import com.nvidia.cuvs.CuVSDeviceMatrix; +import com.nvidia.cuvs.CuVSHostMatrix; +import com.nvidia.cuvs.CuVSMatrix; +import com.nvidia.cuvs.CuVSResources; +import com.nvidia.cuvs.RowView; + +/** Small array-backed matrix used by graph conversion tests. */ +class IntGraphTestMatrix implements CuVSMatrix { + private final int[][] rows; + private final long reportedColumns; + + IntGraphTestMatrix(int[][] rows) { + this(rows, rows.length == 0 ? 0 : rows[0].length); + } + + IntGraphTestMatrix(int[][] rows, long reportedColumns) { + this.rows = rows; + this.reportedColumns = reportedColumns; + } + + @Override + public long size() { + return rows.length; + } + + @Override + public long columns() { + return reportedColumns; + } + + @Override + public DataType dataType() { + return DataType.INT; + } + + @Override + public RowView getRow(long row) { + return new IntRow(rows[Math.toIntExact(row)]); + } + + @Override + public void toArray(int[][] target) { + for (int row = 0; row < rows.length; row++) { + System.arraycopy(rows[row], 0, target[row], 0, rows[row].length); + } + } + + @Override + public void toArray(float[][] target) { + throw new UnsupportedOperationException(); + } + + @Override + public void toArray(byte[][] target) { + throw new UnsupportedOperationException(); + } + + @Override + public void toHost(CuVSHostMatrix target) { + throw new UnsupportedOperationException(); + } + + @Override + public CuVSHostMatrix toHost() { + throw new UnsupportedOperationException(); + } + + @Override + public void toDevice(CuVSDeviceMatrix target, CuVSResources resources) { + throw new UnsupportedOperationException(); + } + + @Override + public CuVSDeviceMatrix toDevice(CuVSResources resources) { + throw new UnsupportedOperationException(); + } + + @Override + public void close() {} + + static class Device extends IntGraphTestMatrix implements CuVSDeviceMatrix { + Device(int[][] rows, long reportedColumns) { + super(rows, reportedColumns); + } + + @Override + public void toHost(CuVSHostMatrix target) { + throw new AssertionError("oversized device graph must not be copied to host"); + } + } + + private record IntRow(int[] values) implements RowView { + @Override + public long size() { + return values.length; + } + + @Override + public int getAsInt(long index) { + return values[Math.toIntExact(index)]; + } + + @Override + public float getAsFloat(long index) { + throw new UnsupportedOperationException(); + } + + @Override + public byte getAsByte(long index) { + throw new UnsupportedOperationException(); + } + + @Override + public void toArray(int[] target) { + System.arraycopy(values, 0, target, 0, values.length); + } + + @Override + public void toArray(float[] target) { + throw new UnsupportedOperationException(); + } + + @Override + public void toArray(byte[] target) { + throw new UnsupportedOperationException(); + } + } +} diff --git a/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestBoundedParallelExecutor.java b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestBoundedParallelExecutor.java new file mode 100644 index 0000000000..8d4c9104f9 --- /dev/null +++ b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestBoundedParallelExecutor.java @@ -0,0 +1,219 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ +package com.nvidia.cuvs.lucene; + +import java.io.IOException; +import java.util.concurrent.AbstractExecutorService; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicIntegerArray; +import java.util.concurrent.atomic.AtomicReference; +import org.apache.lucene.tests.util.LuceneTestCase; +import org.junit.Test; + +public class TestBoundedParallelExecutor extends LuceneTestCase { + + @Test + public void testEffectiveWorkersAreBoundedByCpuAndWork() { + assertEquals(1, BoundedParallelExecutor.effectiveWorkerCount(512, 10, 1)); + assertEquals(1, BoundedParallelExecutor.effectiveWorkerCount(1, 10, 8)); + assertEquals(3, BoundedParallelExecutor.effectiveWorkerCount(512, 3, 8)); + assertEquals(8, BoundedParallelExecutor.effectiveWorkerCount(512, 20, 8)); + } + + @Test + public void testEveryItemIsVisitedExactlyOnce() throws Exception { + int workItems = 10_003; + AtomicIntegerArray visits = new AtomicIntegerArray(workItems); + + try (BoundedParallelExecutor executor = BoundedParallelExecutor.create(512, workItems, 8)) { + executor.invokeRanges( + workItems, + (taskIndex, start, end) -> { + for (int item = start; item < end; item++) { + visits.incrementAndGet(item); + } + }); + } + + for (int item = 0; item < workItems; item++) { + assertEquals("item " + item, 1, visits.get(item)); + } + } + + @Test + public void testFailureWaitsForSiblingBeforeReturning() throws Exception { + CountDownLatch bothStarted = new CountDownLatch(2); + CountDownLatch releaseSibling = new CountDownLatch(1); + AtomicBoolean siblingFinished = new AtomicBoolean(); + + try (BoundedParallelExecutor executor = BoundedParallelExecutor.create(2, 2, 2)) { + IOException thrown = + assertThrows( + IOException.class, + () -> + executor.invokeRanges( + 2, + (taskIndex, start, end) -> { + bothStarted.countDown(); + assertTrue(bothStarted.await(10, TimeUnit.SECONDS)); + if (taskIndex == 0) { + releaseSibling.countDown(); + throw new IOException("expected failure"); + } + assertTrue(releaseSibling.await(10, TimeUnit.SECONDS)); + siblingFinished.set(true); + })); + assertEquals("expected failure", thrown.getMessage()); + assertTrue("sibling task was still running", siblingFinished.get()); + } + } + + @Test + public void testMultipleWorkerFailuresAreSuppressedInTaskOrder() throws Exception { + IOException first = new IOException("first"); + IOException second = new IOException("second"); + + try (BoundedParallelExecutor executor = BoundedParallelExecutor.create(2, 2, 2)) { + IOException thrown = + assertThrows( + IOException.class, + () -> + executor.invokeRanges( + 2, + (taskIndex, start, end) -> { + throw taskIndex == 0 ? first : second; + })); + assertSame(first, thrown); + assertArrayEquals(new Throwable[] {second}, thrown.getSuppressed()); + } + } + + @Test + public void testInterruptionStopsWorkersAndRestoresInterruptStatus() throws Exception { + CountDownLatch workersStarted = new CountDownLatch(2); + CountDownLatch workerExited = new CountDownLatch(2); + CountDownLatch blockWorkers = new CountDownLatch(1); + AtomicReference result = new AtomicReference<>(); + AtomicBoolean interruptRestored = new AtomicBoolean(); + + Thread caller = + new Thread( + () -> { + try (BoundedParallelExecutor executor = BoundedParallelExecutor.create(2, 2, 2)) { + executor.invokeRanges( + 2, + (taskIndex, start, end) -> { + workersStarted.countDown(); + try { + blockWorkers.await(); + } finally { + workerExited.countDown(); + } + }); + result.set(new AssertionError("parallel invocation unexpectedly completed")); + } catch (Throwable failure) { + result.set(failure); + interruptRestored.set(Thread.currentThread().isInterrupted()); + } + }, + "bounded-parallel-interruption-test"); + + caller.start(); + try { + assertTrue(workersStarted.await(10, TimeUnit.SECONDS)); + caller.interrupt(); + caller.join(TimeUnit.SECONDS.toMillis(10)); + assertFalse("caller did not terminate", caller.isAlive()); + assertTrue(result.get() instanceof IOException); + assertTrue("caller interrupt status was not restored", interruptRestored.get()); + assertEquals("worker remained alive after interruption", 0L, workerExited.getCount()); + } finally { + blockWorkers.countDown(); + caller.interrupt(); + caller.join(TimeUnit.SECONDS.toMillis(10)); + } + } + + @Test + public void testPartialSubmissionRejectionStopsAcceptedWorker() throws Exception { + CountDownLatch workerStarted = new CountDownLatch(1); + CountDownLatch workerExited = new CountDownLatch(1); + ExecutorService rejectingExecutor = new RejectAfterFirstExecutor(workerStarted); + + try (BoundedParallelExecutor executor = new BoundedParallelExecutor(2, rejectingExecutor)) { + assertThrows( + RejectedExecutionException.class, + () -> + executor.invokeRanges( + 2, + (taskIndex, start, end) -> { + workerStarted.countDown(); + try { + new CountDownLatch(1).await(); + } finally { + workerExited.countDown(); + } + })); + assertEquals("accepted worker remained alive", 0L, workerExited.getCount()); + } + } + + private static final class RejectAfterFirstExecutor extends AbstractExecutorService { + private final ExecutorService delegate = Executors.newSingleThreadExecutor(); + private final CountDownLatch firstTaskStarted; + private int submissions; + + RejectAfterFirstExecutor(CountDownLatch firstTaskStarted) { + this.firstTaskStarted = firstTaskStarted; + } + + @Override + public void execute(Runnable command) { + if (submissions++ > 0) { + try { + if (!firstTaskStarted.await(10, TimeUnit.SECONDS)) { + throw new AssertionError("accepted task did not start"); + } + } catch (InterruptedException interrupted) { + Thread.currentThread().interrupt(); + throw new RejectedExecutionException( + "interrupted before expected rejection", interrupted); + } + throw new RejectedExecutionException("expected rejection"); + } + delegate.execute(command); + } + + @Override + public void shutdown() { + delegate.shutdown(); + } + + @Override + public java.util.List shutdownNow() { + return delegate.shutdownNow(); + } + + @Override + public boolean isShutdown() { + return delegate.isShutdown(); + } + + @Override + public boolean isTerminated() { + return delegate.isTerminated(); + } + + @Override + public boolean awaitTermination(long timeout, TimeUnit unit) throws InterruptedException { + return delegate.awaitTermination(timeout, unit); + } + } +} diff --git a/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsGraphMaterialization.java b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsGraphMaterialization.java new file mode 100644 index 0000000000..9324280b38 --- /dev/null +++ b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsGraphMaterialization.java @@ -0,0 +1,190 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ +package com.nvidia.cuvs.lucene; + +import com.nvidia.cuvs.CagraIndexParams; +import com.nvidia.cuvs.CuVSDeviceMatrix; +import com.nvidia.cuvs.CuVSHostMatrix; +import com.nvidia.cuvs.CuVSMatrix; +import com.nvidia.cuvs.lucene.AcceleratedHNSWUtils.QuantizationType; +import java.io.IOException; +import java.lang.reflect.Constructor; +import java.lang.reflect.Method; +import java.util.Arrays; +import java.util.List; +import java.util.Random; +import java.util.concurrent.atomic.AtomicInteger; +import org.apache.lucene.index.FieldInfo; +import org.apache.lucene.store.IndexOutput; +import org.apache.lucene.tests.util.LuceneTestCase; +import org.apache.lucene.util.hnsw.HnswGraph.NodesIterator; +import org.apache.lucene.util.hnsw.NeighborArray; +import org.junit.Test; + +public class TestWriterThreadsGraphMaterialization extends LuceneTestCase { + + private static final int NUM_NODES = GPUBuiltHnswGraph.PARALLEL_MIN_NODES + 1_000; + private static final int DEGREE = 12; + + @Test + public void testParallelMaterializationMatchesSerial() throws Exception { + int[][] adjacency = randomAdjacency(NUM_NODES, DEGREE, new Random(1)); + + try (CuVSMatrix matrix = new IntGraphTestMatrix(adjacency)) { + GPUBuiltHnswGraph serial = newSingleLayerGraph(matrix, 1); + GPUBuiltHnswGraph parallel = newSingleLayerGraph(matrix, 4); + assertGraphsEqual(serial, parallel); + } + } + + @Test + public void testGraphCopyBudgetHandlesLargeShapesWithoutOverflow() { + assertTrue(GPUBuiltHnswGraph.fitsParallelGraphCopyBudget(25_000_000L, 32)); + assertFalse(GPUBuiltHnswGraph.fitsParallelGraphCopyBudget(100_000_000L, 32)); + assertFalse(GPUBuiltHnswGraph.fitsParallelGraphCopyBudget(Long.MAX_VALUE, Long.MAX_VALUE)); + assertFalse(GPUBuiltHnswGraph.fitsParallelGraphCopyBudget(-1, 32)); + } + + @Test + public void testOversizedDeviceGraphUsesSerialFallback() throws Exception { + int[][] adjacency = randomAdjacency(NUM_NODES, 1, new Random(2)); + long oversizedColumns = + GPUBuiltHnswGraph.MAX_PARALLEL_GRAPH_COPY_BYTES / Integer.BYTES / NUM_NODES + 1; + + try (CuVSMatrix matrix = new IntGraphTestMatrix.Device(adjacency, oversizedColumns)) { + GPUBuiltHnswGraph graph = newSingleLayerGraph(matrix, 4); + assertEquals(NUM_NODES, graph.size()); + } + } + + @Test + public void testFailedDeviceCopyClosesHostAllocationAndSuppressesCloseFailure() { + RuntimeException copyFailure = new RuntimeException("copy failed"); + RuntimeException closeFailure = new RuntimeException("close failed"); + AtomicInteger hostCloseCount = new AtomicInteger(); + CuVSDeviceMatrix source = + new IntGraphTestMatrix.Device(new int[][] {{0}}, 1) { + @Override + public void toHost(CuVSHostMatrix target) { + throw copyFailure; + } + }; + CuVSHostMatrix hostCopy = new TrackingHostMatrix(hostCloseCount, closeFailure); + + RuntimeException thrown = + assertThrows( + RuntimeException.class, () -> GPUBuiltHnswGraph.copyToHost(source, () -> hostCopy)); + + assertSame(copyFailure, thrown); + assertEquals(1, hostCloseCount.get()); + assertArrayEquals(new Throwable[] {closeFailure}, thrown.getSuppressed()); + } + + @Test + public void testSuccessfulDeviceCopyTransfersHostOwnershipToCaller() { + AtomicInteger hostCloseCount = new AtomicInteger(); + CuVSHostMatrix hostCopy = new TrackingHostMatrix(hostCloseCount, null); + CuVSDeviceMatrix source = + new IntGraphTestMatrix.Device(new int[][] {{0}}, 1) { + @Override + public void toHost(CuVSHostMatrix target) { + assertSame(hostCopy, target); + } + }; + + CuVSHostMatrix returned = GPUBuiltHnswGraph.copyToHost(source, () -> hostCopy); + assertSame(hostCopy, returned); + assertEquals(0, hostCloseCount.get()); + returned.close(); + assertEquals(1, hostCloseCount.get()); + } + + @Test + public void testLegacyPublicDescriptorsRemainAvailable() throws Exception { + Constructor constructor = + GPUBuiltHnswGraph.class.getConstructor(int.class, int.class, List.class, List.class); + assertEquals(0, constructor.getExceptionTypes().length); + + Method createGraph = + AcceleratedHNSWUtils.class.getMethod( + "createMultiLayerHnswGraph", + FieldInfo.class, + int.class, + int.class, + CuVSMatrix.class, + List.class, + int.class, + CagraIndexParams.class, + QuantizationType.class); + assertEquals(GPUBuiltHnswGraph.class, createGraph.getReturnType()); + + Method writeGraph = + AcceleratedHNSWUtils.class.getMethod( + "writeGraph", GPUBuiltHnswGraph.class, IndexOutput.class); + assertArrayEquals(new Class[] {IOException.class}, writeGraph.getExceptionTypes()); + } + + private static GPUBuiltHnswGraph newSingleLayerGraph(CuVSMatrix adjacency, int workers) + throws Exception { + return GPUBuiltHnswGraph.create( + NUM_NODES, + /* dimensions= */ 4, + Arrays.asList((int[]) null), + List.of(adjacency), + workers, + 4); + } + + private static void assertGraphsEqual(GPUBuiltHnswGraph expected, GPUBuiltHnswGraph actual) { + assertEquals(expected.numLevels(), actual.numLevels()); + for (int level = 0; level < expected.numLevels(); level++) { + int[] nodes = NodesIterator.getSortedNodes(expected.getNodesOnLevel(level)); + for (int node : nodes) { + NeighborArray expectedNeighbors = expected.getNeighbors(level, node); + NeighborArray actualNeighbors = actual.getNeighbors(level, node); + assertEquals(expectedNeighbors.size(), actualNeighbors.size()); + assertArrayEquals( + "node " + node + " on level " + level, + Arrays.copyOf(expectedNeighbors.nodes(), expectedNeighbors.size()), + Arrays.copyOf(actualNeighbors.nodes(), actualNeighbors.size())); + } + } + } + + private static int[][] randomAdjacency(int nodes, int degree, Random random) { + int[][] adjacency = new int[nodes][degree]; + for (int[] row : adjacency) { + for (int index = 0; index < row.length; index++) { + row[index] = random.nextInt(nodes); + } + } + return adjacency; + } + + private static final class TrackingHostMatrix extends IntGraphTestMatrix + implements CuVSHostMatrix { + private final AtomicInteger closeCount; + private final RuntimeException closeFailure; + + TrackingHostMatrix(AtomicInteger closeCount, RuntimeException closeFailure) { + super(new int[][] {{0}}); + this.closeCount = closeCount; + this.closeFailure = closeFailure; + } + + @Override + public int get(int row, int column) { + return 0; + } + + @Override + public void close() { + closeCount.incrementAndGet(); + if (closeFailure != null) { + throw closeFailure; + } + } + } +} diff --git a/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsGraphSerialization.java b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsGraphSerialization.java new file mode 100644 index 0000000000..0642fe7768 --- /dev/null +++ b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsGraphSerialization.java @@ -0,0 +1,214 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ +package com.nvidia.cuvs.lucene; + +import com.nvidia.cuvs.CuVSMatrix; +import java.io.IOException; +import java.util.Arrays; +import java.util.List; +import java.util.Random; +import org.apache.lucene.store.ByteBuffersDirectory; +import org.apache.lucene.store.Directory; +import org.apache.lucene.store.IOContext; +import org.apache.lucene.store.IndexInput; +import org.apache.lucene.store.IndexOutput; +import org.apache.lucene.tests.util.LuceneTestCase; +import org.apache.lucene.util.hnsw.NeighborArray; +import org.junit.Test; + +public class TestWriterThreadsGraphSerialization extends LuceneTestCase { + + private static final int NUM_NODES = AcceleratedHNSWUtils.PARALLEL_MIN_NODES + 1_000; + private static final int DEGREE = 12; + private static final int REPORTED_MAX_CONN = 512; + + @Test + public void testParallelSerializationMatchesSerialAcrossWaves() throws Exception { + int[][] adjacency = randomAdjacency(NUM_NODES, DEGREE, new Random(2)); + + try (CuVSMatrix matrix = new IntGraphTestMatrix(adjacency); + Directory directory = new ByteBuffersDirectory()) { + GPUBuiltHnswGraph serialGraph = new ReportedMaxConnGraph(matrix); + GPUBuiltHnswGraph parallelGraph = new ReportedMaxConnGraph(matrix); + + int[][] serialOffsets; + try (IndexOutput output = directory.createOutput("serial", IOContext.DEFAULT)) { + serialOffsets = AcceleratedHNSWUtils.writeGraph(serialGraph, output, 1, 4); + } + int[][] parallelOffsets; + try (IndexOutput output = directory.createOutput("parallel", IOContext.DEFAULT)) { + parallelOffsets = AcceleratedHNSWUtils.writeGraph(parallelGraph, output, 4, 4); + } + + assertEquals(serialOffsets.length, parallelOffsets.length); + for (int level = 0; level < serialOffsets.length; level++) { + assertArrayEquals(serialOffsets[level], parallelOffsets[level]); + } + assertArrayEquals(readAllBytes(directory, "serial"), readAllBytes(directory, "parallel")); + assertTrue( + "test must cross a bounded-wave boundary", + AcceleratedHNSWUtils.nodesPerSerializationWave(REPORTED_MAX_CONN) < NUM_NODES); + } + } + + @Test + public void testNodeEncodingHasStableLiteralBytes() throws Exception { + NeighborArray neighbors = new NeighborArray(3, true); + neighbors.addInOrder(130, 1.0f); + neighbors.addInOrder(0, 0.9f); + neighbors.addInOrder(128, 0.8f); + + SerializedGraph encoded = serialize(new LiteralGraph(131, neighbors)); + assertArrayEquals(new byte[] {3, 0, (byte) 0x80, 1, 2}, encoded.bytes()); + assertArrayEquals(new int[] {5}, encoded.offsets()[0]); + + SerializedGraph empty = serialize(new LiteralGraph(1, new NeighborArray(0, true))); + assertArrayEquals(new byte[] {0}, empty.bytes()); + assertArrayEquals(new int[] {1}, empty.offsets()[0]); + } + + @Test + public void testSerializationWaveHonorsByteBudget() { + for (int maxConn : new int[] {0, 1, 32, 88, 152, 512, Integer.MAX_VALUE}) { + int nodes = AcceleratedHNSWUtils.nodesPerSerializationWave(maxConn); + long maximumBytesPerNode = (Math.max(0L, maxConn) + 1L) * 5L; + assertTrue(nodes > 0); + if (maximumBytesPerNode > AcceleratedHNSWUtils.MAX_PARALLEL_ENCODE_BYTES) { + assertEquals(1, nodes); + } else { + assertTrue( + (long) nodes * maximumBytesPerNode <= AcceleratedHNSWUtils.MAX_PARALLEL_ENCODE_BYTES); + } + } + } + + @Test + public void testSerialSerializationRejectsMissingNeighborsWithoutWritingBytes() throws Exception { + try (Directory directory = new ByteBuffersDirectory(); + IndexOutput output = directory.createOutput("missing", IOContext.DEFAULT)) { + IOException thrown = + assertThrows( + IOException.class, + () -> AcceleratedHNSWUtils.writeGraph(new MissingNeighborsGraph(), output)); + + assertTrue(thrown.getMessage().contains("node 0 on level 0")); + assertEquals(0L, output.getFilePointer()); + } + } + + @Test + public void testParallelSerializationRejectsMissingNeighborsWithoutWritingBytes() + throws Exception { + int[][] adjacency = randomAdjacency(NUM_NODES, DEGREE, new Random(3)); + try (CuVSMatrix matrix = new IntGraphTestMatrix(adjacency); + Directory directory = new ByteBuffersDirectory(); + IndexOutput output = directory.createOutput("missing", IOContext.DEFAULT)) { + GPUBuiltHnswGraph graph = new MissingParallelNeighborsGraph(matrix); + IOException thrown = + assertThrows( + IOException.class, () -> AcceleratedHNSWUtils.writeGraph(graph, output, 4, 4)); + + assertTrue(thrown.getMessage().contains("node " + (NUM_NODES / 2) + " on level 0")); + assertEquals(0L, output.getFilePointer()); + } + } + + private static SerializedGraph serialize(GPUBuiltHnswGraph graph) throws Exception { + try (Directory directory = new ByteBuffersDirectory()) { + int[][] offsets; + try (IndexOutput output = directory.createOutput("encoded", IOContext.DEFAULT)) { + offsets = AcceleratedHNSWUtils.writeGraph(graph, output); + } + return new SerializedGraph(readAllBytes(directory, "encoded"), offsets); + } + } + + private static byte[] readAllBytes(Directory directory, String file) throws Exception { + try (IndexInput input = directory.openInput(file, IOContext.DEFAULT)) { + byte[] bytes = new byte[Math.toIntExact(input.length())]; + input.readBytes(bytes, 0, bytes.length); + return bytes; + } + } + + private static int[][] randomAdjacency(int nodes, int degree, Random random) { + int[][] adjacency = new int[nodes][degree]; + for (int[] row : adjacency) { + for (int index = 0; index < row.length; index++) { + row[index] = random.nextInt(nodes); + } + } + return adjacency; + } + + private static final class ReportedMaxConnGraph extends GPUBuiltHnswGraph { + ReportedMaxConnGraph(CuVSMatrix adjacency) { + super(NUM_NODES, /* dimensions= */ 4, Arrays.asList((int[]) null), List.of(adjacency)); + } + + @Override + public int maxConn() { + return REPORTED_MAX_CONN; + } + } + + private static final class LiteralGraph extends GPUBuiltHnswGraph { + private final int graphSize; + private final NeighborArray neighbors; + + LiteralGraph(int graphSize, NeighborArray neighbors) { + super( + 1, + /* dimensions= */ 1, + Arrays.asList((int[]) null), + List.of(new IntGraphTestMatrix(new int[][] {{}}))); + this.graphSize = graphSize; + this.neighbors = neighbors; + } + + @Override + public int size() { + return graphSize; + } + + @Override + public int maxConn() { + return neighbors.size(); + } + + @Override + public NeighborArray getNeighbors(int level, int node) { + return neighbors; + } + } + + private static final class MissingNeighborsGraph extends GPUBuiltHnswGraph { + MissingNeighborsGraph() { + super( + 1, + /* dimensions= */ 1, + Arrays.asList((int[]) null), + List.of(new IntGraphTestMatrix(new int[][] {{}}))); + } + + @Override + public NeighborArray getNeighbors(int level, int node) { + return null; + } + } + + private static final class MissingParallelNeighborsGraph extends GPUBuiltHnswGraph { + MissingParallelNeighborsGraph(CuVSMatrix adjacency) { + super(NUM_NODES, /* dimensions= */ 4, Arrays.asList((int[]) null), List.of(adjacency)); + } + + @Override + public NeighborArray getNeighbors(int level, int node) { + return node == NUM_NODES / 2 ? null : super.getNeighbors(level, node); + } + } + + private record SerializedGraph(byte[] bytes, int[][] offsets) {} +} diff --git a/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsPersistedIndex.java b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsPersistedIndex.java new file mode 100644 index 0000000000..0002fe032f --- /dev/null +++ b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsPersistedIndex.java @@ -0,0 +1,119 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ +package com.nvidia.cuvs.lucene; + +import static com.nvidia.cuvs.lucene.ThreadLocalCuVSResourcesProvider.isSupported; +import static org.apache.lucene.index.VectorSimilarityFunction.EUCLIDEAN; +import static org.apache.lucene.search.DocIdSetIterator.NO_MORE_DOCS; + +import java.util.Random; +import org.apache.lucene.codecs.Codec; +import org.apache.lucene.codecs.KnnVectorsReader; +import org.apache.lucene.codecs.hnsw.HnswGraphProvider; +import org.apache.lucene.codecs.perfield.PerFieldKnnVectorsFormat; +import org.apache.lucene.document.Document; +import org.apache.lucene.document.Field; +import org.apache.lucene.document.KnnFloatVectorField; +import org.apache.lucene.document.StringField; +import org.apache.lucene.index.CodecReader; +import org.apache.lucene.index.DirectoryReader; +import org.apache.lucene.index.IndexWriter; +import org.apache.lucene.index.IndexWriterConfig; +import org.apache.lucene.index.LeafReader; +import org.apache.lucene.search.IndexSearcher; +import org.apache.lucene.search.KnnFloatVectorQuery; +import org.apache.lucene.store.Directory; +import org.apache.lucene.tests.util.LuceneTestCase; +import org.apache.lucene.tests.util.LuceneTestCase.SuppressSysoutChecks; +import org.apache.lucene.tests.util.TestUtil; +import org.apache.lucene.util.hnsw.HnswGraph; +import org.junit.Test; + +@SuppressSysoutChecks(bugUrl = "") +public class TestWriterThreadsPersistedIndex extends LuceneTestCase { + + private static final String VECTOR_FIELD = "vector"; + private static final int VECTOR_COUNT = AcceleratedHNSWUtils.PARALLEL_MIN_NODES + 1; + private static final int DIMENSIONS = 32; + + @Test + public void testParallelGraphRoundTripAboveThreshold() throws Exception { + assumeTrue("cuVS not supported", isSupported()); + AcceleratedHNSWParams params = + new AcceleratedHNSWParams.Builder() + .withWriterThreads(4) + .withStrategy(AcceleratedHNSWParams.Strategy.CUSTOM) + .withIntermediateGraphDegree(32) + .withGraphDegree(16) + .withHNSWLayer(1) + .build(); + Codec codec = new Lucene101AcceleratedHNSWCodec(params); + float[] query = null; + Random random = new Random(0x2594L); + + try (Directory directory = newDirectory()) { + IndexWriterConfig config = + new IndexWriterConfig() + .setCodec(codec) + .setUseCompoundFile(false) + .setMaxBufferedDocs(VECTOR_COUNT + 1) + .setRAMBufferSizeMB(IndexWriterConfig.DISABLE_AUTO_FLUSH); + try (IndexWriter writer = new IndexWriter(directory, config)) { + for (int id = 0; id < VECTOR_COUNT; id++) { + float[] vector = new float[DIMENSIONS]; + for (int dimension = 0; dimension < DIMENSIONS; dimension++) { + vector[dimension] = random.nextFloat(); + } + if (id == 0) { + query = vector.clone(); + } + Document document = new Document(); + document.add(new StringField("id", Integer.toString(id), Field.Store.YES)); + document.add(new KnnFloatVectorField(VECTOR_FIELD, vector, EUCLIDEAN)); + writer.addDocument(document); + } + } + + TestUtil.checkIndex(directory); + try (DirectoryReader reader = DirectoryReader.open(directory)) { + assertEquals(1, reader.leaves().size()); + assertEquals(VECTOR_COUNT, reader.numDocs()); + HnswGraph graph = graphOf(getOnlyLeafReader(reader)); + assertEquals(VECTOR_COUNT, graph.size()); + int arcs = 0; + HnswGraph.NodesIterator nodes = graph.getNodesOnLevel(0); + while (nodes.hasNext()) { + int node = nodes.nextInt(); + graph.seek(0, node); + for (int neighbor = graph.nextNeighbor(); + neighbor != NO_MORE_DOCS; + neighbor = graph.nextNeighbor()) { + assertTrue(neighbor >= 0); + assertTrue(neighbor < VECTOR_COUNT); + arcs++; + } + } + assertTrue("persisted graph contains no arcs", arcs > 0); + + IndexSearcher searcher = new IndexSearcher(reader); + var hits = searcher.search(new KnnFloatVectorQuery(VECTOR_FIELD, query, 10), 10); + assertEquals(10, hits.scoreDocs.length); + boolean foundExactVector = false; + for (var hit : hits.scoreDocs) { + foundExactVector |= "0".equals(searcher.storedFields().document(hit.doc).get("id")); + } + assertTrue("the indexed vector must be returned for its own query", foundExactVector); + } + } + } + + private static HnswGraph graphOf(LeafReader leaf) throws Exception { + KnnVectorsReader reader = ((CodecReader) leaf).getVectorReader(); + if (reader instanceof PerFieldKnnVectorsFormat.FieldsReader fieldsReader) { + reader = fieldsReader.getFieldReader(VECTOR_FIELD); + } + return ((HnswGraphProvider) reader).getGraph(VECTOR_FIELD); + } +} From f92d07a7de95d3877db5441771f1bf060e06ca71 Mon Sep 17 00:00:00 2001 From: nvzm123 Date: Fri, 18 Sep 2026 13:01:32 +0000 Subject: [PATCH 2/2] Document bounded graph processing and strengthen persistence test --- ...vidia-cuvs-lucene-acceleratedhnswparams.md | 6 +-- ...nvidia-cuvs-lucene-acceleratedhnswutils.md | 18 ++++---- ...om-nvidia-cuvs-lucene-gpubuilthnswgraph.md | 44 ++++++++++++++----- ...ne-lucene99acceleratedhnswvectorswriter.md | 10 ++--- ...leratedhnswbinaryquantizedvectorswriter.md | 10 ++--- ...leratedhnswscalarquantizedvectorswriter.md | 10 ++--- fern/pages/user_guide/lucene.md | 4 +- .../TestWriterThreadsPersistedIndex.java | 23 ++++++---- 8 files changed, 77 insertions(+), 48 deletions(-) diff --git a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-acceleratedhnswparams.md b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-acceleratedhnswparams.md index 27ffe5cb39..7af4ebd43d 100644 --- a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-acceleratedhnswparams.md +++ b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-acceleratedhnswparams.md @@ -18,11 +18,11 @@ public class AcceleratedHNSWParams public int getWriterThreads() ``` -Get the cuVS writer threads parameter +Get the maximum thread count for cuVS writes and post-build graph processing. **Returns** -cuVS writer threads parameter +maximum writer and graph-processing thread count _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWParams.java:149`_ @@ -218,7 +218,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWP public Builder withWriterThreads(int writerThreads) ``` -Set the number of cuVS writer threads while building the index +Set the maximum number of threads used for cuVS writes and post-build graph processing. Valid range - Minimum: \{@value MIN_WRITER_THREADS\}, Maximum: \{@value MAX_WRITER_THREADS\} Default value - \{@value DEFAULT_WRITER_THREADS\} diff --git a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-acceleratedhnswutils.md b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-acceleratedhnswutils.md index 219149bc6d..cff23383c8 100644 --- a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-acceleratedhnswutils.md +++ b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-acceleratedhnswutils.md @@ -21,7 +21,7 @@ public static GPUBuiltHnswGraph createSingleVectorHnswGraph(int size, int dimens Creates a dummy HNSW graph for a single vector. The graph will have 1 level with 1 node and no neighbors. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:54`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:56`_ ### createMultiLayerHnswGraph @@ -35,7 +35,7 @@ M = ceil(cagraGraphDegree / 2), where cagraGraphDegree is the CAGRA adjacency li Each layer contains 1/M nodes from the previous layer Creates layers until the highest layer has ≤ M nodes -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:81`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:83`_ ### writeGraph @@ -62,7 +62,7 @@ a 2D array of offsets | --- | --- | | `IOException` | I/O Exceptions | -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:237`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:263`_ ### writeMeta @@ -91,7 +91,7 @@ Writes the meta information for the index. | --- | --- | | `IOException` | I/O Exceptions | -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:302`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:443`_ ### printInfoStream @@ -107,7 +107,7 @@ A utility method to print info/debugging messages using InfoStream. | --- | --- | | `msg` | the debugging message to print | -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:384`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:525`_ ### writeEmpty @@ -129,7 +129,7 @@ Writes an empty meta information for the field. | --- | --- | | `IOException` | I/O Exceptions | -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:396`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:537`_ ### quantizeFloatVectorsToBinary @@ -152,7 +152,7 @@ Bits are packed: 8 dimensions per byte. A list of byte binary representation for the input vectors -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:409`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:550`_ ### quantizeFloatVectorsToScalar @@ -172,6 +172,6 @@ Scalar quantization. A list of byte scalar representation for the input vectors -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:451`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:592`_ -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:31`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/AcceleratedHNSWUtils.java:33`_ diff --git a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-gpubuilthnswgraph.md b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-gpubuilthnswgraph.md index f8a3aee9dc..d37556c79e 100644 --- a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-gpubuilthnswgraph.md +++ b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-gpubuilthnswgraph.md @@ -31,7 +31,27 @@ Multi-layer constructor that supports arbitrary number of layers. | `layerNodes` | the nodes on the layer | | `layerAdjacencies` | adjacency list | -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:41`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:51`_ + +### create + +```java +static GPUBuiltHnswGraph create( int size, int dimensions, List layerNodes, List layerAdjacencies, int requestedWorkers) throws IOException +``` + +Builds a graph while bounding adjacency materialization to the requested writer threads. + +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:66`_ + +### fitsParallelGraphCopyBudget + +```java +static boolean fitsParallelGraphCopyBudget(long rows, long columns) +``` + +Returns whether an INT32 graph can be copied without exceeding the host-memory budget. + +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:162`_ ### getNodesOnLevel @@ -41,7 +61,7 @@ public NodesIterator getNodesOnLevel(int level) Get all nodes on a given level as node 0th ordinals. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:89`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:234`_ ### getNeighbors @@ -62,7 +82,7 @@ Get the neighbors for the node and the level it resides. an instance of NeighborArray -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:107`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:252`_ ### seek @@ -72,7 +92,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGrap Move the pointer to exactly the given level's target. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:132`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:277`_ ### nextNeighbor @@ -82,7 +102,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGrap Iterates over the neighbor list. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:142`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:287`_ ### entryNode @@ -92,7 +112,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGrap Returns graph's entry point on the top level. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:173`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:318`_ ### maxConn @@ -102,7 +122,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGrap returns M, the maximum number of connections for a node. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:192`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:337`_ ### neighborCount @@ -112,7 +132,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGrap Returns the neighbor count. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:207`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:352`_ ### size @@ -122,7 +142,7 @@ public int size() Returns the number of nodes in the graph. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:282`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:427`_ ### numLevels @@ -136,7 +156,7 @@ Returns the number of levels in the HNSW graph. the number of levels -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:291`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:436`_ ### dimensions @@ -150,6 +170,6 @@ Gets the vector dimension. the vector dimension -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:300`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:445`_ -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:21`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/GPUBuiltHnswGraph.java:25`_ diff --git a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-lucene99acceleratedhnswvectorswriter.md b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-lucene99acceleratedhnswvectorswriter.md index 2f30caeb07..7d06bcadf6 100644 --- a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-lucene99acceleratedhnswvectorswriter.md +++ b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-lucene99acceleratedhnswvectorswriter.md @@ -57,7 +57,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99Accelera Build the indexes and writes it to the disk. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java:203`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java:205`_ ### mergeOneField @@ -67,7 +67,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99Accelera Write field for merging. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java:291`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java:293`_ ### finish @@ -77,7 +77,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99Accelera Called once at the end before close. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java:300`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java:302`_ ### close @@ -87,7 +87,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99Accelera Closes the resources. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java:320`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java:322`_ ### ramBytesUsed @@ -97,6 +97,6 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99Accelera Returns the memory usage of this object in bytes. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java:330`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java:332`_ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/Lucene99AcceleratedHNSWVectorsWriter.java:52`_ diff --git a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-luceneacceleratedhnswbinaryquantizedvectorswriter.md b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-luceneacceleratedhnswbinaryquantizedvectorswriter.md index 91a4e67303..9f7d2436cb 100644 --- a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-luceneacceleratedhnswbinaryquantizedvectorswriter.md +++ b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-luceneacceleratedhnswbinaryquantizedvectorswriter.md @@ -57,7 +57,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAccelerate Build the indexes and writes it to the disk. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java:215`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java:217`_ ### mergeOneField @@ -67,7 +67,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAccelerate Write field for merging. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java:302`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java:304`_ ### finish @@ -77,7 +77,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAccelerate Called once at the end before close. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java:333`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java:335`_ ### close @@ -87,7 +87,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAccelerate Closes the resources. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java:353`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java:355`_ ### ramBytesUsed @@ -97,6 +97,6 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAccelerate Returns the memory usage of this object in bytes. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java:362`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java:364`_ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWBinaryQuantizedVectorsWriter.java:57`_ diff --git a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-luceneacceleratedhnswscalarquantizedvectorswriter.md b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-luceneacceleratedhnswscalarquantizedvectorswriter.md index 9c51cb26e9..9c444cfb71 100644 --- a/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-luceneacceleratedhnswscalarquantizedvectorswriter.md +++ b/fern/pages/lucene_api/lucene-api-com-nvidia-cuvs-lucene-luceneacceleratedhnswscalarquantizedvectorswriter.md @@ -57,7 +57,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAccelerate Build the indexes and writes it to the disk. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java:241`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java:243`_ ### mergeOneField @@ -67,7 +67,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAccelerate Write field for merging. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java:326`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java:328`_ ### finish @@ -77,7 +77,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAccelerate Called once at the end before close. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java:357`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java:359`_ ### close @@ -87,7 +87,7 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAccelerate Closes the resources. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java:377`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java:379`_ ### ramBytesUsed @@ -97,6 +97,6 @@ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAccelerate Returns the memory usage of this object in bytes. -_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java:386`_ +_Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java:388`_ _Source: `java/cuvs-lucene/src/main/java/com/nvidia/cuvs/lucene/LuceneAcceleratedHNSWScalarQuantizedVectorsWriter.java:56`_ diff --git a/fern/pages/user_guide/lucene.md b/fern/pages/user_guide/lucene.md index 34aa9402d4..34ef6fd019 100644 --- a/fern/pages/user_guide/lucene.md +++ b/fern/pages/user_guide/lucene.md @@ -231,7 +231,9 @@ Both `AcceleratedHNSWParams` and `GPUSearchParams` default to a `HEURISTIC` stra - For the accelerated HNSW codecs, set `maxConn` and `beamWidth`, the HNSW parameters you would tune on the CPU. cuVS derives graph degrees and the build algorithm from them. These two values also configure the CPU fallback writer, so one setting covers both paths. - For the GPU search codec, set `buildQuality`. Higher values spend more build time for a higher-quality graph. This codec also passes `graphDegree` into the heuristic, so leave it at its default unless you intend to cap the graph. -Increasing `writerThreads` raises index build concurrency. The accelerated HNSW codecs default to a single writer thread, while the GPU search codec defaults to 32. +`writerThreads` controls cuVS writer concurrency for both parameter types. For the accelerated HNSW codecs, it also bounds the CPU workers used to materialize and serialize the finished graph; the effective worker count cannot exceed the available processors or work items, and the default of one keeps this post-processing serial. The GPU search codec continues to use the setting only for native writer concurrency and defaults to 32. + +With more than one effective worker, parallel device-graph materialization is eligible for adjacency layers with at least 65,536 nodes. It uses one complete temporary native-host copy of an eligible device adjacency when that copy is at most 4 GiB; device adjacencies above that limit use serial materialization, although level-zero serialization can still run in parallel. Switching either class to the `CUSTOM` strategy exposes the underlying CAGRA parameters directly, including `graphDegree`, `intermediateGraphDegree`, the graph build algorithm, and its parameters. Use `CUSTOM` only when you have measurements that justify specific values; the defaults derived by cuVS are a better starting point. For background on the parameters themselves, see the [CAGRA indexing guide](/user-guide/api-guides/indexing-guide/cagra) and the [tuning guide](/getting-started/introduction/tuning-indexes). diff --git a/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsPersistedIndex.java b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsPersistedIndex.java index 0002fe032f..6a5fcb32f1 100644 --- a/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsPersistedIndex.java +++ b/java/cuvs-lucene/src/test/java/com/nvidia/cuvs/lucene/TestWriterThreadsPersistedIndex.java @@ -19,6 +19,7 @@ import org.apache.lucene.document.StringField; import org.apache.lucene.index.CodecReader; import org.apache.lucene.index.DirectoryReader; +import org.apache.lucene.index.FloatVectorValues; import org.apache.lucene.index.IndexWriter; import org.apache.lucene.index.IndexWriterConfig; import org.apache.lucene.index.LeafReader; @@ -50,7 +51,6 @@ public void testParallelGraphRoundTripAboveThreshold() throws Exception { .withHNSWLayer(1) .build(); Codec codec = new Lucene101AcceleratedHNSWCodec(params); - float[] query = null; Random random = new Random(0x2594L); try (Directory directory = newDirectory()) { @@ -66,9 +66,6 @@ public void testParallelGraphRoundTripAboveThreshold() throws Exception { for (int dimension = 0; dimension < DIMENSIONS; dimension++) { vector[dimension] = random.nextFloat(); } - if (id == 0) { - query = vector.clone(); - } Document document = new Document(); document.add(new StringField("id", Integer.toString(id), Field.Store.YES)); document.add(new KnnFloatVectorField(VECTOR_FIELD, vector, EUCLIDEAN)); @@ -80,7 +77,8 @@ public void testParallelGraphRoundTripAboveThreshold() throws Exception { try (DirectoryReader reader = DirectoryReader.open(directory)) { assertEquals(1, reader.leaves().size()); assertEquals(VECTOR_COUNT, reader.numDocs()); - HnswGraph graph = graphOf(getOnlyLeafReader(reader)); + LeafReader leaf = getOnlyLeafReader(reader); + HnswGraph graph = graphOf(leaf); assertEquals(VECTOR_COUNT, graph.size()); int arcs = 0; HnswGraph.NodesIterator nodes = graph.getNodesOnLevel(0); @@ -97,14 +95,23 @@ public void testParallelGraphRoundTripAboveThreshold() throws Exception { } assertTrue("persisted graph contains no arcs", arcs > 0); + int queryNode = graph.entryNode(); + assertTrue(queryNode >= 0); + assertTrue(queryNode < VECTOR_COUNT); + FloatVectorValues values = leaf.getFloatVectorValues(VECTOR_FIELD); + assertNotNull(values); + float[] query = values.vectorValue(queryNode).clone(); + int queryDoc = values.ordToDoc(queryNode); + String queryId = leaf.storedFields().document(queryDoc).get("id"); + IndexSearcher searcher = new IndexSearcher(reader); var hits = searcher.search(new KnnFloatVectorQuery(VECTOR_FIELD, query, 10), 10); assertEquals(10, hits.scoreDocs.length); - boolean foundExactVector = false; + boolean foundQueryNode = false; for (var hit : hits.scoreDocs) { - foundExactVector |= "0".equals(searcher.storedFields().document(hit.doc).get("id")); + foundQueryNode |= queryId.equals(searcher.storedFields().document(hit.doc).get("id")); } - assertTrue("the indexed vector must be returned for its own query", foundExactVector); + assertTrue("the entry-node vector must be returned for its own query", foundQueryNode); } } }