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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 42 additions & 5 deletions benchmarker/cmd/ann_benchmark.go
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,18 @@ func intFromUUID(uuidStr string) int {
}

// Writes a single batch of vectors to Weaviate using gRPC
const payloadAlphabet = "abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789 "

// randomPayload returns a pseudo-random string of n bytes, used to simulate
// non-vector data so the objects-bucket value is dominated by properties.
func randomPayload(n int) string {
b := make([]byte, n)
for i := range b {
b[i] = payloadAlphabet[rand.Intn(len(payloadAlphabet))]
}
return string(b)
}

func writeChunk(chunk *Batch, client *weaviategrpc.WeaviateClient, cfg *Config) {
objects := make([]*weaviategrpc.BatchObject, len(chunk.Vectors))

Expand Down Expand Up @@ -138,12 +150,17 @@ func writeChunk(chunk *Batch, client *weaviategrpc.WeaviateClient, cfg *Config)
} else {
objects[i].VectorBytes = encodeVector(vector)
}
if cfg.Filter {
nonRefProperties, err := structpb.NewStruct(map[string]interface{}{
"category": strconv.Itoa(chunk.Filters[i]),
})
if cfg.Filter || cfg.PayloadBytes > 0 {
props := map[string]interface{}{}
if cfg.Filter {
props["category"] = strconv.Itoa(chunk.Filters[i])
}
if cfg.PayloadBytes > 0 {
props["payload"] = randomPayload(cfg.PayloadBytes)
}
nonRefProperties, err := structpb.NewStruct(props)
if err != nil {
log.Fatalf("Error creating filtered struct: %v", err)
log.Fatalf("Error creating properties struct: %v", err)
}
objects[i].Properties = &weaviategrpc.BatchObject_Properties{
NonRefProperties: nonRefProperties,
Expand Down Expand Up @@ -212,8 +229,20 @@ func createSchema(cfg *Config, client *weaviate.Client) {
multiTenancyEnabled = true
}

var classProperties []*models.Property
if cfg.PayloadBytes > 0 {
indexInverted := false
classProperties = append(classProperties, &models.Property{
Name: "payload",
DataType: []string{"text"},
IndexFilterable: &indexInverted,
IndexSearchable: &indexInverted,
})
}

classObj := &models.Class{
Class: cfg.ClassName,
Properties: classProperties,
Description: fmt.Sprintf("Created by the Weaviate Benchmarker at %s", time.Now().String()),
MultiTenancyConfig: &models.MultiTenancyConfig{
Enabled: multiTenancyEnabled,
Expand Down Expand Up @@ -369,6 +398,10 @@ func createSchema(cfg *Config, client *weaviate.Client) {

vectorIndexConfig["filterStrategy"] = cfg.FilterStrategy

if cfg.CacheSize > 0 {
vectorIndexConfig["vectorCacheMaxObjects"] = cfg.CacheSize
}

if cfg.NamedVector != "" {
if cfg.MultiVectorDimensions > 0 {
vectorIndexConfig["multivector"] = map[string]interface{}{
Expand Down Expand Up @@ -1067,6 +1100,8 @@ func initAnnBenchmark() {
"queryDuration", 0, "Instead of querying the test dataset once, query for the specified duration in seconds (default 0)")
annBenchmarkCommand.PersistentFlags().BoolVar(&globalConfig.BQ,
"bq", false, "Set BQ")
annBenchmarkCommand.PersistentFlags().IntVar(&globalConfig.PayloadBytes,
"payloadBytes", 0, "Attach an incompressible text property of approximately this many bytes to every object, to simulate non-vector data (default 0)")
annBenchmarkCommand.PersistentFlags().BoolVar(&globalConfig.Cache,
"cache", false, "Set cache")
annBenchmarkCommand.PersistentFlags().BoolVar(&globalConfig.WaitForBackground,
Expand Down Expand Up @@ -1113,6 +1148,8 @@ func initAnnBenchmark() {
"indexType", "hnsw", "Index type (hnsw, flat or hfresh)")
annBenchmarkCommand.PersistentFlags().IntVar(&globalConfig.MaxConnections,
"maxConnections", 16, "Set Weaviate efConstruction parameter (default 16)")
annBenchmarkCommand.PersistentFlags().IntVar(&globalConfig.CacheSize,
"cacheSize", 0, "Set vectorCacheMaxObjects in vectorIndexConfig (0 means use Weaviate default)")
annBenchmarkCommand.PersistentFlags().IntVar(&globalConfig.Shards,
"shards", 1, "Set number of Weaviate shards")
annBenchmarkCommand.PersistentFlags().IntVarP(&globalConfig.BatchSize,
Expand Down
4 changes: 3 additions & 1 deletion benchmarker/cmd/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ type Config struct {
Shards int
DistanceMetric string
MaxConnections int
CacheSize int
Labels string
LabelMap map[string]string
EfConstruction int
Expand Down Expand Up @@ -82,6 +83,7 @@ type Config struct {
MaxPostingSizeKB int
Replicas int
RngFactor float64
PayloadBytes int
}

func (c *Config) Validate() error {
Expand All @@ -105,7 +107,7 @@ func (c *Config) Validate() error {
}

func (c *Config) performUpdates() bool {
return c.UpdatePercentage > 0 && c.UpdatePercentage < 1 && c.UpdateIterations > 0
return c.UpdatePercentage > 0 && c.UpdatePercentage <= 1 && c.UpdateIterations > 0
}

func (c *Config) validateCommon() error {
Expand Down
210 changes: 210 additions & 0 deletions benchmarker/scripts/python/reduce-dataset.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,210 @@
import h5py
import numpy as np
import faiss
import argparse
from pathlib import Path


def create_faiss_index(dimensions: int, distance: str):
"""
Create a FAISS index based on the distance metric.

Args:
dimensions: Number of dimensions in the vectors
distance: Distance metric ("angular", "dot", or "euclidean")

Returns:
FAISS index instance
"""
if distance == "angular":
# For angular distance, we normalize vectors and use inner product
index = faiss.IndexFlatIP(dimensions)
elif distance == "dot":
# For dot product, use inner product
index = faiss.IndexFlatIP(dimensions)
elif distance == "euclidean":
# For euclidean distance, use L2
index = faiss.IndexFlatL2(dimensions)
else:
raise ValueError(f"Unsupported distance metric: {distance}. Must be 'angular', 'dot', or 'euclidean'")

return index


def normalize_vectors(vectors: np.ndarray) -> np.ndarray:
"""Normalize vectors to unit length for angular distance."""
norms = np.linalg.norm(vectors, axis=1, keepdims=True)
norms[norms == 0] = 1 # Avoid division by zero
return vectors / norms


def reduce_dataset(
input_file: str,
output_file: str,
train_size: int,
distance: str,
limit: int = 100,
random_sample: bool = False,
seed: int = None
) -> None:
"""
Reduce the training dataset and recompute neighbors using FAISS brute force.

Args:
input_file: Path to the input HDF5 file
output_file: Path to the output HDF5 file
train_size: Target size for the training dataset
distance: Distance metric ("angular", "dot", or "euclidean")
limit: Number of neighbors to compute for each test query (default: 100)
random_sample: If True, randomly sample training data; otherwise take first N
seed: Random seed for sampling (optional)
"""
if seed is not None:
np.random.seed(seed)

print(f"Reading input file: {input_file}")
with h5py.File(input_file, "r") as hf:
# Read original data
original_train = hf["train"][:]
original_test = hf["test"][:]

print(f"Original train dimensions: {original_train.shape}")
print(f"Original test dimensions: {original_test.shape}")

# Determine how many training samples to keep
original_train_size = original_train.shape[0]
target_size = min(train_size, original_train_size)

if target_size < original_train_size:
print(f"Reducing training dataset from {original_train_size} to {target_size} samples")

if random_sample:
# Randomly sample indices
indices = np.random.choice(original_train_size, size=target_size, replace=False)
indices = np.sort(indices) # Sort to maintain some order
print(f"Randomly sampled {target_size} training vectors")
else:
# Take first N samples
indices = np.arange(target_size)
print(f"Taking first {target_size} training vectors")

# Sample the training data
reduced_train = original_train[indices]
else:
print(f"Training dataset size ({original_train_size}) is already <= target size ({train_size}), keeping all")
reduced_train = original_train
indices = np.arange(original_train_size)

dimensions = reduced_train.shape[1]
print(f"Train dimensions after reduction: {reduced_train.shape}")
print(f"Building FAISS flat index for {dimensions} dimensions with {distance} distance...")

# Prepare vectors based on distance metric
if distance == "angular":
# Normalize both train and test vectors for angular distance
train_vectors = normalize_vectors(reduced_train.astype(np.float32))
test_vectors = normalize_vectors(original_test.astype(np.float32))
else:
train_vectors = reduced_train.astype(np.float32)
test_vectors = original_test.astype(np.float32)

# Create FAISS index
index = create_faiss_index(dimensions, distance)
index.add(train_vectors)
print(f"Index built with {index.ntotal} vectors")

# Compute neighbors for all test queries
print(f"Computing neighbors for {len(test_vectors)} test queries (limit={limit})...")
D, I = index.search(test_vectors, limit)

# Map indices back to original training set indices if we sampled
if target_size < original_train_size:
# I contains indices into the reduced training set
# Map them back to original indices
neighbors_data = indices[I].astype(np.int64)
else:
neighbors_data = I.astype(np.int64)

print(f"Neighbors dimensions: {neighbors_data.shape}")
print(f"Sample neighbors[0]: {neighbors_data[0][:5]}...")

# Write output file
print(f"Writing output file: {output_file}")
with h5py.File(output_file, "w") as out_hf:
# Write train dataset (use original dtype)
out_hf.create_dataset("train", data=reduced_train)

# Write test dataset (unchanged)
out_hf.create_dataset("test", data=original_test)

# Write neighbors dataset
out_hf.create_dataset("neighbors", data=neighbors_data)

# Copy any other datasets that might exist (like distance metadata, etc.)
# but skip filter-related datasets as requested
skip_datasets = {"train", "test", "neighbors", "train_categories",
"test_categories", "train_properties", "test_properties", "filters"}

for key in hf.keys():
if key not in skip_datasets:
print(f"Copying dataset: {key}")
out_hf.create_dataset(key, data=hf[key][:])

print(f"Successfully created reduced dataset: {output_file}")
print(f" Train size: {reduced_train.shape[0]} (was {original_train_size})")
print(f" Test size: {original_test.shape[0]} (unchanged)")
print(f" Neighbors: {neighbors_data.shape}")


def main():
parser = argparse.ArgumentParser(
description="Reduce training dataset size and recompute neighbors using FAISS brute force"
)
parser.add_argument("input_file", help="Path to the input HDF5 file")
parser.add_argument("output_file", help="Path to the output HDF5 file")
parser.add_argument("--train-size", type=int, required=True,
help="Target size for the training dataset (e.g., 20000)")
parser.add_argument("--distance", required=True, choices=["angular", "dot", "euclidean"],
help="Distance metric for the dataset")
parser.add_argument("--limit", type=int, default=100,
help="Number of neighbors to compute for each test query (default: 100)")
parser.add_argument("--random-sample", action="store_true",
help="Randomly sample training data instead of taking first N")
parser.add_argument("--seed", type=int, default=None,
help="Random seed for sampling (optional)")

args = parser.parse_args()

# Validate input file exists
if not Path(args.input_file).exists():
print(f"Error: Input file {args.input_file} does not exist")
return 1

# Validate train size
if args.train_size <= 0:
print(f"Error: --train-size must be positive, got {args.train_size}")
return 1

try:
reduce_dataset(
args.input_file,
args.output_file,
args.train_size,
args.distance,
args.limit,
args.random_sample,
args.seed
)
return 0
except Exception as e:
print(f"Error: {e}")
import traceback
traceback.print_exc()
return 1


if __name__ == "__main__":
exit(main())



Loading