From 9ab4a60d45e8e4e6cc68347de16f3199fd168f03 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Thu, 26 Feb 2026 17:27:45 +0000 Subject: [PATCH 01/20] Save point in debugging --- cpp/cmake/modules/ConfigureCUDA.cmake | 7 +++++-- cpp/src/cluster/detail/kmeans.cuh | 17 ++++++++++++++++- cpp/src/cluster/detail/kmeans_mg.cuh | 5 ++++- 3 files changed, 25 insertions(+), 4 deletions(-) diff --git a/cpp/cmake/modules/ConfigureCUDA.cmake b/cpp/cmake/modules/ConfigureCUDA.cmake index 98fe9e1f85..395fffdb0b 100644 --- a/cpp/cmake/modules/ConfigureCUDA.cmake +++ b/cpp/cmake/modules/ConfigureCUDA.cmake @@ -59,5 +59,8 @@ if(OpenMP_FOUND) endif() # Debug options -list(APPEND CUVS_DEBUG_CUDA_FLAGS -G -Xcompiler=-rdynamic --maxrregcount=64) -list(APPEND CUVS_DEBUG_CUDA_FLAGS -Xptxas --suppress-stack-size-warning) +if(CMAKE_BUILD_TYPE MATCHES Debug) + message(VERBOSE "cuVS: Building with debugging flags") + list(APPEND CUVS_CUDA_FLAGS -Xcompiler=-rdynamic) + list(APPEND CUVS_CUDA_FLAGS -Xptxas --suppress-stack-size-warning) +endif() diff --git a/cpp/src/cluster/detail/kmeans.cuh b/cpp/src/cluster/detail/kmeans.cuh index 635e8813bd..dee8565f5c 100644 --- a/cpp/src/cluster/detail/kmeans.cuh +++ b/cpp/src/cluster/detail/kmeans.cuh @@ -73,6 +73,16 @@ void initRandom(raft::resources const& handle, handle, X, centroids, n_clusters, params.rng_state.seed); } +template +__global__ void dump_var(DataT* v, uint64_t len) { + printf("{"); + for (int i = 0; i < 10; i++) { + if (i < len) printf("%f, ", double(v[i])); + } + printf("} \n"); +} + + /* * @brief Selects 'n_clusters' samples from the input X using kmeans++ algorithm. @@ -144,6 +154,7 @@ void kmeansPlusPlus(raft::resources const& handle, } raft::random::RngState rng(params.rng_state.seed, params.rng_state.type); + //raft::random::RngState rng(params.rng_state.seed, raft::random::GeneratorType::GenPhilox); std::mt19937 gen(params.rng_state.seed); std::uniform_int_distribution<> dis(0, n_samples - 1); @@ -181,10 +192,14 @@ void kmeansPlusPlus(raft::resources const& handle, // <<< Step-3 >>> : Sample x in X with probability p_x = d^2(x, C) / phi_X (C) // Choose 'n_trials' centroid candidates from X with probability proportional to the squared // distance to the nearest existing cluster - + // printf("size of index = %lu, size of weights = %lu, sizeof types (%ld, %ld)\n", + // uint64_t(indices_view.extent(0)), uint64_t(const_weights_view.extent(0)), sizeof(DataT), sizeof(IndexT)); raft::random::discrete(handle, rng, indices_view, const_weights_view); raft::matrix::gather(handle, const_X_view, const_indices_view, candidates_view); + cudaDeviceSynchronize(); + dump_var<<<1, 1>>>(indices.data_handle(), uint64_t(n_trials)); + cudaDeviceSynchronize(); // Calculate pairwise distance between X and the centroid candidates // Output - pwd [n_trials x n_samples] auto pwd = distBuffer.view(); diff --git a/cpp/src/cluster/detail/kmeans_mg.cuh b/cpp/src/cluster/detail/kmeans_mg.cuh index ce3ca5a1fe..2125a297ad 100644 --- a/cpp/src/cluster/detail/kmeans_mg.cuh +++ b/cpp/src/cluster/detail/kmeans_mg.cuh @@ -281,6 +281,7 @@ void initKMeansPlusPlus(const raft::resources& handle, niter); // <<<< Step-3 >>> : for O( log(psi) ) times do + auto nvtx_range = nvtx3::start_range("step 3"); for (int iter = 0; iter < niter; ++iter) { CUVS_LOG_KMEANS(handle, "@Rank-%d:KMeans|| - Iteration %d: # potential centroids sampled - " @@ -374,6 +375,7 @@ void initKMeansPlusPlus(const raft::resources& handle, raft::make_device_matrix_view(centroidsBuf.data(), tot_centroids, n_features); /// <<<< End of Step-5 >>> } /// <<<< Step-6 >>> + nvtx3::end_range(nvtx_range); CUVS_LOG_KMEANS(handle, "@Rank-%d:KMeans||: # potential centroids sampled - %d\n", @@ -404,16 +406,17 @@ void initKMeansPlusPlus(const raft::resources& handle, // seed they should generate the same potentialCentroids auto const_centroids = raft::make_device_matrix_view( potentialCentroids.data_handle(), potentialCentroids.extent(0), potentialCentroids.extent(1)); + auto nvtx_range = nvtx3::start_range("init_plus_plus"); cuvs::cluster::kmeans::init_plus_plus( handle, params, const_centroids, centroidsRawData, workspace); + nvtx3::end_range(nvtx_range); auto inertia = raft::make_host_scalar(0); auto n_iter = raft::make_host_scalar(0); auto weight_view = raft::make_device_vector_view(weight.data_handle(), weight.extent(0)); cuvs::cluster::kmeans::params params_copy = params; params_copy.rng_state = default_params.rng_state; - cuvs::cluster::kmeans::fit_main(handle, params_copy, const_centroids, From cbccc53a56b2666084c338945776ec3dff730746 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Sat, 7 Mar 2026 20:13:50 +0000 Subject: [PATCH 02/20] Adding the reproducer file provided in the bug --- repro_kmeans_mg.cu | 292 +++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 292 insertions(+) create mode 100644 repro_kmeans_mg.cu diff --git a/repro_kmeans_mg.cu b/repro_kmeans_mg.cu new file mode 100644 index 0000000000..332ae50b92 --- /dev/null +++ b/repro_kmeans_mg.cu @@ -0,0 +1,292 @@ +/* + * Minimal multi-GPU KMeans example using cuVS. + * + * nvcc -std=c++17 --extended-lambda --expt-relaxed-constexpr \ + * -Xcompiler -Wno-deprecated-declarations -Xcompiler -fopenmp \ + * -o repro_kmeans_mg repro_kmeans_mg.cu \ + * -I$CONDA_PREFIX/include \ + * -I$CONDA_PREFIX/include/rapids \ + * -I$CONDA_PREFIX/include/rapids/libcudacxx \ + * -L$CONDA_PREFIX/lib \ + * -lcuvs -lrmm -lnccl -lucp -lucs -lucxx -lgomp \ + * -DRAFT_SYSTEM_LITTLE_ENDIAN=1 -DLIBCUDACXX_ENABLE_EXPERIMENTAL_MEMORY_RESOURCE \ + * -Xlinker -rpath=$CONDA_PREFIX/lib \ + * -gencode arch=compute_XX,code=sm_XX + */ + +#include + +#include +#include +#include +#include +#include +#include + +#include + +#include +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +// ----- helpers -------------------------------------------------------------- +#define NCCLCHECK(cmd) \ + do { \ + ncclResult_t res = cmd; \ + if (res != ncclSuccess) { \ + std::fprintf( \ + stderr, "NCCL error %s:%d '%s'\n", __FILE__, __LINE__, ncclGetErrorString(res)); \ + std::exit(EXIT_FAILURE); \ + } \ + } while (0) + +#define CUDACHECK(cmd) \ + do { \ + cudaError_t e = cmd; \ + if (e != cudaSuccess) { \ + std::fprintf( \ + stderr, "CUDA error %s:%d '%s'\n", __FILE__, __LINE__, cudaGetErrorString(e)); \ + std::exit(EXIT_FAILURE); \ + } \ + } while (0) + +constexpr int64_t N_CLUSTERS = 1000; +constexpr int N_BENCHMARK_ITERS = 1; + + +int main(int argc, char* argv[]) +{ + // Parse command line arguments + int64_t n_samples_total = 10000000; // Default value + int64_t n_features = 256; // Default value + + for (int i = 1; i < argc; i++) { + std::string arg = argv[i]; + if (arg == "--samples" || arg == "-s") { + if (i + 1 < argc) { + n_samples_total = std::stoll(argv[++i]); + } else { + std::cerr << "--samples requires a value" << std::endl; + return 1; + } + } else if (arg == "--features" || arg == "-f") { + if (i + 1 < argc) { + n_features = std::stoll(argv[++i]); + } else { + std::cerr << "--features requires a value" << std::endl; + return 1; + } + } else if (arg == "--help" || arg == "-h") { + std::cout << "Usage: " << argv[0] << " [options]" << std::endl; + std::cout << "Options:" << std::endl; + std::cout << " -s, --samples N Set total number of samples (default: 10000000)" << std::endl; + std::cout << " -f, --features N Set number of features (default: 256)" << std::endl; + std::cout << " -h, --help Show this help message" << std::endl; + return 0; + } + } + + // Detect available GPUs + int n_devices = 0; + CUDACHECK(cudaGetDeviceCount(&n_devices)); + if (n_devices < 1) { + std::cerr << "No CUDA devices found." << std::endl; + return 1; + } + std::cout << "Found " << n_devices << " GPU(s), running multi-GPU KMeans." + << std::endl; + + // Create NCCL communicators (one per device) + std::vector dev_list(n_devices); + for (int i = 0; i < n_devices; ++i) dev_list[i] = i; + + std::vector nccl_comms(n_devices); + NCCLCHECK(ncclCommInitAll(nccl_comms.data(), n_devices, dev_list.data())); + + // ------------------------------------------------------------------ + // Generate shared cluster centers on the host (BEFORE the parallel region) + // so that every rank's data is clustered around the SAME centers. + // Use seed 42 to match `make_blobs(random_state=42, centers=1000)`. + // ------------------------------------------------------------------ + std::vector h_centers(N_CLUSTERS * n_features); + std::mt19937 gen(42); + { + std::uniform_real_distribution dist(-10.0f, 10.0f); + for (auto& v : h_centers) v = dist(gen); + } + + // Pre-generate per-rank seeds (mimics Dask's per-partition random seeds) + std::vector rank_seeds(n_devices); + for (auto& s : rank_seeds) s = gen(); + + // Launch one OpenMP thread per GPU (OPG model) +#pragma omp parallel for num_threads(n_devices) + for (int rank = 0; rank < n_devices; ++rank) { + // 1. Bind this thread to the correct GPU + CUDACHECK(cudaSetDevice(rank)); + + // 2. Create a raft handle and inject the NCCL communicator. + // This is the *only* extra step compared to single-GPU usage. + raft::resources handle; + raft::comms::build_comms_nccl_only(&handle, nccl_comms[rank], n_devices, rank); + + auto stream = raft::resource::get_cuda_stream(handle); + + // 3. Compute this rank's shard of the dataset. + // Each rank gets a contiguous, non-overlapping slice of the rows. + int64_t n_samples_per_rank = n_samples_total / n_devices; + int64_t leftover = n_samples_total % n_devices; + int64_t n_local = n_samples_per_rank + (rank < leftover ? 1 : 0); + + // Copy the shared cluster centers to this GPU + rmm::device_uvector d_blob_centers(N_CLUSTERS * n_features, stream); + CUDACHECK(cudaMemcpyAsync(d_blob_centers.data(), + h_centers.data(), + h_centers.size() * sizeof(float), + cudaMemcpyHostToDevice, + stream)); + + // Generate synthetic clustered data directly on device. + // Each rank uses a different seed so the data points are distinct, + // but all ranks share the same cluster centers. + rmm::device_uvector d_X(n_local * n_features, stream); + rmm::device_uvector d_labels_true(n_local, stream); + + raft::random::make_blobs(d_X.data(), + d_labels_true.data(), + n_local, + n_features, + N_CLUSTERS, + stream, + /*row_major=*/true, + /*centers=*/d_blob_centers.data(), + /*cluster_std=*/nullptr, + /*cluster_std_scalar=*/0.5f, + /*shuffle=*/true, + /*center_box_min=*/-10.0f, + /*center_box_max=*/10.0f, + /*seed=*/rank_seeds[rank]); + + // 4. Prepare output buffers + rmm::device_uvector d_centroids(N_CLUSTERS * n_features, stream); + rmm::device_uvector d_labels(n_local, stream); + + // Ensure data generation is complete before timing + raft::resource::sync_stream(handle); + + if (rank == 0) { + std::cout << "Data generation complete. " + << n_samples_total << " total samples, " + << n_local << " per rank, " + << n_features << " features, " + << N_CLUSTERS << " clusters." << std::endl; + } + + // 5. Benchmark loop — run both RNG types to compare + struct RngConfig { + raft::random::GeneratorType type; + const char* name; + }; + RngConfig rng_configs[] = { + {raft::random::GeneratorType::GenPC, "GenPC"}, + {raft::random::GeneratorType::GenPhilox, "GenPhilox"}, + }; + + for (const auto& rng_cfg : rng_configs) { + // Create NVTX range for this RNG configuration + nvtx3::scoped_range range(rng_cfg.name); + + if (rank == 0) { + std::printf("\n=== RNG: %s ===\n", rng_cfg.name); + } + + int64_t total_kmeans_iters = 0; + for (int iter = 0; iter < N_BENCHMARK_ITERS; ++iter) { + float inertia = 0; + float pred_inertia = 0; + int64_t n_iter = 0; + + cuvs::cluster::kmeans::params params; + params.n_clusters = static_cast(N_CLUSTERS); + params.max_iter = 300; + params.tol = 1e-4; + params.rng_state.seed = 42; + params.rng_state.type = rng_cfg.type; + params.oversampling_factor = 2.0; + params.n_init = 1; + + auto X_view = raft::make_device_matrix_view( + d_X.data(), n_local, n_features); + auto centroids_view = raft::make_device_matrix_view( + d_centroids.data(), N_CLUSTERS, n_features); + auto centroids_const_view = raft::make_device_matrix_view( + d_centroids.data(), N_CLUSTERS, n_features); + auto labels_view = raft::make_device_vector_view( + d_labels.data(), n_local); + + auto t0 = std::chrono::high_resolution_clock::now(); + + { + nvtx3::scoped_range fit_range("kmeans_fit"); + cuvs::cluster::kmeans::fit(handle, + params, + X_view, + std::nullopt, + centroids_view, + raft::make_host_scalar_view(&inertia), + raft::make_host_scalar_view(&n_iter)); + raft::resource::sync_stream(handle); + } + + auto t1 = std::chrono::high_resolution_clock::now(); + + { + nvtx3::scoped_range predict_range("kmeans_predict"); + cuvs::cluster::kmeans::predict(handle, + params, + X_view, + std::nullopt, + centroids_const_view, + labels_view, + true, + raft::make_host_scalar_view(&pred_inertia)); + raft::resource::sync_stream(handle); + } + + auto t2 = std::chrono::high_resolution_clock::now(); + + double fit_s = std::chrono::duration(t1 - t0).count(); + double total_s = std::chrono::duration(t2 - t0).count(); + + total_kmeans_iters += n_iter; + + if (rank == 0) { + std::printf("%d) Time taken by fit %.2fs and predict %.2fs (iters: %ld)\n", + iter, fit_s, total_s, static_cast(n_iter)); + } + } + + if (rank == 0) { + double avg_iters = static_cast(total_kmeans_iters) / N_BENCHMARK_ITERS; + std::printf("Average KMeans iterations per run (%s): %.1f\n", rng_cfg.name, avg_iters); + } + } + } + + // Clean up NCCL + for (int i = 0; i < n_devices; ++i) { + ncclCommDestroy(nccl_comms[i]); + } + + std::cout << "Done." << std::endl; + return 0; +} From c0f8a5778095d0dd20cfee2422681652f5a01d92 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Sat, 7 Mar 2026 20:16:29 +0000 Subject: [PATCH 03/20] Adding reproducer builder --- build_repro.sh | 13 +++++++++++++ 1 file changed, 13 insertions(+) create mode 100644 build_repro.sh diff --git a/build_repro.sh b/build_repro.sh new file mode 100644 index 0000000000..19e827cad7 --- /dev/null +++ b/build_repro.sh @@ -0,0 +1,13 @@ +#!/bin/bash + +nvcc -std=c++17 --extended-lambda --expt-relaxed-constexpr \ + -Xcompiler -Wno-deprecated-declarations -Xcompiler -fopenmp \ + -o repro_kmeans_mg repro_kmeans_mg.cu \ + -I$CONDA_PREFIX/include \ + -I$CONDA_PREFIX/include/rapids \ + -I$CONDA_PREFIX/include/rapids/libcudacxx \ + -L$CONDA_PREFIX/lib \ + -lcuvs -lrmm -lnccl -lucp -lucs -lucxx -lgomp \ + -DRAFT_SYSTEM_LITTLE_ENDIAN=1 -DLIBCUDACXX_ENABLE_EXPERIMENTAL_MEMORY_RESOURCE \ + -Xlinker -rpath=$CONDA_PREFIX/lib \ + -gencode arch=compute_XX,code=sm_XX From 7b2cf3b4687078b5e9d6fc5bef2d16b64e0f4824 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Sat, 7 Mar 2026 22:50:09 +0000 Subject: [PATCH 04/20] Extracted weights and wrote a RAFT discrete() reproducer --- build_repro.sh | 16 ++++- cpp/src/cluster/detail/kmeans.cuh | 14 +++- cpp/src/cluster/detail/kmeans_mg.cuh | 2 + raft_discrete.cu | 103 +++++++++++++++++++++++++++ repro_kmeans_mg.cu | 2 +- 5 files changed, 133 insertions(+), 4 deletions(-) create mode 100644 raft_discrete.cu diff --git a/build_repro.sh b/build_repro.sh index 19e827cad7..6969b55229 100644 --- a/build_repro.sh +++ b/build_repro.sh @@ -1,8 +1,20 @@ #!/bin/bash +# nvcc -std=c++17 --extended-lambda --expt-relaxed-constexpr \ +# -Xcompiler -Wno-deprecated-declarations -Xcompiler -fopenmp \ +# -o repro_kmeans_mg ../repro_kmeans_mg.cu \ +# -I$CONDA_PREFIX/include \ +# -I$CONDA_PREFIX/include/rapids \ +# -I$CONDA_PREFIX/include/rapids/libcudacxx \ +# -L$CONDA_PREFIX/lib \ +# -lcuvs -lrmm -lnccl -lucp -lucs -lucxx -lgomp \ +# -DRAFT_SYSTEM_LITTLE_ENDIAN=1 -DLIBCUDACXX_ENABLE_EXPERIMENTAL_MEMORY_RESOURCE \ +# -Xlinker -rpath=$CONDA_PREFIX/lib \ +# -gencode arch=compute_100,code=sm_100 + nvcc -std=c++17 --extended-lambda --expt-relaxed-constexpr \ -Xcompiler -Wno-deprecated-declarations -Xcompiler -fopenmp \ - -o repro_kmeans_mg repro_kmeans_mg.cu \ + -o raft_discrete ../raft_discrete.cu \ -I$CONDA_PREFIX/include \ -I$CONDA_PREFIX/include/rapids \ -I$CONDA_PREFIX/include/rapids/libcudacxx \ @@ -10,4 +22,4 @@ nvcc -std=c++17 --extended-lambda --expt-relaxed-constexpr \ -lcuvs -lrmm -lnccl -lucp -lucs -lucxx -lgomp \ -DRAFT_SYSTEM_LITTLE_ENDIAN=1 -DLIBCUDACXX_ENABLE_EXPERIMENTAL_MEMORY_RESOURCE \ -Xlinker -rpath=$CONDA_PREFIX/lib \ - -gencode arch=compute_XX,code=sm_XX + -gencode arch=compute_100,code=sm_100 diff --git a/cpp/src/cluster/detail/kmeans.cuh b/cpp/src/cluster/detail/kmeans.cuh index dee8565f5c..607eb80b01 100644 --- a/cpp/src/cluster/detail/kmeans.cuh +++ b/cpp/src/cluster/detail/kmeans.cuh @@ -52,6 +52,8 @@ #include #include +#include "/home/vinayd/nwork/git-repos/snippets/ketu.cuh" + namespace cuvs::cluster::kmeans::detail { static const std::string CUVS_NAME = "cuvs"; @@ -187,6 +189,7 @@ void kmeansPlusPlus(raft::resources const& handle, RAFT_LOG_DEBUG(" k-means++ - Sampled %d/%d centroids", n_clusters_picked, n_clusters); + static ketu::DataTracker weights_tracker("weight_tracker_phi"); // <<<< Step-2 >>> : while |C| < k while (n_clusters_picked < n_clusters) { // <<< Step-3 >>> : Sample x in X with probability p_x = d^2(x, C) / phi_X (C) @@ -198,7 +201,16 @@ void kmeansPlusPlus(raft::resources const& handle, raft::matrix::gather(handle, const_X_view, const_indices_view, candidates_view); cudaDeviceSynchronize(); - dump_var<<<1, 1>>>(indices.data_handle(), uint64_t(n_trials)); + //dump_var<<<1, 1>>>(indices.data_handle(), uint64_t(n_trials)); + printf("n_samples = %lu\n", uint64_t(n_samples)); + // Create a CPU vector and copy indices from GPU to CPU + std::vector h_weights(n_samples); + raft::copy(h_weights.data(), minClusterDistance.data_handle(), n_samples, stream); + raft::resource::sync_stream(handle, stream); + if (weights_tracker.is_changed(h_weights.data(), n_samples)) { + weights_tracker.update(h_weights.data(), n_samples); + weights_tracker.dump_to_file(); + } cudaDeviceSynchronize(); // Calculate pairwise distance between X and the centroid candidates // Output - pwd [n_trials x n_samples] diff --git a/cpp/src/cluster/detail/kmeans_mg.cuh b/cpp/src/cluster/detail/kmeans_mg.cuh index 2125a297ad..52d084f164 100644 --- a/cpp/src/cluster/detail/kmeans_mg.cuh +++ b/cpp/src/cluster/detail/kmeans_mg.cuh @@ -37,6 +37,8 @@ #include #include +#include "/home/vinayd/nwork/git-repos/snippets/ketu.cuh" + namespace cuvs::cluster::kmeans::mg::detail { #define CUVS_LOG_KMEANS(handle, fmt, ...) \ diff --git a/raft_discrete.cu b/raft_discrete.cu new file mode 100644 index 0000000000..05b5978b34 --- /dev/null +++ b/raft_discrete.cu @@ -0,0 +1,103 @@ +/* + * Demonstration of raft::random::discrete API usage + * + * The discrete API samples indices from a discrete probability distribution + * defined by weights. Each index i is sampled with probability proportional + * to weights[i]. + * + * Compile with: + * nvcc -o raft_discrete raft_discrete.cu -I --expt-extended-lambda + */ + +#include +#include +#include +#include + +#include + +#include +#include +#include +#include +#include +#include + +int main() +{ + // Create RAFT handle (manages CUDA resources) + raft::handle_t handle; + auto stream = raft::resource::get_cuda_stream(handle); + + // Read weights from file (one weight per line) + std::vector h_weights; + { + std::ifstream weights_file("weight_tracker_phi_00001.txt"); + if (!weights_file.is_open()) { + std::cerr << "Error: Could not open weights.txt" << std::endl; + return 1; + } + std::string line; + while (std::getline(weights_file, line)) { + if (!line.empty()) { + h_weights.push_back(std::stof(line)); + } + } + } + + // Number of categories (weights) + const int n_categories = static_cast(h_weights.size()); + // Number of samples to draw + const int n_samples = 9; + rmm::device_uvector d_weights(n_categories, stream); + raft::copy(d_weights.data(), h_weights.data(), n_categories, stream); + + // Create output array for sampled indices + rmm::device_uvector d_indices(n_samples, stream); + + // Create random number generator with a unique seed each run + auto seed = static_cast(std::chrono::high_resolution_clock::now().time_since_epoch().count()) ^ + static_cast(std::random_device{}()); + raft::random::RngState rng(seed, raft::random::GenPC); + + // Create mdspan views for the API + auto indices_view = raft::make_device_vector_view(d_indices.data(), n_samples); + auto weights_view = + raft::make_device_vector_view(d_weights.data(), n_categories); + + // Sample indices according to the weight distribution + // Each index i will be sampled with probability weights[i] / sum(weights) + raft::random::discrete(handle, rng, indices_view, weights_view); + + // Copy results back to host + std::vector h_indices(n_samples); + raft::copy(h_indices.data(), d_indices.data(), n_samples, stream); + raft::resource::sync_stream(handle, stream); + + // Print results + std::cout << "Weights: "; + for (int i = 0; i < n_categories && i < 10; ++i) { + std::cout << h_weights[i] << " "; + } + std::cout << std::endl; + + std::cout << "Sampled indices: "; + for (int i = 0; i < n_samples; ++i) { + std::cout << h_indices[i] << " "; + } + std::cout << std::endl; + + // Count occurrences of each index + std::vector counts(n_categories, 0); + for (int i = 0; i < n_samples; ++i) { + counts[h_indices[i]]++; + } + + // std::cout << "Counts per category: "; + // for (int i = 0; i < n_categories; ++i) { + // std::cout << counts[i] << " "; + // } + // std::cout << std::endl; + + return 0; +} diff --git a/repro_kmeans_mg.cu b/repro_kmeans_mg.cu index 332ae50b92..dfecd31a37 100644 --- a/repro_kmeans_mg.cu +++ b/repro_kmeans_mg.cu @@ -197,7 +197,7 @@ int main(int argc, char* argv[]) const char* name; }; RngConfig rng_configs[] = { - {raft::random::GeneratorType::GenPC, "GenPC"}, + //{raft::random::GeneratorType::GenPC, "GenPC"}, {raft::random::GeneratorType::GenPhilox, "GenPhilox"}, }; From 908d6b8a1bf410b43741b4bda2df5ba876043f39 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 9 Mar 2026 11:28:14 +0000 Subject: [PATCH 05/20] Compare discrete result with weights array --- raft_discrete.cu | 97 ++++++++++++++++++++++++++++++++---------------- 1 file changed, 64 insertions(+), 33 deletions(-) diff --git a/raft_discrete.cu b/raft_discrete.cu index 05b5978b34..470eabe4b5 100644 --- a/raft_discrete.cu +++ b/raft_discrete.cu @@ -17,6 +17,7 @@ #include #include +#include #include #include #include @@ -31,20 +32,21 @@ int main() // Read weights from file (one weight per line) std::vector h_weights; - { - std::ifstream weights_file("weight_tracker_phi_00001.txt"); - if (!weights_file.is_open()) { - std::cerr << "Error: Could not open weights.txt" << std::endl; - return 1; - } - std::string line; - while (std::getline(weights_file, line)) { - if (!line.empty()) { - h_weights.push_back(std::stof(line)); - } + const char* weights_filename_env = std::getenv("WEIGHTS_FILE"); + std::string weights_filename = weights_filename_env ? weights_filename_env : "weight_tracker_phi_00001.txt"; + std::ifstream weights_file(weights_filename); + if (!weights_file.is_open()) { + std::cerr << "Error: Could not open " << weights_filename << std::endl; + return 1; + } + std::string line; + while (std::getline(weights_file, line)) { + if (!line.empty()) { + h_weights.push_back(std::stof(line)); } } + // Number of categories (weights) const int n_categories = static_cast(h_weights.size()); // Number of samples to draw @@ -58,7 +60,9 @@ int main() // Create random number generator with a unique seed each run auto seed = static_cast(std::chrono::high_resolution_clock::now().time_since_epoch().count()) ^ static_cast(std::random_device{}()); - raft::random::RngState rng(seed, raft::random::GenPC); + const char* use_pc_env = std::getenv("USE_PC"); + bool use_pc = use_pc_env && std::string(use_pc_env) == "1"; + raft::random::RngState rng(seed, use_pc ? raft::random::GenPC : raft::random::GenPhilox); // Create mdspan views for the API auto indices_view = raft::make_device_vector_view(d_indices.data(), n_samples); @@ -67,37 +71,64 @@ int main() // Sample indices according to the weight distribution // Each index i will be sampled with probability weights[i] / sum(weights) - raft::random::discrete(handle, rng, indices_view, weights_view); - // Copy results back to host + // Count occurrences of each index + std::vector counts(n_categories, 0); std::vector h_indices(n_samples); - raft::copy(h_indices.data(), d_indices.data(), n_samples, stream); - raft::resource::sync_stream(handle, stream); + + // Number of iterations to run + const int n_iterations = 1000; + + for (int iter = 0; iter < n_iterations; ++iter) { + raft::random::discrete(handle, rng, indices_view, weights_view); + + // Copy results back to host + raft::copy(h_indices.data(), d_indices.data(), n_samples, stream); + raft::resource::sync_stream(handle, stream); + + // Update counts + for (int i = 0; i < n_samples; ++i) { + counts[h_indices[i]]++; + } + } // Print results - std::cout << "Weights: "; - for (int i = 0; i < n_categories && i < 10; ++i) { - std::cout << h_weights[i] << " "; + // std::cout << "Weights: "; + // for (int i = 0; i < n_categories && i < 10; ++i) { + // std::cout << h_weights[i] << " "; + // } + // std::cout << std::endl; + + //std::cout << "Total samples: " << n_iterations * n_samples << std::endl; + + // Normalize counts array + double counts_sum = 0.0; + for (int i = 0; i < n_categories; ++i) { + counts_sum += counts[i]; + } + std::vector normalized_counts(n_categories); + for (int i = 0; i < n_categories; ++i) { + normalized_counts[i] = static_cast(counts[i]) / counts_sum; } - std::cout << std::endl; - std::cout << "Sampled indices: "; - for (int i = 0; i < n_samples; ++i) { - std::cout << h_indices[i] << " "; + // Normalize weights array + double weights_sum = 0.0; + for (int i = 0; i < n_categories; ++i) { + weights_sum += h_weights[i]; + } + std::vector normalized_weights(n_categories); + for (int i = 0; i < n_categories; ++i) { + normalized_weights[i] = static_cast(h_weights[i]) / weights_sum; } - std::cout << std::endl; - // Count occurrences of each index - std::vector counts(n_categories, 0); - for (int i = 0; i < n_samples; ++i) { - counts[h_indices[i]]++; + // Calculate sum of squared differences + double sum_sq_diff = 0.0; + for (int i = 0; i < n_categories; ++i) { + double diff = normalized_counts[i] - normalized_weights[i]; + sum_sq_diff += diff * diff; } - // std::cout << "Counts per category: "; - // for (int i = 0; i < n_categories; ++i) { - // std::cout << counts[i] << " "; - // } - // std::cout << std::endl; + std::cout << weights_filename << " " << "Sum of squared differences: " << sum_sq_diff << std::endl; return 0; } From 70bac348ab95ebe7c1cc688a90257e5a67e9bbca Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 9 Mar 2026 15:18:48 +0000 Subject: [PATCH 06/20] Some more debugging --- build_repro.sh | 26 +++++++++++++------------- cpp/src/cluster/detail/kmeans.cuh | 10 +++++----- cpp/src/cluster/detail/kmeans_mg.cuh | 8 +++----- repro_kmeans_mg.cu | 2 +- 4 files changed, 22 insertions(+), 24 deletions(-) diff --git a/build_repro.sh b/build_repro.sh index 6969b55229..76cec02cc6 100644 --- a/build_repro.sh +++ b/build_repro.sh @@ -1,8 +1,20 @@ #!/bin/bash +nvcc -g -std=c++17 --extended-lambda --expt-relaxed-constexpr \ + -Xcompiler -Wno-deprecated-declarations -Xcompiler -fopenmp \ + -o repro_kmeans_mg ../repro_kmeans_mg.cu \ + -I$CONDA_PREFIX/include \ + -I$CONDA_PREFIX/include/rapids \ + -I$CONDA_PREFIX/include/rapids/libcudacxx \ + -L$CONDA_PREFIX/lib \ + -lcuvs -lrmm -lnccl -lucp -lucs -lucxx -lgomp \ + -DRAFT_SYSTEM_LITTLE_ENDIAN=1 -DLIBCUDACXX_ENABLE_EXPERIMENTAL_MEMORY_RESOURCE \ + -Xlinker -rpath=$CONDA_PREFIX/lib \ + -gencode arch=compute_100,code=sm_100 + # nvcc -std=c++17 --extended-lambda --expt-relaxed-constexpr \ # -Xcompiler -Wno-deprecated-declarations -Xcompiler -fopenmp \ -# -o repro_kmeans_mg ../repro_kmeans_mg.cu \ +# -o raft_discrete ../raft_discrete.cu \ # -I$CONDA_PREFIX/include \ # -I$CONDA_PREFIX/include/rapids \ # -I$CONDA_PREFIX/include/rapids/libcudacxx \ @@ -11,15 +23,3 @@ # -DRAFT_SYSTEM_LITTLE_ENDIAN=1 -DLIBCUDACXX_ENABLE_EXPERIMENTAL_MEMORY_RESOURCE \ # -Xlinker -rpath=$CONDA_PREFIX/lib \ # -gencode arch=compute_100,code=sm_100 - -nvcc -std=c++17 --extended-lambda --expt-relaxed-constexpr \ - -Xcompiler -Wno-deprecated-declarations -Xcompiler -fopenmp \ - -o raft_discrete ../raft_discrete.cu \ - -I$CONDA_PREFIX/include \ - -I$CONDA_PREFIX/include/rapids \ - -I$CONDA_PREFIX/include/rapids/libcudacxx \ - -L$CONDA_PREFIX/lib \ - -lcuvs -lrmm -lnccl -lucp -lucs -lucxx -lgomp \ - -DRAFT_SYSTEM_LITTLE_ENDIAN=1 -DLIBCUDACXX_ENABLE_EXPERIMENTAL_MEMORY_RESOURCE \ - -Xlinker -rpath=$CONDA_PREFIX/lib \ - -gencode arch=compute_100,code=sm_100 diff --git a/cpp/src/cluster/detail/kmeans.cuh b/cpp/src/cluster/detail/kmeans.cuh index 607eb80b01..14e6e51227 100644 --- a/cpp/src/cluster/detail/kmeans.cuh +++ b/cpp/src/cluster/detail/kmeans.cuh @@ -189,7 +189,7 @@ void kmeansPlusPlus(raft::resources const& handle, RAFT_LOG_DEBUG(" k-means++ - Sampled %d/%d centroids", n_clusters_picked, n_clusters); - static ketu::DataTracker weights_tracker("weight_tracker_phi"); + //static ketu::DataTracker weights_tracker("weight_tracker_phi"); // <<<< Step-2 >>> : while |C| < k while (n_clusters_picked < n_clusters) { // <<< Step-3 >>> : Sample x in X with probability p_x = d^2(x, C) / phi_X (C) @@ -207,10 +207,10 @@ void kmeansPlusPlus(raft::resources const& handle, std::vector h_weights(n_samples); raft::copy(h_weights.data(), minClusterDistance.data_handle(), n_samples, stream); raft::resource::sync_stream(handle, stream); - if (weights_tracker.is_changed(h_weights.data(), n_samples)) { - weights_tracker.update(h_weights.data(), n_samples); - weights_tracker.dump_to_file(); - } + // if (weights_tracker.is_changed(h_weights.data(), n_samples)) { + // weights_tracker.update(h_weights.data(), n_samples); + // weights_tracker.dump_to_file(); + // } cudaDeviceSynchronize(); // Calculate pairwise distance between X and the centroid candidates // Output - pwd [n_trials x n_samples] diff --git a/cpp/src/cluster/detail/kmeans_mg.cuh b/cpp/src/cluster/detail/kmeans_mg.cuh index 52d084f164..6bd6ceef14 100644 --- a/cpp/src/cluster/detail/kmeans_mg.cuh +++ b/cpp/src/cluster/detail/kmeans_mg.cuh @@ -37,8 +37,6 @@ #include #include -#include "/home/vinayd/nwork/git-repos/snippets/ketu.cuh" - namespace cuvs::cluster::kmeans::mg::detail { #define CUVS_LOG_KMEANS(handle, fmt, ...) \ @@ -146,7 +144,6 @@ void initKMeansPlusPlus(const raft::resources& handle, auto n_clusters = params.n_clusters; auto metric = params.metric; - raft::random::RngState rng(params.rng_state.seed, raft::random::GeneratorType::GenPhilox); // <<<< Step-1 >>> : C <- sample a point uniformly at random from X // 1.1 - Select a rank r' at random from the available n_rank ranks with a @@ -184,7 +181,7 @@ void initKMeansPlusPlus(const raft::resources& handle, // 1.2 - Rank r' samples a point uniformly at random from the local dataset // X which will be used as the initial centroid for kmeans++ if (my_rank == rp) { - std::mt19937 gen(params.rng_state.seed); + std::mt19937 gen(params.rng_state.seed + 31415926); std::uniform_int_distribution<> dis(0, n_samples - 1); int cIdx = dis(gen); @@ -319,8 +316,9 @@ void initKMeansPlusPlus(const raft::resources& handle, // <<<< Step-4 >>> : Sample each point x in X independently and identify new // potentialCentroids + //raft::random::RngState rng(params.rng_state.seed, raft::random::GeneratorType::GenPhilox); raft::random::uniform( - handle, rng, uniformRands.data_handle(), uniformRands.extent(0), (DataT)0, (DataT)1); + handle, params.rng_state, uniformRands.data_handle(), uniformRands.extent(0), (DataT)0, (DataT)1); cuvs::cluster::kmeans::SamplingOp select_op(psi, params.oversampling_factor, n_clusters, diff --git a/repro_kmeans_mg.cu b/repro_kmeans_mg.cu index dfecd31a37..332ae50b92 100644 --- a/repro_kmeans_mg.cu +++ b/repro_kmeans_mg.cu @@ -197,7 +197,7 @@ int main(int argc, char* argv[]) const char* name; }; RngConfig rng_configs[] = { - //{raft::random::GeneratorType::GenPC, "GenPC"}, + {raft::random::GeneratorType::GenPC, "GenPC"}, {raft::random::GeneratorType::GenPhilox, "GenPhilox"}, }; From 9d728c3d06e34a51686cf6f897af55a133615b37 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 16 Mar 2026 14:14:49 +0000 Subject: [PATCH 07/20] Correctly seeding at multiple places --- cpp/src/cluster/detail/kmeans.cuh | 10 ++++--- cpp/src/cluster/detail/kmeans_mg.cuh | 43 +++++++++++++++++++--------- 2 files changed, 35 insertions(+), 18 deletions(-) diff --git a/cpp/src/cluster/detail/kmeans.cuh b/cpp/src/cluster/detail/kmeans.cuh index 14e6e51227..8bbf2659dd 100644 --- a/cpp/src/cluster/detail/kmeans.cuh +++ b/cpp/src/cluster/detail/kmeans.cuh @@ -155,14 +155,14 @@ void kmeansPlusPlus(raft::resources const& handle, raft::linalg::norm(handle, X, L2NormX.view()); } - raft::random::RngState rng(params.rng_state.seed, params.rng_state.type); + std::mt19937_64 gen_64(params.rng_state.seed); //raft::random::RngState rng(params.rng_state.seed, raft::random::GeneratorType::GenPhilox); - std::mt19937 gen(params.rng_state.seed); + //std::mt19937 gen(params.rng_state.seed); std::uniform_int_distribution<> dis(0, n_samples - 1); // <<< Step-1 >>>: C <-- sample a point uniformly at random from X auto initialCentroid = raft::make_device_matrix_view( - X.data_handle() + dis(gen) * n_features, 1, n_features); + X.data_handle() + dis(gen_64) * n_features, 1, n_features); int n_clusters_picked = 1; // store the chosen centroid in the buffer @@ -197,12 +197,14 @@ void kmeansPlusPlus(raft::resources const& handle, // distance to the nearest existing cluster // printf("size of index = %lu, size of weights = %lu, sizeof types (%ld, %ld)\n", // uint64_t(indices_view.extent(0)), uint64_t(const_weights_view.extent(0)), sizeof(DataT), sizeof(IndexT)); + uint64_t gpu_seed = gen_64(); + raft::random::RngState rng(gpu_seed, params.rng_state.type); raft::random::discrete(handle, rng, indices_view, const_weights_view); raft::matrix::gather(handle, const_X_view, const_indices_view, candidates_view); cudaDeviceSynchronize(); //dump_var<<<1, 1>>>(indices.data_handle(), uint64_t(n_trials)); - printf("n_samples = %lu\n", uint64_t(n_samples)); + //printf("n_samples = %lu\n", uint64_t(n_samples)); // Create a CPU vector and copy indices from GPU to CPU std::vector h_weights(n_samples); raft::copy(h_weights.data(), minClusterDistance.data_handle(), n_samples, stream); diff --git a/cpp/src/cluster/detail/kmeans_mg.cuh b/cpp/src/cluster/detail/kmeans_mg.cuh index 6bd6ceef14..c6de34ae34 100644 --- a/cpp/src/cluster/detail/kmeans_mg.cuh +++ b/cpp/src/cluster/detail/kmeans_mg.cuh @@ -58,6 +58,7 @@ static cuvs::cluster::kmeans::params default_params; template void initRandom(const raft::resources& handle, const cuvs::cluster::kmeans::params& params, + std::mt19937_64& gen_64, raft::device_matrix_view X, raft::device_matrix_view centroids) { @@ -96,8 +97,9 @@ void initRandom(const raft::resources& handle, auto centroidsSampledInRank = raft::make_device_matrix(handle, nCentroidsSampledInRank, n_features); + uint64_t gpu_seed = gen_64(); cuvs::cluster::kmeans::shuffle_and_gather( - handle, X, centroidsSampledInRank.view(), nCentroidsSampledInRank, params.rng_state.seed); + handle, X, centroidsSampledInRank.view(), nCentroidsSampledInRank, gpu_seed); std::vector displs(n_ranks); std::exclusive_scan(nCentroidsElementsToReceiveFromRank.begin(), @@ -130,6 +132,7 @@ void initRandom(const raft::resources& handle, template void initKMeansPlusPlus(const raft::resources& handle, const cuvs::cluster::kmeans::params& params, + std::mt19937_64& gen_64, raft::device_matrix_view X, raft::device_matrix_view centroidsRawData, rmm::device_uvector& workspace) @@ -156,9 +159,9 @@ void initKMeansPlusPlus(const raft::resources& handle, // Choose rp on rank 0 and broadcast to all ranks to guarantee agreement int rp = 0; if (my_rank == KMEANS_COMM_ROOT) { - std::mt19937 gen(params.rng_state.seed); + //std::mt19937 gen(params.rng_state.seed); std::uniform_int_distribution<> dis(0, n_rank - 1); - rp = dis(gen); + rp = dis(gen_64); } { rmm::device_scalar rp_d(stream); @@ -181,10 +184,10 @@ void initKMeansPlusPlus(const raft::resources& handle, // 1.2 - Rank r' samples a point uniformly at random from the local dataset // X which will be used as the initial centroid for kmeans++ if (my_rank == rp) { - std::mt19937 gen(params.rng_state.seed + 31415926); + //std::mt19937 gen(params.rng_state.seed + 31415926); std::uniform_int_distribution<> dis(0, n_samples - 1); - int cIdx = dis(gen); + int cIdx = dis(gen_64); auto centroidsView = raft::make_device_matrix_view( X.data_handle() + cIdx * n_features, 1, n_features); @@ -316,9 +319,11 @@ void initKMeansPlusPlus(const raft::resources& handle, // <<<< Step-4 >>> : Sample each point x in X independently and identify new // potentialCentroids - //raft::random::RngState rng(params.rng_state.seed, raft::random::GeneratorType::GenPhilox); + uint64_t gpu_seed; + gpu_seed = gen_64(); + raft::random::RngState rng(gpu_seed, raft::random::GeneratorType::GenPhilox); raft::random::uniform( - handle, params.rng_state, uniformRands.data_handle(), uniformRands.extent(0), (DataT)0, (DataT)1); + handle, rng, uniformRands.data_handle(), uniformRands.extent(0), (DataT)0, (DataT)1); cuvs::cluster::kmeans::SamplingOp select_op(psi, params.oversampling_factor, n_clusters, @@ -407,16 +412,19 @@ void initKMeansPlusPlus(const raft::resources& handle, auto const_centroids = raft::make_device_matrix_view( potentialCentroids.data_handle(), potentialCentroids.extent(0), potentialCentroids.extent(1)); auto nvtx_range = nvtx3::start_range("init_plus_plus"); + auto params_copy = params; + params_copy.rng_state.seed = gen_64(); cuvs::cluster::kmeans::init_plus_plus( - handle, params, const_centroids, centroidsRawData, workspace); + handle, params_copy, const_centroids, centroidsRawData, workspace); nvtx3::end_range(nvtx_range); auto inertia = raft::make_host_scalar(0); auto n_iter = raft::make_host_scalar(0); auto weight_view = raft::make_device_vector_view(weight.data_handle(), weight.extent(0)); - cuvs::cluster::kmeans::params params_copy = params; - params_copy.rng_state = default_params.rng_state; + //cuvs::cluster::kmeans::params params_copy = params; + //params_copy.rng_state = default_params.rng_state; + params_copy.rng_state.seed = gen_64(); cuvs::cluster::kmeans::fit_main(handle, params_copy, const_centroids, @@ -439,10 +447,11 @@ void initKMeansPlusPlus(const raft::resources& handle, // generate `n_random_clusters` centroids cuvs::cluster::kmeans::params rand_params = params; - rand_params.rng_state = default_params.rng_state; + //rand_params.rng_state = default_params.rng_state; + rand_params.rng_state.seed = gen_64(); rand_params.init = cuvs::cluster::kmeans::params::InitMethod::Random; rand_params.n_clusters = n_random_clusters; - initRandom(handle, rand_params, X, centroidsRawData); + initRandom(handle, rand_params, gen_64, X, centroidsRawData); // copy centroids generated during kmeans|| iteration to the buffer raft::copy( @@ -513,6 +522,12 @@ void fit(const raft::resources& handle, auto n_clusters = params.n_clusters; auto metric = params.metric; + const int my_rank = comm.get_rank(); + const int n_ranks = comm.get_size(); + + std::mt19937_64 gen_64(params.rng_state.seed + (uint64_t(my_rank) << 32)); + printf("I am rank %d of %d total ranks\n", my_rank, n_ranks); + auto weight = raft::make_device_vector(handle, n_samples); if (sample_weight) { raft::copy(handle, weight.view(), sample_weight.value()); @@ -528,11 +543,11 @@ void fit(const raft::resources& handle, CUVS_LOG_KMEANS(handle, "KMeans.fit: initialize cluster centers by randomly choosing from the " "input data.\n"); - initRandom(handle, params, X, centroids); + initRandom(handle, params, gen_64, X, centroids); } else if (params.init == cuvs::cluster::kmeans::params::InitMethod::KMeansPlusPlus) { // default method to initialize is kmeans++ CUVS_LOG_KMEANS(handle, "KMeans.fit: initialize cluster centers using k-means++ algorithm.\n"); - initKMeansPlusPlus(handle, params, X, centroids, workspace); + initKMeansPlusPlus(handle, params, gen_64, X, centroids, workspace); } else if (params.init == cuvs::cluster::kmeans::params::InitMethod::Array) { CUVS_LOG_KMEANS(handle, "KMeans.fit: initialize cluster centers from the ndarray array input " From afaf8f4f8bc5c3df790e8c7d763ec4c6e5d7875e Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 16 Mar 2026 15:25:11 +0000 Subject: [PATCH 08/20] Randomized seeeding --- cpp/src/cluster/detail/kmeans_mg.cuh | 1 + repro_kmeans_mg.cu | 5 ++++- 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/cpp/src/cluster/detail/kmeans_mg.cuh b/cpp/src/cluster/detail/kmeans_mg.cuh index c6de34ae34..f4f8ccacf4 100644 --- a/cpp/src/cluster/detail/kmeans_mg.cuh +++ b/cpp/src/cluster/detail/kmeans_mg.cuh @@ -425,6 +425,7 @@ void initKMeansPlusPlus(const raft::resources& handle, //cuvs::cluster::kmeans::params params_copy = params; //params_copy.rng_state = default_params.rng_state; params_copy.rng_state.seed = gen_64(); + std::printf("params_copy.rng_state.seed = %lu\n", static_cast(params_copy.rng_state.seed)); cuvs::cluster::kmeans::fit_main(handle, params_copy, const_centroids, diff --git a/repro_kmeans_mg.cu b/repro_kmeans_mg.cu index 332ae50b92..e939a8333b 100644 --- a/repro_kmeans_mg.cu +++ b/repro_kmeans_mg.cu @@ -201,6 +201,8 @@ int main(int argc, char* argv[]) {raft::random::GeneratorType::GenPhilox, "GenPhilox"}, }; + std::mt19937_64 gen_64(std::chrono::high_resolution_clock::now().time_since_epoch().count()); + uint64_t last_seed = gen_64(); for (const auto& rng_cfg : rng_configs) { // Create NVTX range for this RNG configuration nvtx3::scoped_range range(rng_cfg.name); @@ -219,7 +221,7 @@ int main(int argc, char* argv[]) params.n_clusters = static_cast(N_CLUSTERS); params.max_iter = 300; params.tol = 1e-4; - params.rng_state.seed = 42; + params.rng_state.seed = last_seed; params.rng_state.type = rng_cfg.type; params.oversampling_factor = 2.0; params.n_init = 1; @@ -277,6 +279,7 @@ int main(int argc, char* argv[]) if (rank == 0) { double avg_iters = static_cast(total_kmeans_iters) / N_BENCHMARK_ITERS; + std::printf("Last seed used: %lu\n", last_seed); std::printf("Average KMeans iterations per run (%s): %.1f\n", rng_cfg.name, avg_iters); } } From da9093855beecf1574a92478151947e8a5e78321 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 23 Mar 2026 13:24:53 +0000 Subject: [PATCH 09/20] Removing debug code --- cpp/src/cluster/detail/kmeans.cuh | 19 ------------------- 1 file changed, 19 deletions(-) diff --git a/cpp/src/cluster/detail/kmeans.cuh b/cpp/src/cluster/detail/kmeans.cuh index 8bbf2659dd..1bdc6a6e65 100644 --- a/cpp/src/cluster/detail/kmeans.cuh +++ b/cpp/src/cluster/detail/kmeans.cuh @@ -52,7 +52,6 @@ #include #include -#include "/home/vinayd/nwork/git-repos/snippets/ketu.cuh" namespace cuvs::cluster::kmeans::detail { @@ -75,16 +74,6 @@ void initRandom(raft::resources const& handle, handle, X, centroids, n_clusters, params.rng_state.seed); } -template -__global__ void dump_var(DataT* v, uint64_t len) { - printf("{"); - for (int i = 0; i < 10; i++) { - if (i < len) printf("%f, ", double(v[i])); - } - printf("} \n"); -} - - /* * @brief Selects 'n_clusters' samples from the input X using kmeans++ algorithm. @@ -202,18 +191,10 @@ void kmeansPlusPlus(raft::resources const& handle, raft::random::discrete(handle, rng, indices_view, const_weights_view); raft::matrix::gather(handle, const_X_view, const_indices_view, candidates_view); - cudaDeviceSynchronize(); - //dump_var<<<1, 1>>>(indices.data_handle(), uint64_t(n_trials)); - //printf("n_samples = %lu\n", uint64_t(n_samples)); // Create a CPU vector and copy indices from GPU to CPU std::vector h_weights(n_samples); raft::copy(h_weights.data(), minClusterDistance.data_handle(), n_samples, stream); raft::resource::sync_stream(handle, stream); - // if (weights_tracker.is_changed(h_weights.data(), n_samples)) { - // weights_tracker.update(h_weights.data(), n_samples); - // weights_tracker.dump_to_file(); - // } - cudaDeviceSynchronize(); // Calculate pairwise distance between X and the centroid candidates // Output - pwd [n_trials x n_samples] auto pwd = distBuffer.view(); From c7854e060a2445834caeed4a9a3dd9dff387bef3 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 23 Mar 2026 13:25:24 +0000 Subject: [PATCH 10/20] Removing trailing spaces --- repro_kmeans_mg.cu | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/repro_kmeans_mg.cu b/repro_kmeans_mg.cu index e939a8333b..14815e4695 100644 --- a/repro_kmeans_mg.cu +++ b/repro_kmeans_mg.cu @@ -68,7 +68,7 @@ int main(int argc, char* argv[]) // Parse command line arguments int64_t n_samples_total = 10000000; // Default value int64_t n_features = 256; // Default value - + for (int i = 1; i < argc; i++) { std::string arg = argv[i]; if (arg == "--samples" || arg == "-s") { @@ -94,7 +94,7 @@ int main(int argc, char* argv[]) return 0; } } - + // Detect available GPUs int n_devices = 0; CUDACHECK(cudaGetDeviceCount(&n_devices)); @@ -206,7 +206,7 @@ int main(int argc, char* argv[]) for (const auto& rng_cfg : rng_configs) { // Create NVTX range for this RNG configuration nvtx3::scoped_range range(rng_cfg.name); - + if (rank == 0) { std::printf("\n=== RNG: %s ===\n", rng_cfg.name); } @@ -236,7 +236,7 @@ int main(int argc, char* argv[]) d_labels.data(), n_local); auto t0 = std::chrono::high_resolution_clock::now(); - + { nvtx3::scoped_range fit_range("kmeans_fit"); cuvs::cluster::kmeans::fit(handle, From 923fae51810e09e9e041b838b220611d0f24449b Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 23 Mar 2026 13:27:00 +0000 Subject: [PATCH 11/20] Removing debug prints --- cpp/src/cluster/detail/kmeans_mg.cuh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cpp/src/cluster/detail/kmeans_mg.cuh b/cpp/src/cluster/detail/kmeans_mg.cuh index f4f8ccacf4..8577b1d365 100644 --- a/cpp/src/cluster/detail/kmeans_mg.cuh +++ b/cpp/src/cluster/detail/kmeans_mg.cuh @@ -424,8 +424,8 @@ void initKMeansPlusPlus(const raft::resources& handle, raft::make_device_vector_view(weight.data_handle(), weight.extent(0)); //cuvs::cluster::kmeans::params params_copy = params; //params_copy.rng_state = default_params.rng_state; + // Update the seed one more time params_copy.rng_state.seed = gen_64(); - std::printf("params_copy.rng_state.seed = %lu\n", static_cast(params_copy.rng_state.seed)); cuvs::cluster::kmeans::fit_main(handle, params_copy, const_centroids, From 9b2a34ff91f415bed258a1a286c5df2f44af92bf Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 23 Mar 2026 13:30:48 +0000 Subject: [PATCH 12/20] Removing debug code --- cpp/src/cluster/detail/kmeans.cuh | 1 - 1 file changed, 1 deletion(-) diff --git a/cpp/src/cluster/detail/kmeans.cuh b/cpp/src/cluster/detail/kmeans.cuh index 1bdc6a6e65..a1378d7533 100644 --- a/cpp/src/cluster/detail/kmeans.cuh +++ b/cpp/src/cluster/detail/kmeans.cuh @@ -178,7 +178,6 @@ void kmeansPlusPlus(raft::resources const& handle, RAFT_LOG_DEBUG(" k-means++ - Sampled %d/%d centroids", n_clusters_picked, n_clusters); - //static ketu::DataTracker weights_tracker("weight_tracker_phi"); // <<<< Step-2 >>> : while |C| < k while (n_clusters_picked < n_clusters) { // <<< Step-3 >>> : Sample x in X with probability p_x = d^2(x, C) / phi_X (C) From 78305d95cafa054ae33744fbb00f810aed6e117b Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 23 Mar 2026 13:31:51 +0000 Subject: [PATCH 13/20] Undoing cmake change --- cpp/cmake/modules/ConfigureCUDA.cmake | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/cpp/cmake/modules/ConfigureCUDA.cmake b/cpp/cmake/modules/ConfigureCUDA.cmake index 395fffdb0b..98fe9e1f85 100644 --- a/cpp/cmake/modules/ConfigureCUDA.cmake +++ b/cpp/cmake/modules/ConfigureCUDA.cmake @@ -59,8 +59,5 @@ if(OpenMP_FOUND) endif() # Debug options -if(CMAKE_BUILD_TYPE MATCHES Debug) - message(VERBOSE "cuVS: Building with debugging flags") - list(APPEND CUVS_CUDA_FLAGS -Xcompiler=-rdynamic) - list(APPEND CUVS_CUDA_FLAGS -Xptxas --suppress-stack-size-warning) -endif() +list(APPEND CUVS_DEBUG_CUDA_FLAGS -G -Xcompiler=-rdynamic --maxrregcount=64) +list(APPEND CUVS_DEBUG_CUDA_FLAGS -Xptxas --suppress-stack-size-warning) From ef430f4d592854eca89859f07dcf1a4e8dde2224 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 23 Mar 2026 13:51:35 +0000 Subject: [PATCH 14/20] Removing irrelevant changes --- cpp/src/cluster/detail/kmeans_mg.cuh | 11 +---------- 1 file changed, 1 insertion(+), 10 deletions(-) diff --git a/cpp/src/cluster/detail/kmeans_mg.cuh b/cpp/src/cluster/detail/kmeans_mg.cuh index 8577b1d365..8e91c91d7f 100644 --- a/cpp/src/cluster/detail/kmeans_mg.cuh +++ b/cpp/src/cluster/detail/kmeans_mg.cuh @@ -159,7 +159,6 @@ void initKMeansPlusPlus(const raft::resources& handle, // Choose rp on rank 0 and broadcast to all ranks to guarantee agreement int rp = 0; if (my_rank == KMEANS_COMM_ROOT) { - //std::mt19937 gen(params.rng_state.seed); std::uniform_int_distribution<> dis(0, n_rank - 1); rp = dis(gen_64); } @@ -184,7 +183,6 @@ void initKMeansPlusPlus(const raft::resources& handle, // 1.2 - Rank r' samples a point uniformly at random from the local dataset // X which will be used as the initial centroid for kmeans++ if (my_rank == rp) { - //std::mt19937 gen(params.rng_state.seed + 31415926); std::uniform_int_distribution<> dis(0, n_samples - 1); int cIdx = dis(gen_64); @@ -283,7 +281,6 @@ void initKMeansPlusPlus(const raft::resources& handle, niter); // <<<< Step-3 >>> : for O( log(psi) ) times do - auto nvtx_range = nvtx3::start_range("step 3"); for (int iter = 0; iter < niter; ++iter) { CUVS_LOG_KMEANS(handle, "@Rank-%d:KMeans|| - Iteration %d: # potential centroids sampled - " @@ -321,7 +318,7 @@ void initKMeansPlusPlus(const raft::resources& handle, // potentialCentroids uint64_t gpu_seed; gpu_seed = gen_64(); - raft::random::RngState rng(gpu_seed, raft::random::GeneratorType::GenPhilox); + raft::random::RngState rng(gpu_seed, params.rng_state.rng_type); raft::random::uniform( handle, rng, uniformRands.data_handle(), uniformRands.extent(0), (DataT)0, (DataT)1); cuvs::cluster::kmeans::SamplingOp select_op(psi, @@ -380,7 +377,6 @@ void initKMeansPlusPlus(const raft::resources& handle, raft::make_device_matrix_view(centroidsBuf.data(), tot_centroids, n_features); /// <<<< End of Step-5 >>> } /// <<<< Step-6 >>> - nvtx3::end_range(nvtx_range); CUVS_LOG_KMEANS(handle, "@Rank-%d:KMeans||: # potential centroids sampled - %d\n", @@ -411,19 +407,15 @@ void initKMeansPlusPlus(const raft::resources& handle, // seed they should generate the same potentialCentroids auto const_centroids = raft::make_device_matrix_view( potentialCentroids.data_handle(), potentialCentroids.extent(0), potentialCentroids.extent(1)); - auto nvtx_range = nvtx3::start_range("init_plus_plus"); auto params_copy = params; params_copy.rng_state.seed = gen_64(); cuvs::cluster::kmeans::init_plus_plus( handle, params_copy, const_centroids, centroidsRawData, workspace); - nvtx3::end_range(nvtx_range); auto inertia = raft::make_host_scalar(0); auto n_iter = raft::make_host_scalar(0); auto weight_view = raft::make_device_vector_view(weight.data_handle(), weight.extent(0)); - //cuvs::cluster::kmeans::params params_copy = params; - //params_copy.rng_state = default_params.rng_state; // Update the seed one more time params_copy.rng_state.seed = gen_64(); cuvs::cluster::kmeans::fit_main(handle, @@ -448,7 +440,6 @@ void initKMeansPlusPlus(const raft::resources& handle, // generate `n_random_clusters` centroids cuvs::cluster::kmeans::params rand_params = params; - //rand_params.rng_state = default_params.rng_state; rand_params.rng_state.seed = gen_64(); rand_params.init = cuvs::cluster::kmeans::params::InitMethod::Random; rand_params.n_clusters = n_random_clusters; From a1de0477dd9118800a00eb7eaef4c95495db5b00 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 23 Mar 2026 14:04:07 +0000 Subject: [PATCH 15/20] Removing irrelevant changes 2 --- cpp/src/cluster/detail/kmeans.cuh | 11 ++--------- 1 file changed, 2 insertions(+), 9 deletions(-) diff --git a/cpp/src/cluster/detail/kmeans.cuh b/cpp/src/cluster/detail/kmeans.cuh index a1378d7533..58fa12fefd 100644 --- a/cpp/src/cluster/detail/kmeans.cuh +++ b/cpp/src/cluster/detail/kmeans.cuh @@ -52,7 +52,6 @@ #include #include - namespace cuvs::cluster::kmeans::detail { static const std::string CUVS_NAME = "cuvs"; @@ -145,8 +144,8 @@ void kmeansPlusPlus(raft::resources const& handle, } std::mt19937_64 gen_64(params.rng_state.seed); - //raft::random::RngState rng(params.rng_state.seed, raft::random::GeneratorType::GenPhilox); - //std::mt19937 gen(params.rng_state.seed); + uint64_t gpu_seed = gen_64(); + raft::random::RngState rng(gpu_seed, params.rng_state.type); std::uniform_int_distribution<> dis(0, n_samples - 1); // <<< Step-1 >>>: C <-- sample a point uniformly at random from X @@ -185,15 +184,9 @@ void kmeansPlusPlus(raft::resources const& handle, // distance to the nearest existing cluster // printf("size of index = %lu, size of weights = %lu, sizeof types (%ld, %ld)\n", // uint64_t(indices_view.extent(0)), uint64_t(const_weights_view.extent(0)), sizeof(DataT), sizeof(IndexT)); - uint64_t gpu_seed = gen_64(); - raft::random::RngState rng(gpu_seed, params.rng_state.type); raft::random::discrete(handle, rng, indices_view, const_weights_view); raft::matrix::gather(handle, const_X_view, const_indices_view, candidates_view); - // Create a CPU vector and copy indices from GPU to CPU - std::vector h_weights(n_samples); - raft::copy(h_weights.data(), minClusterDistance.data_handle(), n_samples, stream); - raft::resource::sync_stream(handle, stream); // Calculate pairwise distance between X and the centroid candidates // Output - pwd [n_trials x n_samples] auto pwd = distBuffer.view(); From 0bbeb30daf454624236f71c7f2bd08ee6f3ca54d Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 23 Mar 2026 14:05:13 +0000 Subject: [PATCH 16/20] Removing irrelevant changes 3 --- cpp/src/cluster/detail/kmeans.cuh | 2 -- 1 file changed, 2 deletions(-) diff --git a/cpp/src/cluster/detail/kmeans.cuh b/cpp/src/cluster/detail/kmeans.cuh index 58fa12fefd..9472fafc23 100644 --- a/cpp/src/cluster/detail/kmeans.cuh +++ b/cpp/src/cluster/detail/kmeans.cuh @@ -182,8 +182,6 @@ void kmeansPlusPlus(raft::resources const& handle, // <<< Step-3 >>> : Sample x in X with probability p_x = d^2(x, C) / phi_X (C) // Choose 'n_trials' centroid candidates from X with probability proportional to the squared // distance to the nearest existing cluster - // printf("size of index = %lu, size of weights = %lu, sizeof types (%ld, %ld)\n", - // uint64_t(indices_view.extent(0)), uint64_t(const_weights_view.extent(0)), sizeof(DataT), sizeof(IndexT)); raft::random::discrete(handle, rng, indices_view, const_weights_view); raft::matrix::gather(handle, const_X_view, const_indices_view, candidates_view); From 4d01c2f3cdf1a985980779ef0f6adf2f759b2779 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 23 Mar 2026 14:05:54 +0000 Subject: [PATCH 17/20] Removing irrelevant changes 4 --- cpp/src/cluster/detail/kmeans.cuh | 1 + 1 file changed, 1 insertion(+) diff --git a/cpp/src/cluster/detail/kmeans.cuh b/cpp/src/cluster/detail/kmeans.cuh index 9472fafc23..5408f7a863 100644 --- a/cpp/src/cluster/detail/kmeans.cuh +++ b/cpp/src/cluster/detail/kmeans.cuh @@ -182,6 +182,7 @@ void kmeansPlusPlus(raft::resources const& handle, // <<< Step-3 >>> : Sample x in X with probability p_x = d^2(x, C) / phi_X (C) // Choose 'n_trials' centroid candidates from X with probability proportional to the squared // distance to the nearest existing cluster + raft::random::discrete(handle, rng, indices_view, const_weights_view); raft::matrix::gather(handle, const_X_view, const_indices_view, candidates_view); From 7fcaf121026ec39ac5e51615fa39247563418cc8 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Mon, 23 Mar 2026 14:31:06 +0000 Subject: [PATCH 18/20] Correcting the name of the rng type --- cpp/src/cluster/detail/kmeans_mg.cuh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cpp/src/cluster/detail/kmeans_mg.cuh b/cpp/src/cluster/detail/kmeans_mg.cuh index 8e91c91d7f..c62f22a35b 100644 --- a/cpp/src/cluster/detail/kmeans_mg.cuh +++ b/cpp/src/cluster/detail/kmeans_mg.cuh @@ -318,7 +318,7 @@ void initKMeansPlusPlus(const raft::resources& handle, // potentialCentroids uint64_t gpu_seed; gpu_seed = gen_64(); - raft::random::RngState rng(gpu_seed, params.rng_state.rng_type); + raft::random::RngState rng(gpu_seed, params.rng_state.type); raft::random::uniform( handle, rng, uniformRands.data_handle(), uniformRands.extent(0), (DataT)0, (DataT)1); cuvs::cluster::kmeans::SamplingOp select_op(psi, From 8c436d435ddfa3c0fde8e1959d33bb2c7fe69cfc Mon Sep 17 00:00:00 2001 From: Vinay D Date: Wed, 15 Apr 2026 10:54:56 +0100 Subject: [PATCH 19/20] Removing debuggin files --- build_repro.sh | 25 ---- raft_discrete.cu | 134 -------------------- repro_kmeans_mg.cu | 295 --------------------------------------------- 3 files changed, 454 deletions(-) delete mode 100644 build_repro.sh delete mode 100644 raft_discrete.cu delete mode 100644 repro_kmeans_mg.cu diff --git a/build_repro.sh b/build_repro.sh deleted file mode 100644 index 76cec02cc6..0000000000 --- a/build_repro.sh +++ /dev/null @@ -1,25 +0,0 @@ -#!/bin/bash - -nvcc -g -std=c++17 --extended-lambda --expt-relaxed-constexpr \ - -Xcompiler -Wno-deprecated-declarations -Xcompiler -fopenmp \ - -o repro_kmeans_mg ../repro_kmeans_mg.cu \ - -I$CONDA_PREFIX/include \ - -I$CONDA_PREFIX/include/rapids \ - -I$CONDA_PREFIX/include/rapids/libcudacxx \ - -L$CONDA_PREFIX/lib \ - -lcuvs -lrmm -lnccl -lucp -lucs -lucxx -lgomp \ - -DRAFT_SYSTEM_LITTLE_ENDIAN=1 -DLIBCUDACXX_ENABLE_EXPERIMENTAL_MEMORY_RESOURCE \ - -Xlinker -rpath=$CONDA_PREFIX/lib \ - -gencode arch=compute_100,code=sm_100 - -# nvcc -std=c++17 --extended-lambda --expt-relaxed-constexpr \ -# -Xcompiler -Wno-deprecated-declarations -Xcompiler -fopenmp \ -# -o raft_discrete ../raft_discrete.cu \ -# -I$CONDA_PREFIX/include \ -# -I$CONDA_PREFIX/include/rapids \ -# -I$CONDA_PREFIX/include/rapids/libcudacxx \ -# -L$CONDA_PREFIX/lib \ -# -lcuvs -lrmm -lnccl -lucp -lucs -lucxx -lgomp \ -# -DRAFT_SYSTEM_LITTLE_ENDIAN=1 -DLIBCUDACXX_ENABLE_EXPERIMENTAL_MEMORY_RESOURCE \ -# -Xlinker -rpath=$CONDA_PREFIX/lib \ -# -gencode arch=compute_100,code=sm_100 diff --git a/raft_discrete.cu b/raft_discrete.cu deleted file mode 100644 index 470eabe4b5..0000000000 --- a/raft_discrete.cu +++ /dev/null @@ -1,134 +0,0 @@ -/* - * Demonstration of raft::random::discrete API usage - * - * The discrete API samples indices from a discrete probability distribution - * defined by weights. Each index i is sampled with probability proportional - * to weights[i]. - * - * Compile with: - * nvcc -o raft_discrete raft_discrete.cu -I --expt-extended-lambda - */ - -#include -#include -#include -#include - -#include - -#include -#include -#include -#include -#include -#include -#include - -int main() -{ - // Create RAFT handle (manages CUDA resources) - raft::handle_t handle; - auto stream = raft::resource::get_cuda_stream(handle); - - // Read weights from file (one weight per line) - std::vector h_weights; - const char* weights_filename_env = std::getenv("WEIGHTS_FILE"); - std::string weights_filename = weights_filename_env ? weights_filename_env : "weight_tracker_phi_00001.txt"; - std::ifstream weights_file(weights_filename); - if (!weights_file.is_open()) { - std::cerr << "Error: Could not open " << weights_filename << std::endl; - return 1; - } - std::string line; - while (std::getline(weights_file, line)) { - if (!line.empty()) { - h_weights.push_back(std::stof(line)); - } - } - - - // Number of categories (weights) - const int n_categories = static_cast(h_weights.size()); - // Number of samples to draw - const int n_samples = 9; - rmm::device_uvector d_weights(n_categories, stream); - raft::copy(d_weights.data(), h_weights.data(), n_categories, stream); - - // Create output array for sampled indices - rmm::device_uvector d_indices(n_samples, stream); - - // Create random number generator with a unique seed each run - auto seed = static_cast(std::chrono::high_resolution_clock::now().time_since_epoch().count()) ^ - static_cast(std::random_device{}()); - const char* use_pc_env = std::getenv("USE_PC"); - bool use_pc = use_pc_env && std::string(use_pc_env) == "1"; - raft::random::RngState rng(seed, use_pc ? raft::random::GenPC : raft::random::GenPhilox); - - // Create mdspan views for the API - auto indices_view = raft::make_device_vector_view(d_indices.data(), n_samples); - auto weights_view = - raft::make_device_vector_view(d_weights.data(), n_categories); - - // Sample indices according to the weight distribution - // Each index i will be sampled with probability weights[i] / sum(weights) - - // Count occurrences of each index - std::vector counts(n_categories, 0); - std::vector h_indices(n_samples); - - // Number of iterations to run - const int n_iterations = 1000; - - for (int iter = 0; iter < n_iterations; ++iter) { - raft::random::discrete(handle, rng, indices_view, weights_view); - - // Copy results back to host - raft::copy(h_indices.data(), d_indices.data(), n_samples, stream); - raft::resource::sync_stream(handle, stream); - - // Update counts - for (int i = 0; i < n_samples; ++i) { - counts[h_indices[i]]++; - } - } - - // Print results - // std::cout << "Weights: "; - // for (int i = 0; i < n_categories && i < 10; ++i) { - // std::cout << h_weights[i] << " "; - // } - // std::cout << std::endl; - - //std::cout << "Total samples: " << n_iterations * n_samples << std::endl; - - // Normalize counts array - double counts_sum = 0.0; - for (int i = 0; i < n_categories; ++i) { - counts_sum += counts[i]; - } - std::vector normalized_counts(n_categories); - for (int i = 0; i < n_categories; ++i) { - normalized_counts[i] = static_cast(counts[i]) / counts_sum; - } - - // Normalize weights array - double weights_sum = 0.0; - for (int i = 0; i < n_categories; ++i) { - weights_sum += h_weights[i]; - } - std::vector normalized_weights(n_categories); - for (int i = 0; i < n_categories; ++i) { - normalized_weights[i] = static_cast(h_weights[i]) / weights_sum; - } - - // Calculate sum of squared differences - double sum_sq_diff = 0.0; - for (int i = 0; i < n_categories; ++i) { - double diff = normalized_counts[i] - normalized_weights[i]; - sum_sq_diff += diff * diff; - } - - std::cout << weights_filename << " " << "Sum of squared differences: " << sum_sq_diff << std::endl; - - return 0; -} diff --git a/repro_kmeans_mg.cu b/repro_kmeans_mg.cu deleted file mode 100644 index 14815e4695..0000000000 --- a/repro_kmeans_mg.cu +++ /dev/null @@ -1,295 +0,0 @@ -/* - * Minimal multi-GPU KMeans example using cuVS. - * - * nvcc -std=c++17 --extended-lambda --expt-relaxed-constexpr \ - * -Xcompiler -Wno-deprecated-declarations -Xcompiler -fopenmp \ - * -o repro_kmeans_mg repro_kmeans_mg.cu \ - * -I$CONDA_PREFIX/include \ - * -I$CONDA_PREFIX/include/rapids \ - * -I$CONDA_PREFIX/include/rapids/libcudacxx \ - * -L$CONDA_PREFIX/lib \ - * -lcuvs -lrmm -lnccl -lucp -lucs -lucxx -lgomp \ - * -DRAFT_SYSTEM_LITTLE_ENDIAN=1 -DLIBCUDACXX_ENABLE_EXPERIMENTAL_MEMORY_RESOURCE \ - * -Xlinker -rpath=$CONDA_PREFIX/lib \ - * -gencode arch=compute_XX,code=sm_XX - */ - -#include - -#include -#include -#include -#include -#include -#include - -#include - -#include -#include - -#include -#include -#include -#include -#include -#include -#include -#include -#include - -// ----- helpers -------------------------------------------------------------- -#define NCCLCHECK(cmd) \ - do { \ - ncclResult_t res = cmd; \ - if (res != ncclSuccess) { \ - std::fprintf( \ - stderr, "NCCL error %s:%d '%s'\n", __FILE__, __LINE__, ncclGetErrorString(res)); \ - std::exit(EXIT_FAILURE); \ - } \ - } while (0) - -#define CUDACHECK(cmd) \ - do { \ - cudaError_t e = cmd; \ - if (e != cudaSuccess) { \ - std::fprintf( \ - stderr, "CUDA error %s:%d '%s'\n", __FILE__, __LINE__, cudaGetErrorString(e)); \ - std::exit(EXIT_FAILURE); \ - } \ - } while (0) - -constexpr int64_t N_CLUSTERS = 1000; -constexpr int N_BENCHMARK_ITERS = 1; - - -int main(int argc, char* argv[]) -{ - // Parse command line arguments - int64_t n_samples_total = 10000000; // Default value - int64_t n_features = 256; // Default value - - for (int i = 1; i < argc; i++) { - std::string arg = argv[i]; - if (arg == "--samples" || arg == "-s") { - if (i + 1 < argc) { - n_samples_total = std::stoll(argv[++i]); - } else { - std::cerr << "--samples requires a value" << std::endl; - return 1; - } - } else if (arg == "--features" || arg == "-f") { - if (i + 1 < argc) { - n_features = std::stoll(argv[++i]); - } else { - std::cerr << "--features requires a value" << std::endl; - return 1; - } - } else if (arg == "--help" || arg == "-h") { - std::cout << "Usage: " << argv[0] << " [options]" << std::endl; - std::cout << "Options:" << std::endl; - std::cout << " -s, --samples N Set total number of samples (default: 10000000)" << std::endl; - std::cout << " -f, --features N Set number of features (default: 256)" << std::endl; - std::cout << " -h, --help Show this help message" << std::endl; - return 0; - } - } - - // Detect available GPUs - int n_devices = 0; - CUDACHECK(cudaGetDeviceCount(&n_devices)); - if (n_devices < 1) { - std::cerr << "No CUDA devices found." << std::endl; - return 1; - } - std::cout << "Found " << n_devices << " GPU(s), running multi-GPU KMeans." - << std::endl; - - // Create NCCL communicators (one per device) - std::vector dev_list(n_devices); - for (int i = 0; i < n_devices; ++i) dev_list[i] = i; - - std::vector nccl_comms(n_devices); - NCCLCHECK(ncclCommInitAll(nccl_comms.data(), n_devices, dev_list.data())); - - // ------------------------------------------------------------------ - // Generate shared cluster centers on the host (BEFORE the parallel region) - // so that every rank's data is clustered around the SAME centers. - // Use seed 42 to match `make_blobs(random_state=42, centers=1000)`. - // ------------------------------------------------------------------ - std::vector h_centers(N_CLUSTERS * n_features); - std::mt19937 gen(42); - { - std::uniform_real_distribution dist(-10.0f, 10.0f); - for (auto& v : h_centers) v = dist(gen); - } - - // Pre-generate per-rank seeds (mimics Dask's per-partition random seeds) - std::vector rank_seeds(n_devices); - for (auto& s : rank_seeds) s = gen(); - - // Launch one OpenMP thread per GPU (OPG model) -#pragma omp parallel for num_threads(n_devices) - for (int rank = 0; rank < n_devices; ++rank) { - // 1. Bind this thread to the correct GPU - CUDACHECK(cudaSetDevice(rank)); - - // 2. Create a raft handle and inject the NCCL communicator. - // This is the *only* extra step compared to single-GPU usage. - raft::resources handle; - raft::comms::build_comms_nccl_only(&handle, nccl_comms[rank], n_devices, rank); - - auto stream = raft::resource::get_cuda_stream(handle); - - // 3. Compute this rank's shard of the dataset. - // Each rank gets a contiguous, non-overlapping slice of the rows. - int64_t n_samples_per_rank = n_samples_total / n_devices; - int64_t leftover = n_samples_total % n_devices; - int64_t n_local = n_samples_per_rank + (rank < leftover ? 1 : 0); - - // Copy the shared cluster centers to this GPU - rmm::device_uvector d_blob_centers(N_CLUSTERS * n_features, stream); - CUDACHECK(cudaMemcpyAsync(d_blob_centers.data(), - h_centers.data(), - h_centers.size() * sizeof(float), - cudaMemcpyHostToDevice, - stream)); - - // Generate synthetic clustered data directly on device. - // Each rank uses a different seed so the data points are distinct, - // but all ranks share the same cluster centers. - rmm::device_uvector d_X(n_local * n_features, stream); - rmm::device_uvector d_labels_true(n_local, stream); - - raft::random::make_blobs(d_X.data(), - d_labels_true.data(), - n_local, - n_features, - N_CLUSTERS, - stream, - /*row_major=*/true, - /*centers=*/d_blob_centers.data(), - /*cluster_std=*/nullptr, - /*cluster_std_scalar=*/0.5f, - /*shuffle=*/true, - /*center_box_min=*/-10.0f, - /*center_box_max=*/10.0f, - /*seed=*/rank_seeds[rank]); - - // 4. Prepare output buffers - rmm::device_uvector d_centroids(N_CLUSTERS * n_features, stream); - rmm::device_uvector d_labels(n_local, stream); - - // Ensure data generation is complete before timing - raft::resource::sync_stream(handle); - - if (rank == 0) { - std::cout << "Data generation complete. " - << n_samples_total << " total samples, " - << n_local << " per rank, " - << n_features << " features, " - << N_CLUSTERS << " clusters." << std::endl; - } - - // 5. Benchmark loop — run both RNG types to compare - struct RngConfig { - raft::random::GeneratorType type; - const char* name; - }; - RngConfig rng_configs[] = { - {raft::random::GeneratorType::GenPC, "GenPC"}, - {raft::random::GeneratorType::GenPhilox, "GenPhilox"}, - }; - - std::mt19937_64 gen_64(std::chrono::high_resolution_clock::now().time_since_epoch().count()); - uint64_t last_seed = gen_64(); - for (const auto& rng_cfg : rng_configs) { - // Create NVTX range for this RNG configuration - nvtx3::scoped_range range(rng_cfg.name); - - if (rank == 0) { - std::printf("\n=== RNG: %s ===\n", rng_cfg.name); - } - - int64_t total_kmeans_iters = 0; - for (int iter = 0; iter < N_BENCHMARK_ITERS; ++iter) { - float inertia = 0; - float pred_inertia = 0; - int64_t n_iter = 0; - - cuvs::cluster::kmeans::params params; - params.n_clusters = static_cast(N_CLUSTERS); - params.max_iter = 300; - params.tol = 1e-4; - params.rng_state.seed = last_seed; - params.rng_state.type = rng_cfg.type; - params.oversampling_factor = 2.0; - params.n_init = 1; - - auto X_view = raft::make_device_matrix_view( - d_X.data(), n_local, n_features); - auto centroids_view = raft::make_device_matrix_view( - d_centroids.data(), N_CLUSTERS, n_features); - auto centroids_const_view = raft::make_device_matrix_view( - d_centroids.data(), N_CLUSTERS, n_features); - auto labels_view = raft::make_device_vector_view( - d_labels.data(), n_local); - - auto t0 = std::chrono::high_resolution_clock::now(); - - { - nvtx3::scoped_range fit_range("kmeans_fit"); - cuvs::cluster::kmeans::fit(handle, - params, - X_view, - std::nullopt, - centroids_view, - raft::make_host_scalar_view(&inertia), - raft::make_host_scalar_view(&n_iter)); - raft::resource::sync_stream(handle); - } - - auto t1 = std::chrono::high_resolution_clock::now(); - - { - nvtx3::scoped_range predict_range("kmeans_predict"); - cuvs::cluster::kmeans::predict(handle, - params, - X_view, - std::nullopt, - centroids_const_view, - labels_view, - true, - raft::make_host_scalar_view(&pred_inertia)); - raft::resource::sync_stream(handle); - } - - auto t2 = std::chrono::high_resolution_clock::now(); - - double fit_s = std::chrono::duration(t1 - t0).count(); - double total_s = std::chrono::duration(t2 - t0).count(); - - total_kmeans_iters += n_iter; - - if (rank == 0) { - std::printf("%d) Time taken by fit %.2fs and predict %.2fs (iters: %ld)\n", - iter, fit_s, total_s, static_cast(n_iter)); - } - } - - if (rank == 0) { - double avg_iters = static_cast(total_kmeans_iters) / N_BENCHMARK_ITERS; - std::printf("Last seed used: %lu\n", last_seed); - std::printf("Average KMeans iterations per run (%s): %.1f\n", rng_cfg.name, avg_iters); - } - } - } - - // Clean up NCCL - for (int i = 0; i < n_devices; ++i) { - ncclCommDestroy(nccl_comms[i]); - } - - std::cout << "Done." << std::endl; - return 0; -} From ea165665573300f3863d65f12ed0e20cf27bf080 Mon Sep 17 00:00:00 2001 From: Vinay D Date: Wed, 15 Apr 2026 10:57:50 +0100 Subject: [PATCH 20/20] Removing debug print --- cpp/src/cluster/detail/kmeans_mg.cuh | 1 - 1 file changed, 1 deletion(-) diff --git a/cpp/src/cluster/detail/kmeans_mg.cuh b/cpp/src/cluster/detail/kmeans_mg.cuh index c62f22a35b..7138bc0851 100644 --- a/cpp/src/cluster/detail/kmeans_mg.cuh +++ b/cpp/src/cluster/detail/kmeans_mg.cuh @@ -518,7 +518,6 @@ void fit(const raft::resources& handle, const int n_ranks = comm.get_size(); std::mt19937_64 gen_64(params.rng_state.seed + (uint64_t(my_rank) << 32)); - printf("I am rank %d of %d total ranks\n", my_rank, n_ranks); auto weight = raft::make_device_vector(handle, n_samples); if (sample_weight) {