| Line | Branch | Exec | Source |
|---|---|---|---|
| 1 | /* | ||
| 2 | * Copyright (c) 2026 Tiger Data, Inc. | ||
| 3 | * Licensed under the PostgreSQL License. See LICENSE for details. | ||
| 4 | * | ||
| 5 | * parallel_build_worker.c - Parallel index build, per-participant bodies | ||
| 6 | * | ||
| 7 | * The leader runs these as participant 0 alongside the workers: | ||
| 8 | * Sampling: cooperative heap scan, each participant fills its own | ||
| 9 | * sample slot. | ||
| 10 | * K-means: leader seeds the initial centroids, then iterative | ||
| 11 | * assignment + accumulation over the per-participant samples, | ||
| 12 | * barrier-synchronized with a leader reduce. | ||
| 13 | * Subtrees: batched largest-first schedule; each participant clusters | ||
| 14 | * its batch's child into a ring slot the leader consumes. | ||
| 15 | * Refine (when subsampled): route the full table page-backed and | ||
| 16 | * accumulate per-leaf means into the shared tiled accumulator. | ||
| 17 | * Posting: cooperative heap scan, page-backed routing + RaBitQ encode | ||
| 18 | * into the cluster-keyed shared sort. | ||
| 19 | */ | ||
| 20 | |||
| 21 | #ifdef VS_STANDALONE | ||
| 22 | #include "standalone/parallel_ctx.h" /* ParallelWorkerNumber */ | ||
| 23 | #include "standalone/pg_compat.h" | ||
| 24 | #else | ||
| 25 | #include <postgres.h> | ||
| 26 | |||
| 27 | #include <access/parallel.h> | ||
| 28 | #include <catalog/index.h> | ||
| 29 | #include <miscadmin.h> | ||
| 30 | #include <storage/barrier.h> | ||
| 31 | #include <storage/bufmgr.h> | ||
| 32 | #include <storage/shm_toc.h> | ||
| 33 | #include <utils/rel.h> | ||
| 34 | #include <utils/wait_event.h> | ||
| 35 | #endif | ||
| 36 | |||
| 37 | #include <inttypes.h> | ||
| 38 | #include <math.h> | ||
| 39 | |||
| 40 | #include "algo/hkmeans.h" | ||
| 41 | #include "algo/kmeans_internal.h" | ||
| 42 | #include "algo/vecops.h" | ||
| 43 | #include "core/log.h" | ||
| 44 | #include "core/memory.h" | ||
| 45 | #include "index/build_progress.h" | ||
| 46 | #include "index/index_build.h" | ||
| 47 | #include "index/parallel_build.h" | ||
| 48 | #include "index/posting_build.h" | ||
| 49 | #include "quant/rabitq.h" | ||
| 50 | |||
| 51 | /* ---------------------------------------------------------------- | ||
| 52 | * Phase 1: Sampling callback | ||
| 53 | * | ||
| 54 | * Collects vectors into a per-worker sample slot (stride-subsampled, | ||
| 55 | * normalized for cosine). Initial centroids are seeded later by the | ||
| 56 | * leader from the pooled samples. | ||
| 57 | * ---------------------------------------------------------------- */ | ||
| 58 | |||
| 59 | void | ||
| 60 | 279842 | prism_sample_cb(void *state, ItemPointerData tid, const float *vec) | |
| 61 | { | ||
| 62 | 279842 | SampleCbState *sc = (SampleCbState *)state; | |
| 63 | |||
| 64 | 109480 | (void)tid; | |
| 65 | |||
| 66 | 279842 | sc->seen++; | |
| 67 | /* The sampling pass walks the whole heap; surface that in the progress | ||
| 68 | * view in coarse batches (the leader publishes exact counts at phase | ||
| 69 | * boundaries). */ | ||
| 70 |
2/2✓ Branch 0 taken 205 times.
✓ Branch 1 taken 279637 times.
|
279842 | if ((sc->seen & 1023) == 0) |
| 71 | 205 | prism_build_progress_incr_tuples(1024); | |
| 72 | |||
| 73 | /* Stride-based subsampling */ | ||
| 74 |
2/2✓ Branch 0 taken 22941 times.
✓ Branch 1 taken 256901 times.
|
279842 | if (sc->stride_counter > 0) |
| 75 | { | ||
| 76 | 22941 | sc->stride_counter--; | |
| 77 | 22941 | return; | |
| 78 | } | ||
| 79 | 256901 | sc->stride_counter = sc->stride - 1; | |
| 80 | |||
| 81 |
2/2✓ Branch 0 taken 88222 times.
✓ Branch 1 taken 168679 times.
|
256901 | if (sc->count >= sc->max_samples) |
| 82 | 21488 | return; | |
| 83 | |||
| 84 | 215608 | Dimension dim = sc->dim; | |
| 85 | 215608 | float *dest = sc->samples + (size_t)sc->count * dim; | |
| 86 | |||
| 87 | /* Normalize for cosine. Zero-norm vectors are dropped from the | ||
| 88 | * sample: they carry no direction, so under cosine they sit at the | ||
| 89 | * origin -- equidistant from every real point -- and k-means drags | ||
| 90 | * centroids toward them, warping the encode references (and so the | ||
| 91 | * distance estimates) for every vector in the affected subtree. | ||
| 92 | * They stay in the posting lists like any other row; they just | ||
| 93 | * don't get a vote on the clustering. The check runs on the source | ||
| 94 | * so a dropped vector costs no copy, and the scale writes straight | ||
| 95 | * into the sample slot, normalizing and copying in one pass. */ | ||
| 96 |
2/2✓ Branch 0 taken 5159 times.
✓ Branch 1 taken 210449 times.
|
215608 | if (sc->metric == DISTANCE_COSINE) |
| 97 | { | ||
| 98 | 5159 | float norm_sq = vs_l2_norm_squared(vec, dim); | |
| 99 | |||
| 100 |
2/2✓ Branch 0 taken 3358 times.
✓ Branch 1 taken 1801 times.
|
5159 | if (norm_sq == 0.0f) |
| 101 | 80 | return; | |
| 102 | 5078 | vec32_scale(vec, 1.0f / sqrtf(norm_sq), dest, dim); | |
| 103 | } | ||
| 104 | else | ||
| 105 | 210449 | memcpy(dest, vec, dim * sizeof(float)); | |
| 106 | |||
| 107 | 215527 | sc->count++; | |
| 108 | } | ||
| 109 | |||
| 110 | /* ---------------------------------------------------------------- | ||
| 111 | * Phase 2: K-means assignment — thin wrappers around shared kernel | ||
| 112 | * ---------------------------------------------------------------- */ | ||
| 113 | |||
| 114 | void | ||
| 115 | 3333 | prism_km_assign_and_accumulate( | |
| 116 | const float *samples, | ||
| 117 | uint32_t nsamples, | ||
| 118 | const float *centroids, | ||
| 119 | const float *norms_c, | ||
| 120 | uint32_t nlist, | ||
| 121 | Dimension dim, | ||
| 122 | DistanceMetric metric, | ||
| 123 | float *out_sums, | ||
| 124 | uint32_t *out_cnts, | ||
| 125 | float *out_cost) | ||
| 126 | { | ||
| 127 | 3333 | memset(out_sums, 0, (size_t)nlist * dim * sizeof(float)); | |
| 128 | 3333 | memset(out_cnts, 0, nlist * sizeof(uint32_t)); | |
| 129 | 3333 | *out_cost = 0.0f; | |
| 130 | 3333 | kmeans_assign_accumulate( | |
| 131 | samples, | ||
| 132 | NULL, | ||
| 133 | 0, | ||
| 134 | nsamples, | ||
| 135 | centroids, | ||
| 136 | norms_c, | ||
| 137 | nlist, | ||
| 138 | dim, | ||
| 139 | metric, | ||
| 140 | NULL, | ||
| 141 | 0, | ||
| 142 | out_sums, | ||
| 143 | out_cnts, | ||
| 144 | out_cost); | ||
| 145 | 3333 | } | |
| 146 | |||
| 147 | /* | ||
| 148 | * Build ONE root-child's subtree into `slot` (a DSM ring slot indexed by | ||
| 149 | * participant, not by child, so the batched streaming build keeps only | ||
| 150 | * nparticipants subtrees resident). The subtree is clustered from this child's | ||
| 151 | * root-assigned samples, indexed in place (no contiguous per-child copy). | ||
| 152 | */ | ||
| 153 | void | ||
| 154 | 239 | prism_build_child_subtree( | |
| 155 | uint32_t child, | ||
| 156 | uint32_t child_count, | ||
| 157 | int nparticipants, | ||
| 158 | PrismDsmSamples *dsm_samples, | ||
| 159 | PrismDsmRootAssign *dsm_ra, | ||
| 160 | const float *root_cents, | ||
| 161 | uint32_t nlist, | ||
| 162 | uint32_t fan_out, | ||
| 163 | Dimension dim, | ||
| 164 | DistanceMetric metric, | ||
| 165 | uint32_t km_max_iterations, | ||
| 166 | char *slot, | ||
| 167 | uint64_t slot_size) | ||
| 168 | { | ||
| 169 | 239 | uint32_t nlist_c = (nlist + fan_out - 1) / fan_out; | |
| 170 | /* The caller's histogram already counted this child's samples (one pass | ||
| 171 | * over the assignments instead of one per child). */ | ||
| 172 | 239 | uint32_t cc = child_count; | |
| 173 | |||
| 174 | 239 | KMeansOptions opts = VS_KMEANS_OPTIONS_DEFAULT; | |
| 175 | 239 | opts.max_iterations = km_max_iterations; | |
| 176 | 239 | opts.algorithm = KMEANS_ALGO_LLOYD; | |
| 177 | 239 | opts.initial_centroids = NULL; | |
| 178 | |||
| 179 | 239 | const size_t vec_nbytes = (size_t)dim * sizeof(float); | |
| 180 | |||
| 181 | 239 | HKMeansResult *sub = NULL; | |
| 182 |
2/2✓ Branch 0 taken 3 times.
✓ Branch 1 taken 236 times.
|
239 | if (cc == 0) |
| 183 | { | ||
| 184 | 3 | float *seed = vs_alloc(vec_nbytes); | |
| 185 | 3 | memcpy(seed, root_cents + (size_t)child * dim, vec_nbytes); | |
| 186 | 3 | sub = vs_hkmeans_f32( | |
| 187 | seed, 1, NULL, dim, nlist_c, fan_out, metric, &opts); | ||
| 188 | 3 | vs_free(seed); | |
| 189 | } | ||
| 190 | else | ||
| 191 | { | ||
| 192 | 236 | const float *vbase = prism_dsm_worker_samples(dsm_samples, 0); | |
| 193 | 236 | uint32_t mpw = dsm_samples->max_per_worker; | |
| 194 | 236 | uint32_t *idx = vs_alloc((size_t)cc * sizeof(uint32_t)); | |
| 195 | 236 | uint32_t g = 0; | |
| 196 |
3/3✓ Branch 0 taken 364 times.
✓ Branch 1 taken 455 times.
✓ Branch 2 taken 110 times.
|
929 | for (int t = 0; t < nparticipants; t++) |
| 197 | { | ||
| 198 | 693 | const uint32_t *ra = prism_dsm_root_assignments(dsm_ra, t); | |
| 199 | 693 | uint32_t n = prism_dsm_sample_counts(dsm_samples)[t]; | |
| 200 |
2/2✓ Branch 0 taken 1229412 times.
✓ Branch 1 taken 693 times.
|
1230105 | for (uint32_t i = 0; i < n; i++) |
| 201 |
2/2✓ Branch 0 taken 140526 times.
✓ Branch 1 taken 1088886 times.
|
1229412 | if (ra[i] == child) |
| 202 | 140526 | idx[g++] = (uint32_t)t * mpw + i; | |
| 203 | } | ||
| 204 | 236 | sub = vs_hkmeans_f32( | |
| 205 | vbase, cc, idx, dim, nlist_c, fan_out, metric, &opts); | ||
| 206 | 236 | vs_free(idx); | |
| 207 | } | ||
| 208 | |||
| 209 |
1/2✓ Branch 0 taken 239 times.
✗ Branch 1 not taken.
|
239 | if (sub != NULL) |
| 210 | { | ||
| 211 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 239 times.
|
239 | if ((uint64_t)sub->total_size > slot_size) |
| 212 | ✗ | vs_error( | |
| 213 | "prism: subtree blob %u exceeds slot " | ||
| 214 | "(%u > %" PRIu64 ")", | ||
| 215 | child, | ||
| 216 | sub->total_size, | ||
| 217 | slot_size); | ||
| 218 | 239 | memcpy(slot, sub, sub->total_size); | |
| 219 | 239 | vs_free(sub); | |
| 220 | } | ||
| 221 | 239 | } | |
| 222 | |||
| 223 | void | ||
| 224 | 191 | prism_pbuild_count_children( | |
| 225 | PrismDsmSamples *dsm_samples, | ||
| 226 | PrismDsmRootAssign *dsm_ra, | ||
| 227 | int nparticipants, | ||
| 228 | uint32_t km_k, | ||
| 229 | uint32_t *out_counts) | ||
| 230 | { | ||
| 231 | 191 | memset(out_counts, 0, (size_t)km_k * sizeof(uint32_t)); | |
| 232 |
2/2✓ Branch 0 taken 575 times.
✓ Branch 1 taken 191 times.
|
766 | for (int t = 0; t < nparticipants; t++) |
| 233 | { | ||
| 234 | 575 | const uint32_t *ra = prism_dsm_root_assignments(dsm_ra, t); | |
| 235 | 575 | uint32_t n = prism_dsm_sample_counts(dsm_samples)[t]; | |
| 236 |
2/2✓ Branch 0 taken 589562 times.
✓ Branch 1 taken 575 times.
|
590137 | for (uint32_t i = 0; i < n; i++) |
| 237 | 589562 | out_counts[ra[i]]++; | |
| 238 | } | ||
| 239 | 191 | } | |
| 240 | |||
| 241 | void | ||
| 242 | 141 | prism_pbuild_stream_subtrees( | |
| 243 | int participant_id, | ||
| 244 | int nparticipants, | ||
| 245 | PrismDsmSamples *dsm_samples, | ||
| 246 | PrismDsmRootAssign *dsm_ra, | ||
| 247 | const float *root_cents, | ||
| 248 | uint32_t km_k, | ||
| 249 | uint32_t nlist, | ||
| 250 | uint32_t fan_out, | ||
| 251 | Dimension dim, | ||
| 252 | DistanceMetric metric, | ||
| 253 | uint32_t km_max_iterations, | ||
| 254 | char *subtrees_base, | ||
| 255 | uint64_t slot_size, | ||
| 256 | Barrier *barrier, | ||
| 257 | PrismBatchCb batch_cb, | ||
| 258 | void *cb_arg, | ||
| 259 | uint32_t *out_child_order) | ||
| 260 | { | ||
| 261 | 141 | uint32_t np = (uint32_t)nparticipants; | |
| 262 | 141 | uint32_t nbatches = (km_k + np - 1) / np; | |
| 263 | 57 | char *slot = | |
| 264 | 141 | prism_dsm_child_subtree(subtrees_base, participant_id, slot_size); | |
| 265 | |||
| 266 | /* One histogram pass over the root assignments feeds every child's | ||
| 267 | * sample count (prism_build_child_subtree needs it, and counting per | ||
| 268 | * child would re-scan the assignments km_k times). */ | ||
| 269 | 141 | uint32_t *child_count = vs_alloc0((size_t)km_k * sizeof(uint32_t)); | |
| 270 | 141 | prism_pbuild_count_children( | |
| 271 | dsm_samples, dsm_ra, nparticipants, km_k, child_count); | ||
| 272 | |||
| 273 | /* Schedule children largest-first (LPT): every batch waits for its | ||
| 274 | * slowest subtree at two barriers, so the skewed children must land in | ||
| 275 | * the full batches, leaving the tail batch the small ones. Counts are | ||
| 276 | * identical for every participant (same shared assignments) and ties | ||
| 277 | * break on the child id, so all participants and the leader's blob | ||
| 278 | * replay derive the same order with no coordination. */ | ||
| 279 | 141 | uint32_t *order = vs_alloc((size_t)km_k * sizeof(uint32_t)); | |
| 280 |
3/3✓ Branch 0 taken 364 times.
✓ Branch 1 taken 424 times.
✓ Branch 2 taken 57 times.
|
845 | for (uint32_t i = 0; i < km_k; i++) |
| 281 | 704 | order[i] = i; | |
| 282 |
2/2✓ Branch 0 taken 563 times.
✓ Branch 1 taken 141 times.
|
704 | for (uint32_t i = 1; i < km_k; i++) |
| 283 | { | ||
| 284 | 563 | uint32_t id = order[i]; | |
| 285 | 563 | uint32_t j = i; | |
| 286 |
6/6✓ Branch 0 taken 1348 times.
✓ Branch 1 taken 209 times.
✓ Branch 2 taken 994 times.
✓ Branch 3 taken 354 times.
✓ Branch 4 taken 64 times.
✓ Branch 5 taken 134 times.
|
1557 | while (j > 0 && (child_count[order[j - 1]] < child_count[id] || |
| 287 |
2/2✓ Branch 0 taken 9 times.
✓ Branch 1 taken 211 times.
|
220 | (child_count[order[j - 1]] == child_count[id] && |
| 288 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 9 times.
|
9 | order[j - 1] > id))) |
| 289 | { | ||
| 290 | 994 | order[j] = order[j - 1]; | |
| 291 | 994 | j--; | |
| 292 | } | ||
| 293 | 563 | order[j] = id; | |
| 294 | } | ||
| 295 |
2/2✓ Branch 0 taken 49 times.
✓ Branch 1 taken 92 times.
|
141 | if (out_child_order != NULL) |
| 296 | 49 | memcpy(out_child_order, order, (size_t)km_k * sizeof(uint32_t)); | |
| 297 | |||
| 298 |
2/2✓ Branch 0 taken 285 times.
✓ Branch 1 taken 139 times.
|
424 | for (uint32_t b = 0; b < nbatches; b++) |
| 299 | { | ||
| 300 | 285 | uint32_t my_idx = b * np + (uint32_t)participant_id; | |
| 301 |
2/2✓ Branch 0 taken 239 times.
✓ Branch 1 taken 46 times.
|
285 | if (my_idx < km_k) |
| 302 | { | ||
| 303 | 239 | uint32_t my_child = order[my_idx]; | |
| 304 | 239 | prism_build_child_subtree( | |
| 305 | my_child, | ||
| 306 | 239 | child_count[my_child], | |
| 307 | nparticipants, | ||
| 308 | dsm_samples, | ||
| 309 | dsm_ra, | ||
| 310 | root_cents, | ||
| 311 | nlist, | ||
| 312 | fan_out, | ||
| 313 | dim, | ||
| 314 | metric, | ||
| 315 | km_max_iterations, | ||
| 316 | slot, | ||
| 317 | slot_size); | ||
| 318 | } | ||
| 319 | |||
| 320 | /* All participants have built this batch's subtrees into their slots. | ||
| 321 | */ | ||
| 322 | 285 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 323 | |||
| 324 | /* Leader consumes the batch before the slots are reused. */ | ||
| 325 |
3/4✓ Branch 0 taken 102 times.
✓ Branch 1 taken 181 times.
✓ Branch 2 taken 56 times.
✗ Branch 3 not taken.
|
283 | if (participant_id == 0 && batch_cb != NULL) |
| 326 | { | ||
| 327 | 102 | uint32_t base = b * np; | |
| 328 | 102 | uint32_t bs = km_k - base; | |
| 329 |
2/2✓ Branch 0 taken 26 times.
✓ Branch 1 taken 30 times.
|
102 | if (bs > np) |
| 330 | 26 | bs = np; | |
| 331 | 102 | batch_cb(cb_arg, order + base, bs, subtrees_base, slot_size); | |
| 332 | } | ||
| 333 | |||
| 334 | /* Leader done with the batch; slots free for the next batch. */ | ||
| 335 | 283 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 336 | } | ||
| 337 | |||
| 338 | 139 | vs_free(order); | |
| 339 | 139 | vs_free(child_count); | |
| 340 | 139 | } | |
| 341 | |||
| 342 | /* ---------------------------------------------------------------- | ||
| 343 | * Phases 1, 2, 2b: per-participant execution, shared by leader and workers | ||
| 344 | * | ||
| 345 | * The leader runs as participant 0 and calls these exactly as a worker does, | ||
| 346 | * rather than duplicating the bodies. The phase barriers live inside each | ||
| 347 | * function, so the leader and workers stay in lockstep by construction; the | ||
| 348 | * leader-only coordination (seeding the initial centroids, reducing the | ||
| 349 | * per-iteration accumulators) is gated on participant_id == 0 and runs at the | ||
| 350 | * barrier rendezvous while the other participants wait. | ||
| 351 | * ---------------------------------------------------------------- */ | ||
| 352 | |||
| 353 | void | ||
| 354 | 318 | prism_pbuild_exec_sampling( | |
| 355 | int participant_id, | ||
| 356 | Relation heap, | ||
| 357 | Relation index, | ||
| 358 | struct IndexInfo *index_info, | ||
| 359 | PrismBuildShared *shared, | ||
| 360 | PrismDsmSamples *dsm_samples, | ||
| 361 | Barrier *barrier) | ||
| 362 | { | ||
| 363 | 318 | Dimension dim = shared->dim; | |
| 364 | |||
| 365 | /* Stride to subsample ~max_samples_per_worker from this participant's | ||
| 366 | * share of the heap. */ | ||
| 367 | 318 | double est_rows = RelationGetNumberOfBlocks(heap) * | |
| 368 | 318 | (BLCKSZ / (double)(dim * sizeof(float) + 32)); | |
| 369 | 318 | double est_per = est_rows / shared->nparticipants; | |
| 370 | 318 | uint32_t stride = 1; | |
| 371 |
2/2✓ Branch 0 taken 99 times.
✓ Branch 1 taken 219 times.
|
318 | if (est_per > shared->max_samples_per_worker) |
| 372 | 99 | stride = (uint32_t)(est_per / shared->max_samples_per_worker); | |
| 373 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 188 times.
|
273 | if (stride < 1) |
| 374 | 130 | stride = 1; | |
| 375 | |||
| 376 | 318 | SampleCbState sc = { | |
| 377 | 318 | .samples = prism_dsm_worker_samples(dsm_samples, participant_id), | |
| 378 | .count = 0, | ||
| 379 | .seen = 0, | ||
| 380 | 188 | .max_samples = shared->max_samples_per_worker, | |
| 381 | .stride = stride, | ||
| 382 | .stride_counter = 0, | ||
| 383 | .dim = dim, | ||
| 384 | 318 | .metric = shared->metric, | |
| 385 | }; | ||
| 386 | |||
| 387 | /* progress: only the leader (participant 0) drives the PG progress view. | ||
| 388 | * The scan count is discarded: the posting pass is the one full scan that | ||
| 389 | * feeds shared->reltuples, so the extra passes must not double-count. */ | ||
| 390 | 318 | (void)prism_build_scan( | |
| 391 | heap, | ||
| 392 | index, | ||
| 393 | index_info, | ||
| 394 | shared, | ||
| 395 | true, | ||
| 396 | participant_id == 0, | ||
| 397 | prism_sample_cb, | ||
| 398 | &sc); | ||
| 399 | |||
| 400 | 318 | prism_dsm_sample_counts(dsm_samples)[participant_id] = sc.count; | |
| 401 | 318 | prism_dsm_sample_seen(dsm_samples)[participant_id] = sc.seen; | |
| 402 | |||
| 403 | /* Flush the sub-batch remainder so the participant's contribution to | ||
| 404 | * the progress view is exact (and so every participant makes at least | ||
| 405 | * one flush per scan, whatever share the work-stealing gave it). */ | ||
| 406 | 318 | prism_build_progress_incr_tuples((int64_t)(sc.seen & 1023)); | |
| 407 | |||
| 408 | /* Barrier: all participants done sampling. */ | ||
| 409 | 318 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 410 | 318 | } | |
| 411 | |||
| 412 | uint32_t | ||
| 413 | 318 | prism_pbuild_exec_kmeans( | |
| 414 | int participant_id, | ||
| 415 | PrismBuildShared *shared, | ||
| 416 | PrismDsmSamples *dsm_samples, | ||
| 417 | char *centroids_base, | ||
| 418 | char *km_workers_base, | ||
| 419 | Barrier *barrier) | ||
| 420 | { | ||
| 421 | 318 | Dimension dim = shared->dim; | |
| 422 | 318 | uint32_t km_k = shared->km_k; | |
| 423 | 318 | int nparticipants = shared->nparticipants; | |
| 424 | 318 | float *cents = prism_dsm_centroids(centroids_base); | |
| 425 |
2/2✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
|
318 | float *norms_c = prism_dsm_norms_c(centroids_base, km_k, dim); |
| 426 | /* One seed vector, and the whole km_k-centroid array. */ | ||
| 427 | 318 | const size_t vec_nbytes = (size_t)dim * sizeof(float); | |
| 428 | 318 | const size_t cents_nbytes = (size_t)km_k * vec_nbytes; | |
| 429 | |||
| 430 | /* | ||
| 431 | * Leader-only seed: pick km_k initial centroids from the pooled samples, | ||
| 432 | * spread evenly across the concatenation of every participant's slot. | ||
| 433 | * Centralizing it (rather than slicing by participant) keeps init | ||
| 434 | * independent of how many workers launched and robust to a sample-starved | ||
| 435 | * one. The result varies run to run (work-stealing scan order), so the | ||
| 436 | * parallel build is not bit-reproducible. | ||
| 437 | */ | ||
| 438 |
2/2✓ Branch 0 taken 126 times.
✓ Branch 1 taken 192 times.
|
318 | if (participant_id == 0) |
| 439 | { | ||
| 440 | 82 | uint32_t total_ns = 0; | |
| 441 |
2/2✓ Branch 0 taken 318 times.
✓ Branch 1 taken 126 times.
|
444 | for (int t = 0; t < nparticipants; t++) |
| 442 | 318 | total_ns += prism_dsm_sample_counts(dsm_samples)[t]; | |
| 443 | |||
| 444 |
3/4✓ Branch 0 taken 124 times.
✓ Branch 1 taken 2 times.
✓ Branch 2 taken 82 times.
✗ Branch 3 not taken.
|
126 | uint32_t step = (km_k > 0 && total_ns >= km_k) ? total_ns / km_k : 1; |
| 445 |
2/2✓ Branch 0 taken 742 times.
✓ Branch 1 taken 126 times.
|
868 | for (uint32_t i = 0; i < km_k; i++) |
| 446 | { | ||
| 447 |
2/2✓ Branch 0 taken 741 times.
✓ Branch 1 taken 1 times.
|
742 | uint32_t gidx = (total_ns > 0) ? (i * step) % total_ns : 0; |
| 448 | |||
| 449 | 742 | int t = 0; | |
| 450 | 742 | uint32_t base = 0; | |
| 451 |
2/2✓ Branch 0 taken 1466 times.
✓ Branch 1 taken 1 times.
|
1467 | while (t < nparticipants && |
| 452 |
3/3✓ Branch 0 taken 509 times.
✓ Branch 1 taken 777 times.
✓ Branch 2 taken 180 times.
|
1466 | base + prism_dsm_sample_counts(dsm_samples)[t] <= gidx) |
| 453 | { | ||
| 454 | 725 | base += prism_dsm_sample_counts(dsm_samples)[t]; | |
| 455 | 725 | t++; | |
| 456 | } | ||
| 457 |
2/2✓ Branch 0 taken 741 times.
✓ Branch 1 taken 1 times.
|
742 | if (t < nparticipants) |
| 458 | 741 | memcpy(cents + (size_t)i * dim, | |
| 459 | 741 | prism_dsm_worker_samples(dsm_samples, t) + | |
| 460 | 741 | (size_t)(gidx - base) * dim, | |
| 461 | vec_nbytes); | ||
| 462 | } | ||
| 463 | |||
| 464 |
2/2✓ Branch 0 taken 117 times.
✓ Branch 1 taken 9 times.
|
126 | if (shared->metric == DISTANCE_L2) |
| 465 |
2/2✓ Branch 0 taken 685 times.
✓ Branch 1 taken 117 times.
|
802 | for (uint32_t j = 0; j < km_k; j++) |
| 466 | 685 | norms_c[j] = vs_l2_norm_squared(cents + (size_t)j * dim, dim); | |
| 467 | } | ||
| 468 | |||
| 469 | /* Barrier: initial centroids + norms ready; guards the reads below against | ||
| 470 | * the leader's seed write. */ | ||
| 471 | 318 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 472 | |||
| 473 | 318 | float *my_samples = prism_dsm_worker_samples(dsm_samples, participant_id); | |
| 474 |
2/2✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
|
318 | uint32_t my_n = prism_dsm_sample_counts(dsm_samples)[participant_id]; |
| 475 |
2/2✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
|
318 | float *my_sums = prism_dsm_km_worker_sums( |
| 476 | km_workers_base, km_k, dim, participant_id); | ||
| 477 |
2/2✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
|
318 | uint32_t *my_cnts = prism_dsm_km_worker_cnts( |
| 478 | km_workers_base, km_k, dim, participant_id); | ||
| 479 |
2/2✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
|
318 | float *my_cost = prism_dsm_km_worker_cost( |
| 480 | km_workers_base, km_k, dim, participant_id); | ||
| 481 | |||
| 482 | /* Reduce scratch is leader-only. */ | ||
| 483 |
2/2✓ Branch 0 taken 126 times.
✓ Branch 1 taken 192 times.
|
318 | float *old_cents = participant_id == 0 ? vs_alloc(cents_nbytes) : NULL; |
| 484 | 318 | Size km_sz = prism_dsm_km_workers_size(nparticipants, km_k, dim); | |
| 485 | |||
| 486 | 318 | uint32_t iters = 0; | |
| 487 |
2/2✓ Branch 0 taken 3333 times.
✓ Branch 1 taken 107 times.
|
3440 | for (uint32_t iter = 0; iter < shared->km_max_iterations; iter++) |
| 488 | { | ||
| 489 | 3333 | iters++; | |
| 490 | |||
| 491 | /* Every participant assigns + accumulates over its own samples. */ | ||
| 492 | 3333 | prism_km_assign_and_accumulate( | |
| 493 | my_samples, | ||
| 494 | my_n, | ||
| 495 | cents, | ||
| 496 | norms_c, | ||
| 497 | km_k, | ||
| 498 | dim, | ||
| 499 | shared->metric, | ||
| 500 | my_sums, | ||
| 501 | my_cnts, | ||
| 502 | my_cost); | ||
| 503 | |||
| 504 | /* Barrier: all accumulators written; the leader reduces. */ | ||
| 505 | 3333 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 506 | |||
| 507 |
2/2✓ Branch 0 taken 1174 times.
✓ Branch 1 taken 2159 times.
|
3333 | if (participant_id == 0) |
| 508 | { | ||
| 509 | 1174 | memcpy(old_cents, cents, cents_nbytes); | |
| 510 | |||
| 511 | 1174 | const float **all_sums = vs_alloc(nparticipants * sizeof(float *)); | |
| 512 | 1174 | const uint32_t **all_cnts = vs_alloc( | |
| 513 | nparticipants * sizeof(uint32_t *)); | ||
| 514 | 1174 | float *all_costs = vs_alloc(nparticipants * sizeof(float)); | |
| 515 |
3/3✓ Branch 0 taken 1840 times.
✓ Branch 1 taken 2173 times.
✓ Branch 2 taken 494 times.
|
4507 | for (int t = 0; t < nparticipants; t++) |
| 516 | { | ||
| 517 | 3333 | all_sums[t] = prism_dsm_km_worker_sums( | |
| 518 | km_workers_base, km_k, dim, t); | ||
| 519 | 3333 | all_cnts[t] = prism_dsm_km_worker_cnts( | |
| 520 | km_workers_base, km_k, dim, t); | ||
| 521 | 3333 | all_costs[t] = *prism_dsm_km_worker_cost( | |
| 522 | km_workers_base, km_k, dim, t); | ||
| 523 | } | ||
| 524 | |||
| 525 | 494 | float total_cost; | |
| 526 | 1174 | float shift_sq = kmeans_merge_centroids( | |
| 527 | cents, | ||
| 528 | norms_c, | ||
| 529 | old_cents, | ||
| 530 | all_sums, | ||
| 531 | all_cnts, | ||
| 532 | all_costs, | ||
| 533 | nparticipants, | ||
| 534 | km_k, | ||
| 535 | dim, | ||
| 536 | shared->metric, | ||
| 537 | &total_cost); | ||
| 538 | |||
| 539 | 1174 | float tol_sq = shared->km_tolerance * shared->km_tolerance; | |
| 540 | 1174 | shared->km_converged = (shift_sq < tol_sq); | |
| 541 | 1174 | memset(km_workers_base, 0, km_sz); | |
| 542 | |||
| 543 | 1174 | vs_free(all_sums); | |
| 544 | 1174 | vs_free(all_cnts); | |
| 545 | 1174 | vs_free(all_costs); | |
| 546 | } | ||
| 547 | |||
| 548 | /* Barrier: updated centroids + convergence flag visible to all. */ | ||
| 549 | 3333 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 550 | |||
| 551 |
2/2✓ Branch 0 taken 1506 times.
✓ Branch 1 taken 1827 times.
|
3333 | if (shared->km_converged) |
| 552 | 112 | break; | |
| 553 | } | ||
| 554 | |||
| 555 |
2/2✓ Branch 0 taken 126 times.
✓ Branch 1 taken 192 times.
|
318 | if (old_cents != NULL) |
| 556 | 126 | vs_free(old_cents); | |
| 557 | |||
| 558 | 318 | return iters; | |
| 559 | } | ||
| 560 | |||
| 561 | void | ||
| 562 | 318 | prism_pbuild_exec_root_assign( | |
| 563 | int participant_id, | ||
| 564 | PrismBuildShared *shared, | ||
| 565 | PrismDsmSamples *dsm_samples, | ||
| 566 | PrismDsmRootAssign *dsm_ra, | ||
| 567 | char *centroids_base, | ||
| 568 | Barrier *barrier) | ||
| 569 | { | ||
| 570 | 318 | Dimension dim = shared->dim; | |
| 571 | 318 | uint32_t km_k = shared->km_k; | |
| 572 | 318 | float *cents = prism_dsm_centroids(centroids_base); | |
| 573 | 318 | float *norms_c = prism_dsm_norms_c(centroids_base, km_k, dim); | |
| 574 | |||
| 575 | 506 | kmeans_assign( | |
| 576 | 318 | prism_dsm_worker_samples(dsm_samples, participant_id), | |
| 577 | 0, | ||
| 578 | 318 | prism_dsm_sample_counts(dsm_samples)[participant_id], | |
| 579 | cents, | ||
| 580 | norms_c, | ||
| 581 | km_k, | ||
| 582 | dim, | ||
| 583 | shared->metric, | ||
| 584 | prism_dsm_root_assignments(dsm_ra, participant_id)); | ||
| 585 | |||
| 586 | /* Barrier: all root assignments written; subtree gather can read them. */ | ||
| 587 | 318 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 588 | 318 | } | |
| 589 | |||
| 590 | /* ---------------------------------------------------------------- | ||
| 591 | * Phase 2.5: parallel full-table leaf-encode-reference refinement | ||
| 592 | * | ||
| 593 | * Page-backed mirror of the serial serial_refine_heads: route every row | ||
| 594 | * exactly as the query/insert do (prism_query_route k=1 over the centroid | ||
| 595 | * pages, then head -> leaf), accumulate per-leaf means into the tiled DSM | ||
| 596 | * accumulator, and rewrite each leaf's head-page pt_centroid to the full-table | ||
| 597 | * mean. No in-RAM tree. | ||
| 598 | * ---------------------------------------------------------------- */ | ||
| 599 | |||
| 600 | typedef struct RefineCbState | ||
| 601 | { | ||
| 602 | PrismBuildShared *shared; | ||
| 603 | PrismQueryState *qs; /* page-backed router (workers) */ | ||
| 604 | BlockNumber first_posting; /* head -> leaf: leaf = head - first_posting */ | ||
| 605 | double *sums; /* shared accumulator, indexed leaf - tile_lo */ | ||
| 606 | uint64_t *counts; /* shared accumulator */ | ||
| 607 | Dimension dim; | ||
| 608 | bool cosine; | ||
| 609 | uint32_t tile_lo; /* accumulate only leaves in [tile_lo, tile_hi) */ | ||
| 610 | uint32_t tile_hi; | ||
| 611 | float *scratch; /* per-participant normalized copy (cosine) */ | ||
| 612 | } RefineCbState; | ||
| 613 | |||
| 614 | static void | ||
| 615 | 12600 | prism_refine_cb(void *state, ItemPointerData tid, const float *vec) | |
| 616 | { | ||
| 617 | 12600 | RefineCbState *rs = (RefineCbState *)state; | |
| 618 | 12600 | Dimension dim = rs->dim; | |
| 619 | |||
| 620 | 12600 | (void)tid; | |
| 621 | |||
| 622 | 12600 | uint32_t idx; | |
| 623 | 25200 | const float *v = prism_refine_route_row( | |
| 624 | 12600 | rs->qs, | |
| 625 | rs->first_posting, | ||
| 626 | vec, | ||
| 627 | dim, | ||
| 628 | 12600 | rs->cosine, | |
| 629 | rs->scratch, | ||
| 630 | rs->tile_lo, | ||
| 631 | rs->tile_hi, | ||
| 632 | &idx); | ||
| 633 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 12600 times.
|
12600 | if (v == NULL) |
| 634 | ✗ | return; | |
| 635 | |||
| 636 | 12600 | uint32_t stripe = idx % PRISM_REFINE_LOCK_STRIPES; | |
| 637 | |||
| 638 | 12600 | prism_pbuild_accum_lock(rs->shared, stripe); | |
| 639 | 12600 | double *sum = rs->sums + (size_t)idx * dim; | |
| 640 |
2/2✓ Branch 0 taken 1228800 times.
✓ Branch 1 taken 12600 times.
|
1241400 | for (Dimension j = 0; j < dim; j++) |
| 641 | 1228800 | sum[j] += v[j]; | |
| 642 | 12600 | rs->counts[idx]++; | |
| 643 | 12600 | prism_pbuild_accum_unlock(rs->shared, stripe); | |
| 644 | } | ||
| 645 | |||
| 646 | void | ||
| 647 | 6 | prism_pbuild_exec_refine_paged( | |
| 648 | int participant_id, | ||
| 649 | Relation heap, | ||
| 650 | Relation index, | ||
| 651 | struct IndexInfo *index_info, | ||
| 652 | PrismBuildShared *shared, | ||
| 653 | struct PrismQueryState *qs, | ||
| 654 | BlockNumber first_posting, | ||
| 655 | PrismDsmRefineAccum *accum, | ||
| 656 | Barrier *barrier, | ||
| 657 | PrismLeafWriteFn write_head, | ||
| 658 | void *write_head_ctx) | ||
| 659 | { | ||
| 660 | 6 | Dimension dim = shared->dim; | |
| 661 | 6 | uint32_t nleaves = shared->nlist; /* actual leaf count (published) */ | |
| 662 | 6 | double *sums = prism_dsm_refine_sums(accum); | |
| 663 | 6 | uint64_t *counts = prism_dsm_refine_counts(accum); | |
| 664 | |||
| 665 | /* The accumulator holds at most accum->nleaves leaves (the bounded tile | ||
| 666 | * capacity), so leaves are processed in tiles, re-scanning the heap per | ||
| 667 | * tile. nleaves <= capacity is a single tile (the common case). Leader and | ||
| 668 | * workers derive `tile` identically from the DSM capacity, so the barrier | ||
| 669 | * sequence below stays in lockstep. */ | ||
| 670 | 6 | uint32_t tile = accum->nleaves; | |
| 671 | |||
| 672 | 12 | RefineCbState rs = { | |
| 673 | .shared = shared, | ||
| 674 | .qs = qs, | ||
| 675 | .first_posting = first_posting, | ||
| 676 | .sums = sums, | ||
| 677 | .counts = counts, | ||
| 678 | .dim = dim, | ||
| 679 | 6 | .cosine = (shared->metric == DISTANCE_COSINE), | |
| 680 | 6 | .scratch = vs_alloc((size_t)dim * sizeof(float)), | |
| 681 | }; | ||
| 682 | |||
| 683 | /* Single pass: routing reads only the centroid pages, which refine | ||
| 684 | * never rewrites, so the assignment is a fixed point and one pass | ||
| 685 | * reaches it. */ | ||
| 686 | { | ||
| 687 |
2/2✓ Branch 0 taken 6 times.
✓ Branch 1 taken 6 times.
|
12 | for (uint32_t lo = 0; lo < nleaves; lo += tile) |
| 688 | { | ||
| 689 | 6 | uint32_t hi = (lo + tile < nleaves) ? lo + tile : nleaves; | |
| 690 | 6 | rs.tile_lo = lo; | |
| 691 | 6 | rs.tile_hi = hi; | |
| 692 | |||
| 693 | /* Leader zeroes the (tile-sized) accumulator and resets the scan. | ||
| 694 | */ | ||
| 695 |
2/2✓ Branch 0 taken 2 times.
✓ Branch 1 taken 4 times.
|
6 | if (participant_id == 0) |
| 696 | { | ||
| 697 | 2 | memset(sums, 0, (size_t)(hi - lo) * dim * sizeof(double)); | |
| 698 | 2 | memset(counts, 0, (size_t)(hi - lo) * sizeof(uint64_t)); | |
| 699 | 2 | prism_pbuild_rescan(heap, shared); | |
| 700 | } | ||
| 701 | /* Barrier: accumulator cleared + scan reset before anyone scans. | ||
| 702 | */ | ||
| 703 | 6 | BarrierArriveAndWait( | |
| 704 | barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | ||
| 705 | |||
| 706 | /* Workers cooperatively scan the heap and accumulate page-backed; | ||
| 707 | * the leader does not route (it has no qs), it only | ||
| 708 | * clears/divides, mirroring the phase-3 division of labor. */ | ||
| 709 |
2/2✓ Branch 0 taken 4 times.
✓ Branch 1 taken 2 times.
|
6 | if (participant_id != 0) |
| 710 | 4 | (void)prism_build_scan( | |
| 711 | heap, | ||
| 712 | index, | ||
| 713 | index_info, | ||
| 714 | shared, | ||
| 715 | true, | ||
| 716 | false, | ||
| 717 | prism_refine_cb, | ||
| 718 | &rs); | ||
| 719 | |||
| 720 | /* Barrier: every row accumulated before the leader divides. */ | ||
| 721 | 6 | BarrierArriveAndWait( | |
| 722 | barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | ||
| 723 | |||
| 724 | /* Leader rewrites this tile's leaf head pages = per-leaf means. */ | ||
| 725 |
2/2✓ Branch 0 taken 2 times.
✓ Branch 1 taken 4 times.
|
6 | if (participant_id == 0) |
| 726 | 2 | prism_refine_write_means( | |
| 727 | sums, | ||
| 728 | counts, | ||
| 729 | lo, | ||
| 730 | hi, | ||
| 731 | dim, | ||
| 732 | rs.scratch, | ||
| 733 | write_head, | ||
| 734 | write_head_ctx); | ||
| 735 | /* Barrier: refined heads written before the next tile/pass. */ | ||
| 736 | 6 | BarrierArriveAndWait( | |
| 737 | barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | ||
| 738 | } | ||
| 739 | } | ||
| 740 | |||
| 741 | 6 | vs_free(rs.scratch); | |
| 742 | 6 | } | |
| 743 | |||
| 744 | /* ---------------------------------------------------------------- | ||
| 745 | * Phase 3: posting scan -> cluster-keyed sort (sort-seam path) | ||
| 746 | * | ||
| 747 | * Each worker routes every vector page-backed (the same prism_query_route the | ||
| 748 | * query and insert paths use), RaBitQ-encodes against the target list's head | ||
| 749 | * pt_centroid, and feeds the compact entry into the shared cluster-keyed | ||
| 750 | * sorter (primary + optional SOAR / boundary secondary), via the shared | ||
| 751 | * PrismBuildRouteCtx helper. The leader merges and builds the pages. Memory is | ||
| 752 | * bounded by maintenance_work_mem inside the sorter. | ||
| 753 | * ---------------------------------------------------------------- */ | ||
| 754 | static void | ||
| 755 | 278842 | route_scan_cb(void *state, ItemPointerData tid, const float *vec) | |
| 756 | { | ||
| 757 | 278842 | PrismBuildRouteCtx *ctx = (PrismBuildRouteCtx *)state; | |
| 758 | |||
| 759 | 278842 | prism_build_route_emit(ctx, vec, tid); | |
| 760 |
2/2✓ Branch 0 taken 204 times.
✓ Branch 1 taken 278638 times.
|
278842 | if (((uint64_t)ctx->indtuples & 1023) == 0) |
| 761 | 204 | prism_build_progress_incr_tuples(1024); | |
| 762 | 278842 | } | |
| 763 | |||
| 764 | /* ---------------------------------------------------------------- | ||
| 765 | * Worker entry point — multi-phase build | ||
| 766 | * ---------------------------------------------------------------- */ | ||
| 767 | |||
| 768 | void | ||
| 769 | 192 | prism_parallel_build_main(dsm_segment *seg, shm_toc *toc) | |
| 770 | { | ||
| 771 | 86 | PrismPBuildWorker w; | |
| 772 | 192 | prism_pbuild_worker_attach(toc, &w); | |
| 773 | |||
| 774 | 192 | PrismBuildShared *shared = w.shared; | |
| 775 | 192 | Barrier *barrier = w.barrier; | |
| 776 | 192 | Relation heapRel = w.heapRel; | |
| 777 | 192 | Relation indexRel = w.indexRel; | |
| 778 | 192 | int worker_id = w.worker_id; | |
| 779 | 192 | Dimension dim = w.dim; | |
| 780 | |||
| 781 | /* ---- Phases 1, 2, 2b: the shared per-participant bodies (the leader runs | ||
| 782 | * the very same code as participant 0). Each call includes its phase | ||
| 783 | * barrier(s). ---- */ | ||
| 784 | 192 | void *sample_seg = NULL; | |
| 785 | 86 | PrismDsmSamples *dsm_samples = | |
| 786 | 192 | prism_pbuild_samples_attach(toc, shared, &sample_seg); | |
| 787 | 192 | char *centroids_base = shm_toc_lookup(toc, PRISM_DSM_KEY_CENTROIDS, false); | |
| 788 | 192 | float *cents = prism_dsm_centroids(centroids_base); | |
| 789 | 86 | char *km_workers_base = | |
| 790 | 192 | shm_toc_lookup(toc, PRISM_DSM_KEY_KM_WORKERS, false); | |
| 791 | 86 | PrismDsmRootAssign *dsm_ra = | |
| 792 | 192 | shm_toc_lookup(toc, PRISM_DSM_KEY_ROOT_ASSIGN, false); | |
| 793 | 192 | IndexInfo *indexInfo = BuildIndexInfo(indexRel); | |
| 794 | #ifndef VS_STANDALONE | ||
| 795 | /* The leader marked the build concurrent and scans with an MVCC snapshot; | ||
| 796 | * the worker's freshly built IndexInfo defaults to non-concurrent, so it | ||
| 797 | * must be aligned or heapam's snapshot/OldestXmin assert trips in the scan | ||
| 798 | * (a valid OldestXmin paired with an MVCC snapshot). Matches nbtsort. */ | ||
| 799 | 86 | indexInfo->ii_Concurrent = shared->concurrent; | |
| 800 | #endif | ||
| 801 | |||
| 802 | 192 | prism_pbuild_exec_sampling( | |
| 803 | worker_id, | ||
| 804 | heapRel, | ||
| 805 | indexRel, | ||
| 806 | indexInfo, | ||
| 807 | shared, | ||
| 808 | dsm_samples, | ||
| 809 | barrier); | ||
| 810 | 192 | prism_pbuild_exec_kmeans( | |
| 811 | worker_id, | ||
| 812 | shared, | ||
| 813 | dsm_samples, | ||
| 814 | centroids_base, | ||
| 815 | km_workers_base, | ||
| 816 | barrier); | ||
| 817 | 192 | prism_pbuild_exec_root_assign( | |
| 818 | worker_id, shared, dsm_samples, dsm_ra, centroids_base, barrier); | ||
| 819 | |||
| 820 | /* ---- Phase 2c: batched streaming subtree build (page-backed) ---- | ||
| 821 | * | ||
| 822 | * For a hierarchical tree (>= 2 levels) the workers build per-root-child | ||
| 823 | * subtrees into a bounded ring of slots (subtree DSM is nparticipants | ||
| 824 | * slots, independent of the partition count); the leader consumes each | ||
| 825 | * batch, recording layout counts and keeping the blob in a spillable | ||
| 826 | * store, and later streams every subtree's pages by itself from that | ||
| 827 | * store. Workers pass no callback (leader-only consumes each batch). The | ||
| 828 | * barrier sequence is inside prism_pbuild_stream_subtrees, identical for | ||
| 829 | * leader and workers. A flat (1-level) build has no subtrees — the leader | ||
| 830 | * writes the single level directly, and neither side runs the subtree | ||
| 831 | * barriers. */ | ||
| 832 |
2/2✓ Branch 1 taken 92 times.
✓ Branch 2 taken 100 times.
|
192 | if (vs_hkmeans_nlevels(shared->nlist, shared->fan_out) >= 2) |
| 833 | { | ||
| 834 | /* Ring barrier: the leader creates the subtree ring (sized from the | ||
| 835 | * per-child sample counts) and publishes its handle + slot size | ||
| 836 | * before arriving; attach only after. */ | ||
| 837 | 92 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 838 | 92 | void *ring_seg = NULL; | |
| 839 | 38 | char *subtrees_base = | |
| 840 | 92 | prism_pbuild_subtree_ring_attach(shared, &ring_seg); | |
| 841 | 92 | prism_pbuild_stream_subtrees( | |
| 842 | worker_id, | ||
| 843 | shared->nparticipants, | ||
| 844 | dsm_samples, | ||
| 845 | dsm_ra, | ||
| 846 | cents, | ||
| 847 | shared->km_k, | ||
| 848 | shared->nlist, | ||
| 849 | shared->fan_out, | ||
| 850 | dim, | ||
| 851 | shared->metric, | ||
| 852 | shared->km_max_iterations, | ||
| 853 | subtrees_base, | ||
| 854 | shared->subtree_slot_size, | ||
| 855 | barrier, | ||
| 856 | NULL, | ||
| 857 | NULL, | ||
| 858 | NULL); | ||
| 859 | 90 | prism_pbuild_subtree_ring_release(ring_seg); | |
| 860 | } | ||
| 861 | |||
| 862 | /* Barrier: leader finished streaming the centroid tree + published the | ||
| 863 | * routing state + initialized the sorter; workers build their page-backed | ||
| 864 | * router next (from the just-published shared state). */ | ||
| 865 | 190 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 866 | |||
| 867 | /* ---- Phase 3 setup: page-backed router shared by refine + posting scan | ||
| 868 | * --- */ | ||
| 869 | 190 | void *sortshared = shm_toc_lookup(toc, PRISM_DSM_KEY_SORTSHARED, false); | |
| 870 | 190 | BlockNumber first_posting = shared->first_posting; | |
| 871 | 84 | const float *global_mean = | |
| 872 | 190 | shm_toc_lookup(toc, PRISM_DSM_KEY_GLOBAL_MEAN, false); | |
| 873 | |||
| 874 | /* Exact internal-node centroids, published by the leader before the | ||
| 875 | * tree-ready barrier above: the build descent scores the internal | ||
| 876 | * tree levels exactly (an empty collection leaves the hook inert). */ | ||
| 877 | 190 | void *exact_seg = NULL; | |
| 878 | 84 | PrismExactInternalCentroids exact_centroids; | |
| 879 | 190 | prism_exact_centroid_collection_view( | |
| 880 | 190 | prism_pbuild_exact_centroids_attach(shared, &exact_seg), | |
| 881 | &exact_centroids); | ||
| 882 | |||
| 883 | 190 | uint32_t entry_size = (uint32_t)prism_posting_entry_size(dim); | |
| 884 | 190 | RaBitQParams *rq_params = vs_rabitq_create(dim, shared->rabitq_seed); | |
| 885 | |||
| 886 | /* Per-worker storage over the index for page-backed head/centroid reads | ||
| 887 | * (PG opens one on the worker's indexRel; standalone shares the leader's). | ||
| 888 | */ | ||
| 889 | 190 | VsStorage *storage = prism_pbuild_worker_storage(&w); | |
| 890 | |||
| 891 | /* Routing base — the same PrismIndexBase the query/insert build, so the | ||
| 892 | * worker routes each row identically. nlevels + first_centroid (the | ||
| 893 | * streamed tree's root block) come from the shared state the leader | ||
| 894 | * published; the scales + global mean + fastscan bits also from shared. */ | ||
| 895 | 84 | PrismIndexBase base; | |
| 896 | 274 | prism_build_router_base_init( | |
| 897 | &base, | ||
| 898 | rq_params, | ||
| 899 | storage, | ||
| 900 | dim, | ||
| 901 | 190 | shared->nlevels, | |
| 902 | shared->first_centroid, | ||
| 903 | shared->metric, | ||
| 904 | shared->centroid_format, | ||
| 905 | shared->fastscan_bits, | ||
| 906 | 190 | shared->centroid_error_scale, | |
| 907 | 190 | shared->centroid_beam_scale, | |
| 908 | shared->fan_out, | ||
| 909 | shared->nlist, | ||
| 910 | shared->rabitq_seed, | ||
| 911 | global_mean, | ||
| 912 | 190 | vs_alloc((size_t)dim * sizeof(float))); | |
| 913 | /* Build-only accuracy hook: exact scoring of the internal tree levels | ||
| 914 | * (the query and insert paths never set this). */ | ||
| 915 | 190 | base.exact_internal = &exact_centroids; | |
| 916 | |||
| 917 | 84 | PrismQueryState qs; | |
| 918 | 190 | prism_query_state_init(&qs, &base, 1, PRISM_SECONDARY_TOPK); | |
| 919 | |||
| 920 | /* ---- Phase 2.5: page-backed full-table refine (only when subsampled) | ||
| 921 | * ---- Workers route + accumulate; the leader clears/divides and rewrites | ||
| 922 | * heads. Gated on shared->refine (identical on both sides) so the | ||
| 923 | * barrier sequence stays in lockstep. */ | ||
| 924 |
2/2✓ Branch 0 taken 4 times.
✓ Branch 1 taken 186 times.
|
190 | if (shared->refine) |
| 925 | { | ||
| 926 | /* The accumulator overlays the sample region (dead since the subtree | ||
| 927 | * phase); the leader initialized its header before the tree-ready | ||
| 928 | * barrier above. */ | ||
| 929 | 4 | PrismDsmRefineAccum *accum = prism_pbuild_refine_overlay(dsm_samples); | |
| 930 | 4 | prism_pbuild_exec_refine_paged( | |
| 931 | worker_id, | ||
| 932 | heapRel, | ||
| 933 | indexRel, | ||
| 934 | indexInfo, | ||
| 935 | shared, | ||
| 936 | &qs, | ||
| 937 | first_posting, | ||
| 938 | accum, | ||
| 939 | barrier, | ||
| 940 | NULL, | ||
| 941 | NULL); | ||
| 942 | } | ||
| 943 | |||
| 944 | /* The samples (and the refine overlay riding in them) are dead; hand the | ||
| 945 | * segment back before the posting sort claims its own memory budget. */ | ||
| 946 | 190 | prism_pbuild_samples_release(dsm_samples, sample_seg); | |
| 947 | 190 | dsm_samples = NULL; | |
| 948 | |||
| 949 | /* ---- Phase 3: posting scan -> cluster-keyed sort (page-backed) ---- */ | ||
| 950 | |||
| 951 | /* worker_id is 1..N for launched workers; the sorter's 0-based worker | ||
| 952 | * index is worker_id - 1. Only the workers sort concurrently -- the | ||
| 953 | * leader's merge-only sorter opens after every worker sort finishes and | ||
| 954 | * uses the full budget -- so the budget splits across the workers | ||
| 955 | * alone; splitting across all participants would strand one share for | ||
| 956 | * the whole sort phase. */ | ||
| 957 | 190 | int nsorters = shared->nparticipants > 1 ? shared->nparticipants - 1 : 1; | |
| 958 | 190 | int worker_wm = shared->work_mem_kb / nsorters; | |
| 959 |
1/2✓ Branch 0 taken 106 times.
✗ Branch 1 not taken.
|
190 | if (worker_wm < 64) |
| 960 | 106 | worker_wm = 64; | |
| 961 | 190 | PrismSorter *sorter = prism_pbuild_sort_begin( | |
| 962 | sortshared, seg, worker_id - 1, 0, false, entry_size, worker_wm); | ||
| 963 | |||
| 964 | 84 | PrismBuildRouteCtx route; | |
| 965 | 190 | prism_build_route_ctx_init( | |
| 966 | &route, | ||
| 967 | &qs, | ||
| 968 | sorter, | ||
| 969 | rq_params, | ||
| 970 | storage, | ||
| 971 | first_posting, | ||
| 972 | dim, | ||
| 973 | shared->soar_lambda, | ||
| 974 | shared->boundary_epsilon); | ||
| 975 | |||
| 976 | /* Barrier: the leader reset the scan for the posting phase (after refine | ||
| 977 | * consumed it); workers may now scan. */ | ||
| 978 | 190 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 979 | |||
| 980 | 190 | double heap_tuples = prism_build_scan( | |
| 981 | heapRel, | ||
| 982 | indexRel, | ||
| 983 | indexInfo, | ||
| 984 | shared, | ||
| 985 | true, | ||
| 986 | false, | ||
| 987 | route_scan_cb, | ||
| 988 | &route); | ||
| 989 | |||
| 990 | /* Sub-batch progress remainder, as in the sampling scan. */ | ||
| 991 | 190 | prism_build_progress_incr_tuples( | |
| 992 | 190 | (int64_t)((uint64_t)route.indtuples & 1023)); | |
| 993 | |||
| 994 | 190 | prism_pbuild_sort_performsort(sorter); | |
| 995 | |||
| 996 | /* Barrier: every worker has finished sorting; the leader merges next. */ | ||
| 997 | 190 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 998 | 190 | BarrierDetach(barrier); | |
| 999 | |||
| 1000 | 190 | prism_pbuild_worker_add_counts( | |
| 1001 | shared, route.indtuples, route.soar_dupes, heap_tuples); | ||
| 1002 | |||
| 1003 | 190 | prism_pbuild_sort_end(sorter); | |
| 1004 | 190 | prism_build_route_ctx_cleanup(&route); | |
| 1005 | 190 | prism_query_state_cleanup(&qs); | |
| 1006 | 190 | vs_free(base.pt_global_mean); | |
| 1007 | 190 | prism_pbuild_exact_centroids_release(exact_seg); | |
| 1008 | 190 | prism_pbuild_worker_storage_release(storage); | |
| 1009 | |||
| 1010 | 190 | prism_pbuild_worker_detach(toc, &w); | |
| 1011 | 190 | } | |
| 1012 |