From 912ee706cae34742e09a0101963385b66e8fdf31 Mon Sep 17 00:00:00 2001 From: suyeong Date: Mon, 28 Sep 2026 09:29:52 +0900 Subject: [PATCH 1/3] Sync enVector client and scripts from private main --- README.md | 31 +-- install.py | 7 +- pyproject.toml | 2 +- scripts/e2e.sh | 38 ++++ scripts/get_kmeans_centroids.py | 100 +++++++++ scripts/prepare_dataset.py | 68 ++++-- scripts/prepare_laion_dataset.py | 205 ++++++++++++++++++ scripts/prepare_random_dataset.py | 146 +++++++++++++ scripts/requirements.txt | 4 +- vectordb_bench/backend/clients/api.py | 8 +- .../backend/clients/elastic_cloud/config.py | 4 +- .../backend/clients/envector/cli.py | 51 ++++- .../backend/clients/envector/config.py | 27 ++- .../backend/clients/envector/envector.py | 120 +++++++--- vectordb_bench/backend/dataset.py | 18 +- .../backend/runner/serial_runner.py | 39 ++-- .../config-files/envector_food_config.yml | 46 ++++ .../config-files/envector_laion_config.yml | 36 +++ .../config-files/envector_products_config.yml | 2 +- 19 files changed, 844 insertions(+), 108 deletions(-) create mode 100644 scripts/e2e.sh create mode 100644 scripts/get_kmeans_centroids.py create mode 100644 scripts/prepare_laion_dataset.py create mode 100644 scripts/prepare_random_dataset.py create mode 100644 vectordb_bench/config-files/envector_food_config.yml create mode 100644 vectordb_bench/config-files/envector_laion_config.yml diff --git a/README.md b/README.md index 419095aa4..e1694d81a 100644 --- a/README.md +++ b/README.md @@ -63,8 +63,9 @@ enVector supports two types of benchmark cases: ├── README.md ├── scripts │ ├── get_kmeans_centroids.py # create kmeans centroids +│ ├── prepare_dataset.py # download and prepare ground truth neighbors for GAS dataset +│ ├── prepare_random_dataset.py # download and prepare ground truth neighbors for random dataset │ ├── requirements.txt # python requirements -│ ├── prepare_dataset.py # download and prepare ground truth neighbors for dataset │ └── run_benchmark.sh # benchmark script └── vectordb_bench/config-files # benchmark config file └── envector_{benchmark_case}_config.yml @@ -116,7 +117,7 @@ Run the following commands to run enVector with VectorDBBench's built-in benchma ./scripts/run_benchmark.sh --index-type IVF_FLAT --config-file envector_{benchmark_case}_config.yml # IVF-FLAT ``` -For more details, please refer to `envector_{benchmark_case}_config.yml` in scripts directory for benchmarks with enVector, or you can use the following command: +For more details, please refer to `envector_{benchmark_case}_config.yml` in `vectordb_bench/config-files` directory for benchmarks with enVector, or you can use the following command: ```bash python -m vectordb_bench.cli.vectordbbench envectorflat \ @@ -147,12 +148,11 @@ Prepare the following artifacts for the ANN benchmark with `scripts/prepare_data - prepare ground-truth neighbors - download centroids for the GAS index for corresponding to the embedding model -For the ANN benchmark, we provide two datasets via HuggingFace: -- `PUBMED768D400K`: [cryptolab-playground/pubmed-arxiv-abstract-embedding-gemma-300m](https://huggingface.co/datasets/cryptolab-playground/pubmed-arxiv-abstract-embedding-gemma-300m) -- `BLOOMBERG768D368K`: [cryptolab-playground/Bloomberg-Financial-News-embedding-gemma-300m](https://huggingface.co/datasets/cryptolab-playground/Bloomberg-Financial-News-embedding-gemma-300m) -- `PRODUCTS512D400K` -- `FASHION512D200K` -- `FOOD512D75K` +For the ANN benchmark, we provide four datasets via HuggingFace: +- `pubmed768d400k`: [cryptolab-playground/pubmed-arxiv-abstract-embedding-gemma-300m](https://huggingface.co/datasets/cryptolab-playground/pubmed-arxiv-abstract-embedding-gemma-300m) +- `bloomberg768d368k`: [cryptolab-playground/Bloomberg-Financial-News-embedding-gemma-300m](https://huggingface.co/datasets/cryptolab-playground/Bloomberg-Financial-News-embedding-gemma-300m) +- `products512d400k`: [cryptolab-playground/amazon-products-clip-vit-b-32](https://huggingface.co/datasets/cryptolab-playground/amazon-products-clip-vit-b-32) +- `food512d101k`: [cryptolab-playground/food101-clip-vit-b-32](https://huggingface.co/datasets/cryptolab-playground/food101-clip-vit-b-32) Also, we provide centroids for the corresponding embedding model used in the ANN benchmark: - GAS Centroids: [cryptolab-playground/gas-centroids](https://huggingface.co/datasets/cryptolab-playground/gas-centroids) @@ -165,8 +165,7 @@ pip install -r ./scripts/requirements.txt # Prepare GAS dataset python ./scripts/prepare_dataset.py \ - -d cryptolab-playground/pubmed-arxiv-abstract-embedding-gemma-300m \ - -e embeddinggemma-300m + -d pubmed768d400k ``` Then, you can find the generated files as follows: @@ -196,12 +195,12 @@ Run the provided shell scripts (`./scripts/run_benchmark.sh`) as the following: For more details, please refer to `run_benchmark.sh` or `envector_{benchmark_case}_config.yml` in scripts directory for benchmarks with enVector with ANN (GAS), or you can use the following command: ```bash -python -m vectordb_bench.cli.vectordbbench envectorivfflat \ +python -m vectordb_bench.cli.vectordbbench envectorivfgas \ --config-file envector_pubmed_config.yml # or -python -m vectordb_bench.cli.vectordbbench envectorivfflat \ +python -m vectordb_bench.cli.vectordbbench envectorivfgas \ --uri "localhost:50050" \ --eval-mode mm \ ... \ @@ -211,9 +210,6 @@ python -m vectordb_bench.cli.vectordbbench envectorivfflat \ --nprobe 6 ``` -Note that, **`NUM_PER_BATCH` should be set to the database size** when using IVF-based ANN index for enVector currently. -We will support adjustable `NUM_PER_BATCH` for ANN soon. - ## 🎯 Advanced Usage ### Prepare Other Datasets @@ -225,6 +221,11 @@ If you want to test on other benchmark datasets regardless ANN benchmark, please python ./scripts/prepare_random_dataset.py \ --dataset-dir ./dataset/random512d1m \ --dataset-size 1_000_000 + +# (Optional) Prepare laion dataset +python ./scripts/prepare_laion_dataset.py \ + --dataset-dir ./dataset/laion512d500k \ + --dataset-size 500_000 ``` ### enVector VectorDBBench CLI Options diff --git a/install.py b/install.py index 51202445a..bf43e964a 100644 --- a/install.py +++ b/install.py @@ -49,7 +49,10 @@ def build_multiprocess(args): print("Building base image...") - subprocess.check_call("docker build \ - --rm -t %s -f %s ." % (docker_tag_base(), dockerfile_path_base()), shell=True) + subprocess.check_call( + "docker build \ + --rm -t %s -f %s ." % (docker_tag_base(), dockerfile_path_base()), + shell=True, + ) print("Building end.") diff --git a/pyproject.toml b/pyproject.toml index 79f70759e..b0aa137fa 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -51,7 +51,7 @@ dynamic = ["version"] [project.optional-dependencies] test = [ - "black", + "black>=26.0.0", "ruff", "pytest", ] diff --git a/scripts/e2e.sh b/scripts/e2e.sh new file mode 100644 index 000000000..47c680cfd --- /dev/null +++ b/scripts/e2e.sh @@ -0,0 +1,38 @@ +#!/bin/bash + +set -e + +# Virtual Environment +python -m venv .venv +source .venv/bin/activate + +# Install Python Dependencies +pip install -e . +pip install pyenvector==1.3.0a1 +pip install -r scripts/requirements.txt + +## Dataset Preparation +python ./scripts/prepare_random_dataset.py \ + --dataset-dir random512d10k \ + --dataset-size 10000 + +## Centroid Generation +python ./scripts/get_kmeans_centroids.py \ + --nlist 256 \ + --file-path random512d10k \ + --out-path centroids + +## Benchmark Run +export NUM_PER_BATCH=4096 + +./scripts/run_benchmark.sh \ + --index-type FLAT \ + --config-file envector_sample_config.yml + +./scripts/run_benchmark.sh \ + --index-type IVF_FLAT \ + --config-file envector_sample_config.yml + +./scripts/run_benchmark.sh \ + --index-type IVF_VCT \ + --config-file envector_sample_config.yml \ No newline at end of file diff --git a/scripts/get_kmeans_centroids.py b/scripts/get_kmeans_centroids.py new file mode 100644 index 000000000..372fa4cb4 --- /dev/null +++ b/scripts/get_kmeans_centroids.py @@ -0,0 +1,100 @@ +""" +Get KMeans centroids for a given dataset. +""" + +import argparse +import os + +import faiss +import numpy as np +import pandas as pd +from numpy.linalg import norm + + +def get_args(): + parser = argparse.ArgumentParser(description="KMeans Centroid Calculation") + parser.add_argument( + "--nlist", + type=int, + default=256, + help="Number of clusters for KMeans", + ) + parser.add_argument( + "--file-path", + type=str, + default="/tmp/vectordb_bench/dataset/openai/openai_medium_500k", + help="Path to the dataset directory", + ) + parser.add_argument( + "--out-path", + type=str, + default="/tmp/vectordb_bench/centroids/kmeans-centroids/openai_medium_500k", + help="Path to the output directory", + ) + parser.add_argument( + "--dim", + type=int, + default=512, + help="Dimension of the embeddings.", + ) + return parser.parse_args() + + +def load_dataset(file_path): + print("Loading dataset from:", file_path) + + # load parquet files + train_vectors = pd.read_parquet(f"{file_path}/train.parquet") + + # sort by id + train_vectors.sort_values(by="id", inplace=True) + train_ids = train_vectors["id"].to_numpy(dtype=np.int64) + train_vectors = np.vstack(train_vectors["emb"].values) + train_vectors /= norm(train_vectors, axis=1, keepdims=True) + print(f"train_vectors shape: {train_vectors.shape}") + + return train_vectors.astype(np.float32) + + +def main(): + args = get_args() + + nlist = args.nlist + seed = 42 + + # check existence + if os.path.exists(os.path.join(args.out_path, f"centroids_{nlist}.npy")): + print(f"Centroids for nlist={nlist} already exist. Skipping computation.") + return + + # prepare dataset + train_vectors = load_dataset(args.file_path) + dim = train_vectors.shape[1] + assert dim == args.dim, f"Expected dimension {args.dim}, but got {dim}" + print("✅ Load dataset complete.") + + # kmeans using faiss + kmeans = faiss.Kmeans(dim, nlist, niter=25, seed=seed, verbose=True, gpu=True) + kmeans.train(train_vectors) + print("✅ KMeans training complete.") + + # allocate + _, labels = kmeans.index.search(train_vectors, 1) + labels = labels.flatten() + centroids = kmeans.centroids + print(f"Labels shape: {labels.shape}") + print(f"Centroids shape: {centroids.shape}") + + # normalize + centroids /= norm(centroids, axis=1, keepdims=True) + print(f"Norm: {norm(centroids, axis=1)}") + + # save centroids + os.makedirs(args.out_path, exist_ok=True) + file_name = os.path.join(args.out_path, f"centroids_{nlist}.npy") + np.save(file_name, centroids) + print(f"✅ Centroids saved to {file_name}") + + +if __name__ == "__main__": + main() diff --git a/scripts/prepare_dataset.py b/scripts/prepare_dataset.py index df0e6d798..d8f8801cc 100644 --- a/scripts/prepare_dataset.py +++ b/scripts/prepare_dataset.py @@ -13,6 +13,23 @@ import wget from datasets import load_dataset +SUPPORTED_CASES = { + "pubmed768d400k": { + "dataset_name": "cryptolab-playground/pubmed-arxiv-abstract-embedding-gemma-300m", + "embedding_model": "embeddinggemma-300m", + }, + "bloomberg768d368k": { + "dataset_name": "cryptolab-playground/Bloomberg-Financial-News-embedding-gemma-300m", + "embedding_model": "embeddinggemma-300m", + }, + "products512d400k": { + "dataset_name": "cryptolab-playground/amazon-products-clip-vit-b-32", + "embedding_model": "clip-vit-b-32", + }, + "food512d101k": {"dataset_name": "cryptolab-playground/food101-clip-vit-b-32", "embedding_model": "clip-vit-b-32"}, +} +SUPPORTED_EMBEDDING_MODELS = ["embeddinggemma-300m", "clip-vit-b-32"] + def get_args(): parser = argparse.ArgumentParser(description="Prepare dataset and ground truth neighbors for benchmarking.") @@ -20,25 +37,15 @@ def get_args(): "-d", "--dataset-name", type=str, - default="cryptolab-playground/pubmed-arxiv-abstract-embedding-gemma-300m", + default="pubmed768d400k", help="Huggingface dataset name to download.", - choices=[ - "cryptolab-playground/pubmed-arxiv-abstract-embedding-gemma-300m", - "cryptolab-playground/Bloomberg-Financial-News-embedding-gemma-300m", - ], + choices=list(SUPPORTED_CASES.keys()), ) parser.add_argument( "--dataset-dir", type=str, - default=os.path.join(os.environ.get("DATASET_LOCAL_DIR", "/tmp/vectordb_bench/dataset"), "pubmed768d400k"), - help="Dataset directory to save the dataset and neighbors. Default: 'pubmed768d400k' in DATASET_LOCAL_DIR.", - ) - parser.add_argument( - "-e", - "--embedding-model", - type=str, - default="embeddinggemma-300m", - help="Embedding model name to download centroids for.", + default=None, + help="Dataset directory to save the dataset and neighbors. Default: in DATASET_LOCAL_DIR.", ) parser.add_argument( "--centroids-dir", @@ -51,8 +58,15 @@ def get_args(): def download_dataset(dataset_name: str, output_dir: str = "./dataset/pubmed768d400k") -> None: """Download dataset from Huggingface and save as Parquet files.""" + # check existence + if os.path.exists(os.path.join(output_dir, "train.parquet")) and os.path.exists( + os.path.join(output_dir, "test.parquet") + ): + print(f"Dataset already exists at {output_dir}. Skipping download.") + return + # load dataset - ds = load_dataset(dataset_name) + ds = load_dataset(SUPPORTED_CASES[dataset_name]["dataset_name"]) train = ds["train"].to_pandas() test = ds["test"].to_pandas() @@ -62,12 +76,18 @@ def download_dataset(dataset_name: str, output_dir: str = "./dataset/pubmed768d4 test_table = pa.Table.from_pandas(test) pq.write_table(test_table, f"{output_dir}/test.parquet") + print(f"Saved train and test parquet data to {output_dir}.") def prepare_neighbors( data_dir: str = "./dataset/pubmed768d400k", ) -> None: """Prepare ground truth neighbors using brute-force flat search and save as Parquet.""" + # check neighbor file existence + if os.path.exists(os.path.join(data_dir, "neighbors.parquet")): + print(f"Neighbors data already exists at {data_dir}. Skipping neighbor preparation.") + return + # load dataset train = pd.read_parquet(f"{data_dir}/train.parquet") test = pd.read_parquet(f"{data_dir}/test.parquet") @@ -89,12 +109,17 @@ def prepare_neighbors( table = pa.Table.from_pandas(df) pq.write_table(table, f"{data_dir}/neighbors.parquet") + print(f"Saved neighbors data to {data_dir}.") def download_centroids(embedding_model: str, dataset_dir: str) -> None: """Download pre-computed centroids and for IVF_GAS index.""" + # check existence + if os.path.exists(os.path.join(dataset_dir, embedding_model, "centroids.npy")): + print(f"Centroids data already exists at {os.path.join(dataset_dir, embedding_model)}. Skipping download.") + return - if embedding_model != "embeddinggemma-300m": + if embedding_model not in SUPPORTED_EMBEDDING_MODELS: raise ValueError(f"Centroids for {embedding_model} currently not available.") # BASE URL: https://huggingface.co/datasets/cryptolab-playground/gas-centroids @@ -103,13 +128,20 @@ def download_centroids(embedding_model: str, dataset_dir: str) -> None: # download os.makedirs(os.path.join(dataset_dir, embedding_model), exist_ok=True) wget.download(f"{dataset_link}/centroids.npy", out=os.path.join(dataset_dir, embedding_model, "centroids.npy")) - print(f"\nDownloaded centroids to {os.path.join(dataset_dir, embedding_model)}") + print(f"\nSaved centroids data to {os.path.join(dataset_dir, embedding_model)}") if __name__ == "__main__": args = get_args() + + base_dataset_dir = ( + os.environ.get("DATASET_LOCAL_DIR", "/tmp/vectordb_bench/dataset") + if args.dataset_dir is None + else args.dataset_dir + ) + args.dataset_dir = os.path.join(base_dataset_dir, args.dataset_name) os.makedirs(args.dataset_dir, exist_ok=True) download_dataset(args.dataset_name, args.dataset_dir) prepare_neighbors(args.dataset_dir) - download_centroids(args.embedding_model, args.centroids_dir) + download_centroids(SUPPORTED_CASES[args.dataset_name]["embedding_model"], args.centroids_dir) diff --git a/scripts/prepare_laion_dataset.py b/scripts/prepare_laion_dataset.py new file mode 100644 index 000000000..acd3101cb --- /dev/null +++ b/scripts/prepare_laion_dataset.py @@ -0,0 +1,205 @@ +import argparse +import os + +import faiss +import numpy as np +import pandas as pd +import pyarrow as pa +import pyarrow.parquet as pq +import wget + + +def get_args(): + parser = argparse.ArgumentParser(description="Prepare LAION dataset for benchmarking.") + parser.add_argument( + "--dataset-dir", + type=str, + default=os.path.join(os.environ.get("DATASET_LOCAL_DIR", "/tmp/vectordb_bench/dataset"), "laion512d500k"), + help="Directory to save the LAION dataset embeddings.", + ) + parser.add_argument( + "--dataset-size", + type=int, + default=500_000, + help="Number of dataset embeddings to use.", + ) + parser.add_argument( + "--query-size", + type=int, + default=1_000, + help="Number of query embeddings to use. 1,000 is recommended in VectorDBBench.", + ) + + return parser.parse_args() + + +def download_dataset( + number: int = 1, + download_dir: str = "./dataset/LAION512D500K", +) -> None: + """Download LAION dataset embeddings from LAION deployment.""" + dataset_base_url = "https://deploy.laion.ai/8f83b608504d46bb81708ec86e912220/embeddings" + + # file urls + os.makedirs(download_dir, exist_ok=True) + img_emb_path = f"{dataset_base_url}/img_emb/img_emb_{number}.npy" + txt_emb_path = f"{dataset_base_url}/text_emb/text_emb_{number}.npy" + metadata_path = f"{dataset_base_url}/metadata/metadata_{number}.parquet" + + # download + if not os.path.exists(img_emb_path): + wget.download(img_emb_path, out=download_dir) + + if not os.path.exists(txt_emb_path): + wget.download(txt_emb_path, out=download_dir) + + if not os.path.exists(metadata_path): + wget.download(metadata_path, out=download_dir) + + +def combine_npy_files( + numbers: list[int] = [1], + modal: str = "img", + dataset_dir: str = "./dataset/LAION512D500K", + dataset_size: int = 500_000, +) -> str: + """Combine multiple .npy embedding files into a single .npy file.""" + if modal not in ["img", "text"]: + raise ValueError("Modal should be either 'img' or 'text'.") + + if len(numbers) == 1: + return f"{modal}_emb_{numbers[0]}.npy" + + all_arrays = [] + for i in numbers: + arr = np.load(os.path.join(dataset_dir, f"{modal}_emb_{i}.npy")) + all_arrays.append(arr) + + combined_arr = np.vstack(all_arrays)[:dataset_size] + np.save(os.path.join(dataset_dir, f"{modal}_emb_combined_{combined_arr.shape[0]}.npy"), combined_arr) + + return f"{modal}_emb_combined_{combined_arr.shape[0]}.npy" + + +def npy_to_parquet( + dataset_file: str = "img_emb_1.npy", + modal: str = "img", + dataset_dir: str = "./dataset/LAION512D500K", + dataset_size: int = 500_000, + query_size: int = 1_000, +) -> None: + """Convert downloaded .npy embeddings to Parquet format.""" + print(f"Preparing {modal} embeddings from {dataset_file}...") + + if modal not in ["img", "text"]: + raise ValueError("Modal should be either 'img' or 'text'.") + + if modal == "img": + out_file = "train.parquet" + elif modal == "text": + out_file = "test.parquet" + + arr = np.load(f"{dataset_dir}/{dataset_file}") + print(f"\tShape: {arr.shape}", end=" -> ") + arr = arr[:dataset_size] + print(f"{arr.shape}") + + norm = np.linalg.norm(arr, axis=1, keepdims=True) + arr = arr / norm + print(f"\tNorm: {np.linalg.norm(arr[0])}") + + if modal == "text": + np.random.seed(42) + test_indices = np.random.randint(0, dataset_size, size=query_size) + arr = arr[test_indices] + + test_out_path = os.path.join(dataset_dir, "test_indices.npy") + np.save(test_out_path, test_indices) + print(f"Saved test indices to {test_out_path}") + + ids = np.arange(len(arr)) + id_array = pa.array(ids, type=pa.int64()) + + list_arrays = [arr[i].tolist() for i in range(len(arr))] + vector_array = pa.array(list_arrays, type=pa.list_(pa.float64())) + + assert len(id_array) == len(vector_array) + + table = pa.Table.from_arrays([id_array, vector_array], names=["id", "emb"]) + + out_path = os.path.join(dataset_dir, out_file) + print(f"Saving parquet to {out_path}") + pq.write_table(table, out_path) + + +def prepare_neighbors( + data_dir: str = "./dataset/LAION512D500K", +) -> None: + """Prepare ground truth neighbors using brute-force flat search and save as Parquet.""" + # load dataset + train = pd.read_parquet(f"{data_dir}/train.parquet") + test = pd.read_parquet(f"{data_dir}/test.parquet") + + train = np.stack(train["emb"].to_list()).astype("float32") + test = np.stack(test["emb"].to_list()).astype("float32") + dim = train.shape[1] + + # flat search + index = faiss.IndexFlatIP(dim) + index.add(train) + + k = len(test) + distances, indices = index.search(test, k) + print(distances.shape, indices.shape) + + # save flat search result as neighbors + df = pd.DataFrame({"id": np.arange(len(indices)), "neighbors_id": indices.tolist()}) + + table = pa.Table.from_pandas(df) + pq.write_table(table, f"{data_dir}/neighbors.parquet") + + +if __name__ == "__main__": + args = get_args() + + num_dataset_idx = (args.dataset_size - 1) // 1_000_000 + 1 + print(f"Downloading {num_dataset_idx} dataset files...") + + # download dataset files + data_idx = [i for i in range(1, num_dataset_idx + 1)] + for n in data_idx: + download_dataset(number=n, download_dir=args.dataset_dir) + print("\nDownload completed.") + + # combine npy files + dataset_file_img = combine_npy_files( + numbers=data_idx, + modal="img", + dataset_dir=args.dataset_dir, + dataset_size=args.dataset_size, + ) + dataset_file_text = combine_npy_files( + numbers=data_idx, + modal="text", + dataset_dir=args.dataset_dir, + dataset_size=args.dataset_size, + ) + + # combine npy files + npy_to_parquet( + dataset_file=dataset_file_img, + modal="img", + dataset_dir=args.dataset_dir, + dataset_size=args.dataset_size, + query_size=args.query_size, + ) + npy_to_parquet( + dataset_file=dataset_file_text, + modal="text", + dataset_dir=args.dataset_dir, + dataset_size=args.query_size, + query_size=args.query_size, + ) + + # prepare neighbors + prepare_neighbors(data_dir=args.dataset_dir) diff --git a/scripts/prepare_random_dataset.py b/scripts/prepare_random_dataset.py new file mode 100644 index 000000000..c1d42b2c5 --- /dev/null +++ b/scripts/prepare_random_dataset.py @@ -0,0 +1,146 @@ +import argparse +import os + +import faiss +import numpy as np +import pandas as pd +import pyarrow as pa +import pyarrow.parquet as pq + + +def get_args(): + parser = argparse.ArgumentParser(description="Prepare random dataset for benchmarking.") + parser.add_argument( + "--dataset-dir", + type=str, + default=os.path.join(os.environ.get("DATASET_LOCAL_DIR", "/tmp/vectordb_bench/dataset"), "random512d1m"), + help="Directory to save the random vectors.", + ) + parser.add_argument( + "--dataset-size", + type=int, + default=1_000_000, + help="Number of dataset embeddings to use.", + ) + parser.add_argument( + "--query-size", + type=int, + default=1_000, + help="Number of query embeddings to use. 1,000 is recommended in VectorDBBench.", + ) + parser.add_argument( + "--dim", + type=int, + default=512, + help="Dimension of the embeddings.", + ) + + return parser.parse_args() + + +def get_random_data(num_data, dim, seed): + rng = np.random.default_rng(seed) + + data = rng.uniform(low=-1.0, high=1.0, size=(num_data, dim)) + + # L2 normalize + norm = np.linalg.norm(data, axis=1, keepdims=True) + norm = np.maximum(norm, 1e-10) + data /= norm + + print(data.shape) + return data.astype(np.float32) + + +def npy_to_parquet( + vector: np.ndarray, + dataset_dir: str = "./dataset/random512d1m", + mode: str = "train", +) -> None: + """Convert downloaded .npy embeddings to Parquet format.""" + # check already exists + if os.path.exists(os.path.join(dataset_dir, f"{mode}.parquet")): + print(f"{mode}.parquet already exists in {dataset_dir}, skipping...") + return + + print("Preparing embeddings from numpy array...") + os.makedirs(dataset_dir, exist_ok=True) + + ids = np.arange(len(vector)) + id_array = pa.array(ids, type=pa.int64()) + + list_arrays = [vector[i].tolist() for i in range(len(vector))] + vector_array = pa.array(list_arrays, type=pa.list_(pa.float64())) + + assert len(id_array) == len(vector_array) + + table = pa.Table.from_arrays([id_array, vector_array], names=["id", "emb"]) + + out_path = os.path.join(dataset_dir, f"{mode}.parquet") + print(f"Saving parquet to {out_path}") + pq.write_table(table, out_path) + + +def prepare_neighbors( + data_dir: str = "./dataset/random512d1m", +) -> None: + """Prepare ground truth neighbors using brute-force flat search and save as Parquet.""" + # check already exists + if os.path.exists(os.path.join(data_dir, "neighbors.parquet")): + print(f"neighbors.parquet already exists in {data_dir}, skipping...") + return + + # load dataset + train = pd.read_parquet(f"{data_dir}/train.parquet") + test = pd.read_parquet(f"{data_dir}/test.parquet") + + train = np.stack(train["emb"].to_list()).astype("float32") + test = np.stack(test["emb"].to_list()).astype("float32") + dim = train.shape[1] + + # flat search + index = faiss.IndexFlatIP(dim) + index.add(train) + + k = len(test) + distances, indices = index.search(test, k) + print(f"Distances: {distances.shape}, Indices: {indices.shape}") + + assert all(indices[:, 0] == np.arange(len(test))) ### first N vectors + + # save flat search result as neighbors + df = pd.DataFrame({"id": np.arange(len(indices)), "neighbors_id": indices.tolist()}) + + table = pa.Table.from_pandas(df) + pq.write_table(table, f"{data_dir}/neighbors.parquet") + print(f"Saving parquet to {data_dir}/neighbors.parquet") + + +if __name__ == "__main__": + args = get_args() + + # generate random data and save as .npy + vectors = get_random_data( + num_data=args.dataset_size, + dim=args.dim, + seed=42, + ) + + # prepare train parquet file from numpy arrays + npy_to_parquet( + vector=vectors, + dataset_dir=args.dataset_dir, + mode="train", + ) + + # prepare test set + test_vectors = vectors[: args.query_size] ### first N vectors + + npy_to_parquet( + vector=test_vectors, + dataset_dir=args.dataset_dir, + mode="test", + ) + + # prepare neighbors + prepare_neighbors(data_dir=args.dataset_dir) diff --git a/scripts/requirements.txt b/scripts/requirements.txt index 00dd6eb66..d70880acd 100644 --- a/scripts/requirements.txt +++ b/scripts/requirements.txt @@ -1,4 +1,6 @@ wget==3.2 pandas==2.3.3 pyarrow==22.0.0 -faiss-cpu==1.13.2 \ No newline at end of file +faiss-cpu==1.13.2 +# sklearn +# cuml \ No newline at end of file diff --git a/vectordb_bench/backend/clients/api.py b/vectordb_bench/backend/clients/api.py index a2232df17..f385ce610 100644 --- a/vectordb_bench/backend/clients/api.py +++ b/vectordb_bench/backend/clients/api.py @@ -1,13 +1,13 @@ from abc import ABC, abstractmethod from contextlib import contextmanager -from enum import Enum +from enum import StrEnum from pydantic import BaseModel, SecretStr, validator from vectordb_bench.backend.filter import Filter, FilterOp -class MetricType(str, Enum): +class MetricType(StrEnum): L2 = "L2" COSINE = "COSINE" IP = "IP" @@ -16,7 +16,7 @@ class MetricType(str, Enum): JACCARD = "JACCARD" -class IndexType(str, Enum): +class IndexType(StrEnum): HNSW = "HNSW" HNSW_SQ = "HNSW_SQ" HNSW_BQ = "HNSW_BQ" @@ -46,7 +46,7 @@ class IndexType(str, Enum): IVFGAS = "IVF_VCT" # enVector custom index type -class SQType(str, Enum): +class SQType(StrEnum): SQ6 = "SQ6" SQ8 = "SQ8" BF16 = "BF16" diff --git a/vectordb_bench/backend/clients/elastic_cloud/config.py b/vectordb_bench/backend/clients/elastic_cloud/config.py index 4d9ec32d4..6a3560126 100644 --- a/vectordb_bench/backend/clients/elastic_cloud/config.py +++ b/vectordb_bench/backend/clients/elastic_cloud/config.py @@ -1,4 +1,4 @@ -from enum import Enum +from enum import strEnum from pydantic import BaseModel, SecretStr @@ -16,7 +16,7 @@ def to_dict(self) -> dict: } -class ESElementType(str, Enum): +class ESElementType(strEnum): float = "float" # 4 byte byte = "byte" # 1 byte, -128 to 127 diff --git a/vectordb_bench/backend/clients/envector/cli.py b/vectordb_bench/backend/clients/envector/cli.py index 3b88ae76c..bb9ee0097 100644 --- a/vectordb_bench/backend/clients/envector/cli.py +++ b/vectordb_bench/backend/clients/envector/cli.py @@ -23,12 +23,46 @@ class EnVectorTypedDict(TypedDict): ] eval_mode: Annotated[ str, - click.option("--eval-mode", help="Evaluation mode", type=click.Choice(["mm", "rmp"]), default="mm"), + click.option( + "--eval-mode", + help="Evaluation mode", + type=click.Choice(["mm", "mms", "mm32", "mms32", "rmp"]), + default="mm32", + ), + ] + preset: Annotated[ + str, + click.option( + "--preset", + help="Parameter preset (must match --eval-mode: mm/mms->ip1, mm32/mms32->ip2 or ip3). " + "Empty => derived from --eval-mode.", + type=str, + default="", + ), ] index_name: Annotated[ str, click.option("--index-name", help="Index name", type=str, default="vdbbench"), ] + key_id: Annotated[ + str, + click.option( + "--key-id", help="enVector key id (KMS-managed or local keys//)", type=str, default="default_key" + ), + ] + kms_address: Annotated[ + str, + click.option( + "--kms-address", + type=str, + default="", + help="KMS gateway host:port; enables KMS-managed keys (empty = local keys)", + ), + ] + kms_secure: Annotated[ + bool, + click.option("--kms-secure", type=bool, default=False, help="Use TLS to the KMS gateway"), + ] class EnVectorFlatIndexTypedDict(CommonTypedDict, EnVectorTypedDict): ... @@ -44,11 +78,14 @@ def EnVectorFlat(**parameters: Unpack[EnVectorFlatIndexTypedDict]): db_config=EnVectorConfig( db_label=parameters["db_label"], uri=SecretStr(parameters["uri"]), + key_id=parameters["key_id"], eval_mode=parameters["eval_mode"], collection_name=parameters["index_name"], + kms_address=parameters["kms_address"], + kms_secure=parameters["kms_secure"], index_params={}, ), - db_case_config=FlatIndexConfig(), + db_case_config=FlatIndexConfig(eval_mode=parameters["eval_mode"], preset=parameters["preset"]), **parameters, ) @@ -84,11 +121,16 @@ def EnVectorIVFFlat(**parameters: Unpack[EnVectorIVFFlatIndexTypedDict]): db_config=EnVectorConfig( db_label=parameters["db_label"], uri=SecretStr(parameters["uri"]), + key_id=parameters["key_id"], eval_mode=parameters["eval_mode"], collection_name=parameters["index_name"], + kms_address=parameters["kms_address"], + kms_secure=parameters["kms_secure"], index_params={"nlist": parameters["nlist"], "nprobe": parameters["nprobe"]}, ), db_case_config=IVFFlatIndexConfig( + eval_mode=parameters["eval_mode"], + preset=parameters["preset"], nlist=parameters["nlist"], nprobe=parameters["nprobe"], train_centroids=parameters["train_centroids"], @@ -111,11 +153,16 @@ def EnVectorIVFGAS(**parameters: Unpack[EnVectorIVFGASIndexTypedDict]): db_config=EnVectorConfig( db_label=parameters["db_label"], uri=SecretStr(parameters["uri"]), + key_id=parameters["key_id"], eval_mode=parameters["eval_mode"], collection_name=parameters["index_name"], + kms_address=parameters["kms_address"], + kms_secure=parameters["kms_secure"], index_params={"nlist": parameters["nlist"], "nprobe": parameters["nprobe"]}, ), db_case_config=IVFGASIndexConfig( + eval_mode=parameters["eval_mode"], + preset=parameters["preset"], nlist=parameters["nlist"], nprobe=parameters["nprobe"], train_centroids=parameters["train_centroids"], diff --git a/vectordb_bench/backend/clients/envector/config.py b/vectordb_bench/backend/clients/envector/config.py index 62e2c4f9e..153a1b7c6 100644 --- a/vectordb_bench/backend/clients/envector/config.py +++ b/vectordb_bench/backend/clients/envector/config.py @@ -1,4 +1,4 @@ -from pydantic import BaseModel, SecretStr +from pydantic import BaseModel, SecretStr, validator from ..api import DBCaseConfig, DBConfig, IndexType, MetricType @@ -10,14 +10,33 @@ class EnVectorConfig(DBConfig): key_path: str = "keys" key_id: str = "default_key" collection_name: str = "vdbbench" + # When set (host:port), keys are managed through the enVector KMS gateway + # instead of local key files. Empty string keeps the local-key behavior. + kms_address: str = "" + kms_secure: bool = False # use TLS when talking to the KMS gateway def to_dict(self) -> dict: return { "uri": self.uri.get_secret_value(), "key_path": self.key_path, "key_id": self.key_id, + "collection_name": self.collection_name, + "kms_address": self.kms_address, + "kms_secure": self.kms_secure, } + @validator("*") + def not_empty_field(cls, v: any, field: any): + if ( + field.name in cls.common_short_configs() + or field.name in cls.common_long_configs() + or field.name == "kms_address" # empty = local keys (no KMS gateway) + ): + return v + if not v and isinstance(v, str | SecretStr): + raise ValueError("Empty string!") + return v + class EnVectorIndexConfig(BaseModel): """Base index config for envector""" @@ -25,7 +44,11 @@ class EnVectorIndexConfig(BaseModel): index: IndexType metric_type: MetricType = MetricType.COSINE # envector supports cosine similarity only use_partition_key: bool = True # for label-filter - eval_mode: str = "mm" # default eval_mode + eval_mode: str = "mm32" # default eval_mode + # Parameter preset. Empty => pyenvector derives the per-eval_mode default + # (mm/mms->ip1, mm32/mms32->ip2). Set explicitly to override, e.g. "ip3" + # for mm32/mms32. + preset: str = "" @property def is_gpu_index(self) -> bool: diff --git a/vectordb_bench/backend/clients/envector/envector.py b/vectordb_bench/backend/clients/envector/envector.py index fafe89347..e2e36c35d 100644 --- a/vectordb_bench/backend/clients/envector/envector.py +++ b/vectordb_bench/backend/clients/envector/envector.py @@ -1,6 +1,8 @@ """Wrapper around the EnVector vector database over VectorDB""" import logging +import os +import time from collections.abc import Iterable from contextlib import contextmanager from pathlib import Path @@ -16,6 +18,15 @@ log = logging.getLogger(__name__) +DROP_WAIT_TIMEOUT = float(os.environ.get("ENVECTOR_DROP_WAIT_TIMEOUT", "120")) +DROP_WAIT_POLL_INTERVAL = float(os.environ.get("ENVECTOR_DROP_WAIT_POLL_INTERVAL", "2")) + +# Time to wait for the post-insert merge to reach MERGED_SAVED during optimize(). +# Default 1 day: large index merges can run for many minutes/hours, and search +# must not start until the merge has fully cut over (see _optimize). +OPTIMIZE_WAIT_TIMEOUT = float(os.environ.get("ENVECTOR_OPTIMIZE_WAIT_TIMEOUT", "86400")) +OPTIMIZE_WAIT_POLL_INTERVAL = float(os.environ.get("ENVECTOR_OPTIMIZE_WAIT_POLL_INTERVAL", "2")) + class EnVector(VectorDB): supported_filter_types: list[FilterOp] = [ @@ -39,7 +50,8 @@ def __init__( self.name = name self.db_config = db_config self.case_config = db_case_config - self.collection_name = collection_name + # Get collection_name from db_config if available, otherwise use the parameter + self.collection_name = db_config.get("collection_name", collection_name) self._primary_field = "pk" self._scalar_id_field = "id" @@ -49,21 +61,17 @@ def __init__( self._scalar_id_index_name = "id_sort_idx" self._scalar_labels_index_name = "labels_idx" self.col: ev.Index | None = None + self._request_ids = [] # Initialize the EnVector client - ev.init( - address=self.db_config.get("uri"), - key_path=self.db_config.get("key_path"), - key_id=self.db_config.get("key_id"), - eval_mode=self.case_config.eval_mode, - preset="ip1" if self.case_config.eval_mode == "mm" else "ip", - ) + ev.init(**self._client_init_kwargs()) # Drop old index if specified if drop_old: log.info(f"{self.name} client drop_old index: {self.collection_name}") if self.collection_name in ev.get_index_list(): ev.drop_index(self.collection_name) + self._wait_until_index_deleted() # Check index type index_param = self.case_config.index_param().get("params", {}) @@ -76,6 +84,33 @@ def __init__( ev.disconnect() + def _client_init_kwargs(self) -> dict[str, Any]: + """Common ev.init() kwargs. When kms_address is configured, the client + routes key setup/decryption through the enVector KMS gateway instead of + local key files; otherwise the original local-key behavior is kept.""" + kwargs: dict[str, Any] = { + "address": self.db_config.get("uri"), + "key_id": self.db_config.get("key_id"), + "eval_mode": self.case_config.eval_mode, + } + # preset is optional: when unset pyenvector derives the per-eval_mode + # default (mm/mms->ip1, mm32/mms32->ip2). Pass it only when the config + # specifies one, so callers can override (e.g. "ip3" for mm32/mms32). + if self.case_config.preset: + kwargs["preset"] = self.case_config.preset + kms_address = self.db_config.get("kms_address") + if kms_address: + # KMS manages the key material via auto_key_setup. Do NOT pass + # key_path, or the client builds the cipher from local key files + # (keys//EncKey.json) instead of the KMS-provided keys. + kwargs["kms_address"] = kms_address + # pyenvector uses a dedicated kms_secure flag (default True). Pass it + # explicitly so a no-TLS KMS gateway isn't dialed over TLS. + kwargs["kms_secure"] = bool(self.db_config.get("kms_secure", False)) + else: + kwargs["key_path"] = self.db_config.get("key_path") + return kwargs + def _ensure_index(self, dim: int, index_kwargs: dict[str, Any]): # Check if the collection already exists if self.collection_name in ev.get_index_list(): @@ -95,15 +130,33 @@ def _create_index(self, dim: int, index_kwargs: dict[str, Any]): if index_type in ["IVF_FLAT", "IVF_VCT"] and train_centroids: self._configure_centroids(index_param, index_kwargs) - ev.create_index( - index_name=self.collection_name, - dim=dim, - key_path=self.db_config.get("key_path"), - key_id=self.db_config.get("key_id"), - index_params=index_param, - eval_mode=self.case_config.eval_mode, - **index_kwargs, - ) + create_kwargs: dict[str, Any] = { + "index_name": self.collection_name, + "dim": dim, + "key_id": self.db_config.get("key_id"), + "index_params": index_param, + "eval_mode": self.case_config.eval_mode, + } + # preset optional: unset => pyenvector derives the per-eval_mode default. + if self.case_config.preset: + create_kwargs["preset"] = self.case_config.preset + # In KMS mode the key material comes from the gateway; passing key_path + # would force local key-file lookup. Keep it only for local-key runs. + if not self.db_config.get("kms_address"): + create_kwargs["key_path"] = self.db_config.get("key_path") + ev.create_index(**create_kwargs, **index_kwargs) + + def _wait_until_index_deleted(self): + deadline = time.monotonic() + DROP_WAIT_TIMEOUT + while time.monotonic() < deadline: + if self.collection_name not in ev.get_index_list(): + log.info(f"{self.name} index {self.collection_name} deletion completed") + return + log.debug(f"{self.name} index {self.collection_name} still deleting; waiting...") + time.sleep(DROP_WAIT_POLL_INTERVAL) + + msg = f"Timed out waiting for index deletion: {self.collection_name} (timeout={DROP_WAIT_TIMEOUT}s)" + raise TimeoutError(msg) def _configure_centroids(self, index_param: dict[str, Any], index_kwargs: dict[str, Any]): # Load centroids @@ -129,13 +182,7 @@ def init(self): >>> self.insert_embeddings() >>> self.search_embedding() """ - ev.init( - address=self.db_config.get("uri"), - key_path=self.db_config.get("key_path"), - key_id=self.db_config.get("key_id"), - eval_mode=self.case_config.eval_mode, - preset="ip1" if self.case_config.eval_mode == "mm" else "ip", - ) + ev.init(**self._client_init_kwargs()) try: self.col = ev.Index(self.collection_name) yield @@ -147,7 +194,20 @@ def create_index(self): pass def _optimize(self): - pass + log.debug("Triggering indexing") + self.col.indexing() + log.info("Waiting for merge completion (target_stage=segmentation -> MERGED_SAVED)") + self.col.wait_for_insert_stage( + request_ids=self._request_ids, + target_stage="segmentation", + timeout_s=OPTIMIZE_WAIT_TIMEOUT, + poll_interval_s=OPTIMIZE_WAIT_POLL_INTERVAL, + ) + # clear request_ids after waiting + self._request_ids = [] + log.info("Load index") + self.col.load() + log.info("enVector Indexing completed") def _post_insert(self): pass @@ -172,13 +232,16 @@ def insert_embeddings( assert self.col is not None assert len(embeddings) == len(metadata) - request_ids = kwargs.pop("request_ids", []) # extract request_ids from kwargs for tracking insert operations + # out-list filled with server-generated request_ids; _optimize waits on them + request_ids = kwargs.pop("request_ids", []) insert_count = 0 try: metadata = list(map(str, metadata)) log.debug(f"Inserting {len(embeddings)} embeddings...") - self.col.insert(embeddings, metadata, request_ids=request_ids, await_completion=False) + self.col.insert( + embeddings, metadata, request_ids=request_ids, await_completion=False, execute_until="flush", load=False + ) insert_count += len(embeddings) log.debug(f"Insert successful, count={insert_count}") except Exception as e: @@ -186,6 +249,9 @@ def insert_embeddings( return insert_count, e return insert_count, None + def set_request_ids(self, request_ids: list[str]): + self._request_ids = request_ids + def prepare_filter(self, filters: Filter): pass diff --git a/vectordb_bench/backend/dataset.py b/vectordb_bench/backend/dataset.py index 552b73417..f68d7986f 100644 --- a/vectordb_bench/backend/dataset.py +++ b/vectordb_bench/backend/dataset.py @@ -120,14 +120,16 @@ def file_count(self) -> int: @property def train_files(self) -> list[str]: - train_file = self.train_file - prefix = f"{train_file}" - train_files = [] - prefix_s = [item.strip() for item in prefix.split(",") if item.strip()] - for i in range(len(prefix_s)): - sub_file = f"{prefix_s[i]}.parquet" - train_files.append(sub_file) - return train_files + # Explicit comma-separated base names take precedence (e.g. "a,b" -> a.parquet,b.parquet). + prefix_s = [item.strip() for item in self.train_file.split(",") if item.strip()] + if len(prefix_s) > 1: + return [f"{name}.parquet" for name in prefix_s] + # A single base name is sharded by file_num: train_00.parquet .. train_{N-1}.parquet + # (2-digit zero-padded). file_num <= 1 keeps the original single-file behavior. + base = prefix_s[0] if prefix_s else "train" + if self.file_num and self.file_num > 1: + return [f"{base}_{i:02d}.parquet" for i in range(self.file_num)] + return [f"{base}.parquet"] class LAION(BaseDataset): diff --git a/vectordb_bench/backend/runner/serial_runner.py b/vectordb_bench/backend/runner/serial_runner.py index 235855c48..9b090f507 100644 --- a/vectordb_bench/backend/runner/serial_runner.py +++ b/vectordb_bench/backend/runner/serial_runner.py @@ -2,7 +2,6 @@ import logging import math import multiprocessing as mp -import os import time import traceback @@ -21,8 +20,6 @@ NUM_PER_BATCH = config.NUM_PER_BATCH LOAD_MAX_TRY_COUNT = config.LOAD_MAX_TRY_COUNT -INSERT_TIMEOUT = max(int(os.environ.get("INSERT_TIMEOUT", "300")), config.LOAD_TIMEOUT_DEFAULT) -INSERT_POLL_INTERVAL = int(os.environ.get("INSERT_POLL_INTERVAL", "60")) log = logging.getLogger(__name__) @@ -54,13 +51,15 @@ def retry_insert(self, db: api.VectorDB, retry_idx: int = 0, **kwargs): msg = f"Insert failed and retried more than {config.MAX_INSERT_RETRY} times" raise RuntimeError(msg) from None - def task(self) -> int: + def task(self) -> tuple[int, list[str]]: count = 0 with self.db.init(): log.info(f"({mp.current_process().name:16}) Start inserting embeddings in batch {config.NUM_PER_BATCH}") start = time.perf_counter() request_ids = [] for data_df in self.dataset: + request_id = [] + kwargs = {} all_metadata = data_df[self.dataset.data.train_id_field].tolist() emb_np = np.stack(data_df[self.dataset.data.train_vector_field]) @@ -77,8 +76,10 @@ def task(self) -> int: labels_data = self.dataset.scalar_labels[self.filters.label_field][all_metadata].to_list() else: labels_data = data_df[self.filters.label_field].tolist() + if isinstance(self.db, EnVector): + kwargs["request_ids"] = request_id insert_count, error = self.db.insert_embeddings( - embeddings=all_embeddings, metadata=all_metadata, labels_data=labels_data, request_ids=request_ids + embeddings=all_embeddings, metadata=all_metadata, labels_data=labels_data, **kwargs ) if error is not None: self.retry_insert( @@ -86,32 +87,20 @@ def task(self) -> int: embeddings=all_embeddings, metadata=all_metadata, labels_data=labels_data, - request_ids=request_ids, + **kwargs, ) - + request_ids.extend(request_id) assert insert_count == len(all_metadata) count += insert_count if count % 100_000 == 0: log.info(f"({mp.current_process().name:16}) Loaded {count} embeddings into VectorDB") - if isinstance(self.db, EnVector): - log.info( - "Waiting for inserted rows to become searchable (Index Operation Status v0)... " - f"(requests={len(request_ids)}, timeout={INSERT_TIMEOUT}s)" - ) - self.db.col.indexer.wait_for_inserts_searchable( - index_name=self.db.collection_name, - request_ids=request_ids, - timeout_s=INSERT_TIMEOUT, - poll_interval_s=INSERT_POLL_INTERVAL, - ) - log.info( f"({mp.current_process().name:16}) Finish loading all dataset into VectorDB, " f"dur={time.perf_counter() - start}" ) - return count + return count, request_ids def endless_insert_data(self, all_embeddings: list, all_metadata: list, left_id: int = 0) -> int: with self.db.init(): @@ -163,7 +152,7 @@ def endless_insert_data(self, all_embeddings: list, all_metadata: list, left_id: return count @utils.time_it - def _insert_all_batches(self) -> int: + def _insert_all_batches(self) -> tuple[int, list[str]]: """Performance case only""" with concurrent.futures.ProcessPoolExecutor( mp_context=mp.get_context("spawn"), @@ -171,7 +160,7 @@ def _insert_all_batches(self) -> int: ) as executor: future = executor.submit(self.task) try: - count = future.result(timeout=self.timeout) + return future.result(timeout=self.timeout) except TimeoutError as e: msg = f"VectorDB load dataset timeout in {self.timeout}" log.warning(msg) @@ -181,8 +170,6 @@ def _insert_all_batches(self) -> int: except Exception as e: log.warning(f"VectorDB load dataset error: {e}") raise e from e - else: - return count def run_endlessness(self) -> int: """run forever util DB raises exception or crash""" @@ -220,7 +207,9 @@ def run_endlessness(self) -> int: raise LoadTimeoutError(self.timeout) def run(self) -> int: - count, _ = self._insert_all_batches() + (count, request_ids), _ = self._insert_all_batches() + if isinstance(self.db, EnVector): + self.db.set_request_ids(request_ids) return count diff --git a/vectordb_bench/config-files/envector_food_config.yml b/vectordb_bench/config-files/envector_food_config.yml new file mode 100644 index 000000000..9d0e20881 --- /dev/null +++ b/vectordb_bench/config-files/envector_food_config.yml @@ -0,0 +1,46 @@ +# Custom Case +_base_dataset: &base_dataset + case_type: PerformanceCustomDataset + custom_case_name: FOOD512D101K + custom_case_description: FOOD512D101K benchmark (512D, 101K vectors) + custom_dataset_name: FOOD512D101K + custom_dataset_dir: "" + custom_dataset_size: 101000 + custom_dataset_dim: 512 + custom_dataset_file_count: 1 + custom_dataset_use_shuffled: false + custom_dataset_with_gt: true + k: 10 + +# envector server settings +_base_envector: &base_envector + uri: localhost:50050 + eval_mode: mm + drop_old: true + load: true + +# FLAT +envectorflat: + <<: [*base_dataset, *base_envector] + index_name: food_flat + db_label: FOOD512D101K-FLAT + +# IVF-FLAT with trained k-means centroids +envectorivfflat: + <<: [*base_dataset, *base_envector] + index_name: food_ivfflat + db_label: FOOD512D101K-IVFFLAT + nlist: 128 + nprobe: 6 + train_centroids: true + centroids_path: food/centroids/centroids_128.npy + +# GAS: enVector-customized ANN +envectorivfgas: + <<: [*base_dataset, *base_envector] + index_name: food_ivfgas + db_label: FOOD512D101K-IVFGAS + nlist: 1024 + nprobe: 6 + train_centroids: true + centroids_path: centroids/clip-vit-b-32/centroids.npy \ No newline at end of file diff --git a/vectordb_bench/config-files/envector_laion_config.yml b/vectordb_bench/config-files/envector_laion_config.yml new file mode 100644 index 000000000..950f2ef76 --- /dev/null +++ b/vectordb_bench/config-files/envector_laion_config.yml @@ -0,0 +1,36 @@ +# Custom Case +_base_dataset: &base_dataset + case_type: PerformanceCustomDataset + custom_case_name: LAION512D500K + custom_case_description: LAION512D500K benchmark (512D, 500K vectors) + custom_dataset_name: LAION512D500K + custom_dataset_dir: "" + custom_dataset_size: 500000 + custom_dataset_dim: 512 + custom_dataset_file_count: 1 + custom_dataset_use_shuffled: false + custom_dataset_with_gt: true + k: 10 + +# envector server settings +_base_envector: &base_envector + uri: localhost:50050 + eval_mode: mm + drop_old: true + load: true + +# FLAT +envectorflat: + <<: [*base_dataset, *base_envector] + index_name: laion_flat + db_label: LAION512D500K-FLAT + +# IVF-FLAT with trained k-means centroids +envectorivfflat: + <<: [*base_dataset, *base_envector] + index_name: laion_ivfflat + db_label: LAION512D500K-IVF + nlist: 256 + nprobe: 16 + train_centroids: true + centroids_path: centroids/laion512d500k/centroids.npy \ No newline at end of file diff --git a/vectordb_bench/config-files/envector_products_config.yml b/vectordb_bench/config-files/envector_products_config.yml index 550404208..80b189de8 100644 --- a/vectordb_bench/config-files/envector_products_config.yml +++ b/vectordb_bench/config-files/envector_products_config.yml @@ -40,7 +40,7 @@ envectorivfgas: <<: [*base_dataset, *base_envector] index_name: products_ivfgas db_label: PRODUCTS512D400K-IVFGAS - nlist: 32768 + nlist: 1024 nprobe: 6 train_centroids: true centroids_path: centroids/clip-vit-b-32/centroids.npy \ No newline at end of file From 10be7e7a8d005b8ce0d9dc050d90a2f3cbc9b2a8 Mon Sep 17 00:00:00 2001 From: suyeong Date: Mon, 28 Sep 2026 09:39:24 +0900 Subject: [PATCH 2/3] Drop broken scripts/e2e.sh --- scripts/e2e.sh | 38 -------------------------------------- 1 file changed, 38 deletions(-) delete mode 100644 scripts/e2e.sh diff --git a/scripts/e2e.sh b/scripts/e2e.sh deleted file mode 100644 index 47c680cfd..000000000 --- a/scripts/e2e.sh +++ /dev/null @@ -1,38 +0,0 @@ -#!/bin/bash - -set -e - -# Virtual Environment -python -m venv .venv -source .venv/bin/activate - -# Install Python Dependencies -pip install -e . -pip install pyenvector==1.3.0a1 -pip install -r scripts/requirements.txt - -## Dataset Preparation -python ./scripts/prepare_random_dataset.py \ - --dataset-dir random512d10k \ - --dataset-size 10000 - -## Centroid Generation -python ./scripts/get_kmeans_centroids.py \ - --nlist 256 \ - --file-path random512d10k \ - --out-path centroids - -## Benchmark Run -export NUM_PER_BATCH=4096 - -./scripts/run_benchmark.sh \ - --index-type FLAT \ - --config-file envector_sample_config.yml - -./scripts/run_benchmark.sh \ - --index-type IVF_FLAT \ - --config-file envector_sample_config.yml - -./scripts/run_benchmark.sh \ - --index-type IVF_VCT \ - --config-file envector_sample_config.yml \ No newline at end of file From 3002792a41958b1f5634f2938ca388121fbfc631 Mon Sep 17 00:00:00 2001 From: suyeong Date: Mon, 28 Sep 2026 09:56:00 +0900 Subject: [PATCH 3/3] Bump pyenvector to 1.6.2 in README --- README.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index e1694d81a..f06886066 100644 --- a/README.md +++ b/README.md @@ -22,7 +22,7 @@ cd envector-deployment/docker-compose ```bash # Install Python Dependencies pip install -e . -pip install pyenvector==1.3.0a1 +pip install pyenvector==1.6.2 ``` 3. Run Benchmark @@ -84,7 +84,7 @@ pip install -e . # 3. Install pyenvector # pip uninstall pyenvector # if installed -pip install pyenvector==1.3.0a1 +pip install pyenvector==1.6.2 ``` ### 2. Prepare enVector Server