| 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_leader.c - PG parallel index build, leader side | ||
| 6 | * | ||
| 7 | * The leader sets up the DSM, launches workers, participates in the | ||
| 8 | * sampling + k-means phases as worker_id 0, streams the centroid tree to | ||
| 9 | * pages, then (phase 3) merges the workers' cluster-keyed sorted runs and | ||
| 10 | * builds each posting list. Head blocks are formula-derived (cluster c's head | ||
| 11 | * is first_posting + c), so no per-partition head array is materialized. | ||
| 12 | * Returns false if parallelism could not start, so the caller falls back to a | ||
| 13 | * serial build. | ||
| 14 | */ | ||
| 15 | |||
| 16 | #ifdef VS_STANDALONE | ||
| 17 | #include "standalone/instr_time.h" | ||
| 18 | #include "standalone/parallel_ctx.h" /* ParallelContext + lifecycle */ | ||
| 19 | #include "standalone/pg_compat.h" | ||
| 20 | #else | ||
| 21 | #include <postgres.h> | ||
| 22 | |||
| 23 | #include <access/parallel.h> | ||
| 24 | #include <access/table.h> | ||
| 25 | #include <access/tableam.h> | ||
| 26 | #include <access/xloginsert.h> | ||
| 27 | #include <catalog/index.h> | ||
| 28 | #include <commands/progress.h> | ||
| 29 | #include <common/pg_prng.h> | ||
| 30 | #include <executor/instrument.h> | ||
| 31 | #include <miscadmin.h> | ||
| 32 | #include <pgstat.h> | ||
| 33 | #include <portability/instr_time.h> | ||
| 34 | #include <tcop/tcopprot.h> | ||
| 35 | #include <utils/backend_progress.h> | ||
| 36 | #include <utils/memutils.h> | ||
| 37 | #include <utils/rel.h> | ||
| 38 | #include <utils/sampling.h> | ||
| 39 | #include <utils/wait_event.h> | ||
| 40 | #endif | ||
| 41 | |||
| 42 | #include <inttypes.h> | ||
| 43 | #include <math.h> | ||
| 44 | |||
| 45 | #include "algo/distance.h" | ||
| 46 | #include "algo/hkmeans.h" | ||
| 47 | #include "algo/kmeans.h" | ||
| 48 | #include "algo/kmeans_internal.h" | ||
| 49 | #include "algo/vecops.h" | ||
| 50 | #include "core/log.h" | ||
| 51 | #include "core/memory.h" | ||
| 52 | #include "index/build_progress.h" | ||
| 53 | #include "index/centroid_build.h" | ||
| 54 | #include "index/centroid_page.h" | ||
| 55 | #include "index/index_build.h" | ||
| 56 | #include "index/parallel_build.h" | ||
| 57 | #include "index/posting_build.h" | ||
| 58 | #include "index/posting_page.h" | ||
| 59 | #include "quant/fastscan.h" | ||
| 60 | #include "types/vec16.h" | ||
| 61 | #include "types/vec32.h" | ||
| 62 | |||
| 63 | #ifndef VS_STANDALONE | ||
| 64 | #include "build.h" | ||
| 65 | #include "meta.h" | ||
| 66 | #include "storage.h" | ||
| 67 | #include "support_pg.h" | ||
| 68 | #endif | ||
| 69 | |||
| 70 | /* | ||
| 71 | * Build every cluster's posting list from a populated (not yet performsorted) | ||
| 72 | * cluster-keyed sorter. Shared by the serial build and the parallel leader: | ||
| 73 | * performsort, then read entries grouped by cluster and write each list with a | ||
| 74 | * single resident page builder (head at first_posting + c, continuations | ||
| 75 | * appended at the relation's end and chained). Ends the sorter. | ||
| 76 | * | ||
| 77 | * Why iterate clusters 0..nlist rather than just draining the sorter until it | ||
| 78 | * is empty: EVERY cluster needs a posting-list head, including clusters that | ||
| 79 | * received no vectors. The centroid tree's leaf entries reference | ||
| 80 | * first_posting | ||
| 81 | * + c for every c in [0, nlist) (see prism_write_centroid_tree), so an empty | ||
| 82 | * cluster still needs a valid (empty) head block to point at. A loop over the | ||
| 83 | * cluster index emits a head for every cluster uniformly — an empty cluster's | ||
| 84 | * inner while simply does not run and finish() returns an empty head. A | ||
| 85 | * drain-until-empty loop would only produce heads for clusters present in the | ||
| 86 | * stream and would have to separately backfill empty heads for gap clusters | ||
| 87 | * and for all trailing clusters past the last one seen. | ||
| 88 | * | ||
| 89 | * This REQUIRES (and assumes) the sorter returns entries in ascending cluster | ||
| 90 | * order, matching the 0..nlist iteration order: the sorter is keyed on the | ||
| 91 | * cluster id (PG tuplesort on the int4 key; standalone qsort by cluster), so | ||
| 92 | * all entries for a cluster are contiguous and clusters appear in increasing | ||
| 93 | * order. Thus while walking c upward, any remaining entry has cur_cluster >= c | ||
| 94 | * (it equals c for a non-empty cluster, or is greater when c is empty). The | ||
| 95 | * asserts below make that contract explicit: an out-of-order key would pair | ||
| 96 | * entries with the wrong cluster, and an out-of-range id (>= nlist) would | ||
| 97 | * never match any c and be silently dropped (leaving the sorter non-empty at | ||
| 98 | * the end). | ||
| 99 | */ | ||
| 100 | void | ||
| 101 | 276 | prism_posting_build_lists( | |
| 102 | PrismSorter *sorter, | ||
| 103 | VsStorage *storage, | ||
| 104 | uint32_t nlist, | ||
| 105 | Dimension dim, | ||
| 106 | bool fastscan, | ||
| 107 | const RaBitQParams *rq_params, | ||
| 108 | BlockNumber first_posting) | ||
| 109 | { | ||
| 110 | 276 | prism_pbuild_sort_performsort(sorter); | |
| 111 | |||
| 112 | /* | ||
| 113 | * Each list's head was pre-written with its pt_centroid (P^T * centroid, | ||
| 114 | * the RaBitQ encode reference) during the centroid build: adopt it and | ||
| 115 | * append the sorted entries in place. No second head construction, no | ||
| 116 | * in-RAM float centroids — the tree can be freed before this runs — and | ||
| 117 | * a cluster with no entries keeps its on-disk head untouched. | ||
| 118 | */ | ||
| 119 | 276 | uint32_t cur_cluster = 0; | |
| 120 | 276 | const void *entry = NULL; | |
| 121 | 276 | bool have = prism_pbuild_sort_getnext(sorter, &cur_cluster, &entry); | |
| 122 |
3/3✓ Branch 0 taken 1114 times.
✓ Branch 1 taken 9633 times.
✓ Branch 2 taken 194 times.
|
10941 | for (uint32_t c = 0; c < nlist; c++) |
| 123 | { | ||
| 124 | /* Head block of cluster c is the formula first_posting + c; the head | ||
| 125 | * region [first_posting, first_posting + nlist) was pre-extended and | ||
| 126 | * its pt_centroid written during the centroid streaming. Continuation | ||
| 127 | * pages are appended at the relation's end and chained (no per-cluster | ||
| 128 | * reserve), so no O(nlist) reserve arrays are needed. Adopt the | ||
| 129 | * pre-written head and append in place: no second head construction, | ||
| 130 | * and a cluster with no entries keeps its on-disk head untouched. */ | ||
| 131 | 10665 | BlockNumber head_blk = first_posting + c; | |
| 132 | 9551 | PrismPostingBuilder hb; | |
| 133 | 10665 | prism_posting_builder_adopt_head( | |
| 134 | &hb, storage, rq_params, dim, c, head_blk, fastscan); | ||
| 135 | |||
| 136 | /* Ascending, grouped order: a pending entry is never for a cluster we | ||
| 137 | * already finished. If it were < c we would have skipped its head. */ | ||
| 138 |
4/4✓ Branch 0 taken 1149 times.
✓ Branch 1 taken 9516 times.
✓ Branch 2 taken 9516 times.
✓ Branch 3 taken 1114 times.
|
10665 | Assert(!have || cur_cluster >= c); |
| 139 | |||
| 140 |
4/4✓ Branch 0 taken 574869 times.
✓ Branch 1 taken 302 times.
✓ Branch 2 taken 564506 times.
✓ Branch 3 taken 10363 times.
|
575171 | while (have && cur_cluster == c) |
| 141 | { | ||
| 142 | 564506 | prism_posting_entry_add(&hb, entry, dim); | |
| 143 | 564506 | have = prism_pbuild_sort_getnext(sorter, &cur_cluster, &entry); | |
| 144 | } | ||
| 145 | |||
| 146 | 10665 | (void)prism_posting_builder_finish(&hb); | |
| 147 | 10665 | prism_posting_builder_cleanup(&hb); | |
| 148 | } | ||
| 149 | |||
| 150 | /* Every entry must have landed in some cluster's list. A leftover entry | ||
| 151 | * means an id >= nlist (out of range) that matched no c and would | ||
| 152 | * otherwise be silently dropped. */ | ||
| 153 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 276 times.
|
276 | Assert(!have); |
| 154 | 276 | prism_pbuild_sort_end(sorter); | |
| 155 | 276 | } | |
| 156 | |||
| 157 | /* ---------------------------------------------------------------- | ||
| 158 | * Batched streaming tree build — leader-side callbacks | ||
| 159 | * ---------------------------------------------------------------- */ | ||
| 160 | |||
| 161 | /* PLAN-pass batch callback: record each subtree's leaf count + centroid-page | ||
| 162 | * count (no writes) so the leader can size the reserve + block layout, and | ||
| 163 | * keep the blob in the spillable store for the streaming pass to read back | ||
| 164 | * in the same (child) order. */ | ||
| 165 | typedef struct PlanCbArg | ||
| 166 | { | ||
| 167 | uint32_t *nleaves_arr; /* [km_k] */ | ||
| 168 | uint32_t *pages_arr; /* [km_k] */ | ||
| 169 | uint32_t max_ent; | ||
| 170 | uint32_t subtree_nlevels; /* actual depth of the (uniform) subtrees */ | ||
| 171 | Dimension dim; | ||
| 172 | double *leaf_sum; /* [dim] leaf-centroid sum across all subtrees, for | ||
| 173 | * the leaf_mean the leader uses as global_mean */ | ||
| 174 | PrismBlobStore *store; /* subtree blobs, in child order */ | ||
| 175 | } PlanCbArg; | ||
| 176 | |||
| 177 | static void | ||
| 178 | 102 | plan_batch_cb( | |
| 179 | void *arg, | ||
| 180 | const uint32_t *children, | ||
| 181 | uint32_t bs, | ||
| 182 | char *base, | ||
| 183 | uint64_t slot_size) | ||
| 184 | { | ||
| 185 | 102 | PlanCbArg *a = (PlanCbArg *)arg; | |
| 186 |
2/2✓ Branch 0 taken 237 times.
✓ Branch 1 taken 102 times.
|
339 | for (uint32_t s = 0; s < bs; s++) |
| 187 | { | ||
| 188 | 111 | const HKMeansResult *sub = (const HKMeansResult *) | |
| 189 | 237 | prism_dsm_child_subtree(base, s, slot_size); | |
| 190 | 237 | uint32_t child = children[s]; | |
| 191 | 237 | BlockNumber *nfb = vs_alloc((size_t)sub->nnodes * sizeof(BlockNumber)); | |
| 192 | /* Subtrees are built to a uniform depth, so any non-empty one gives | ||
| 193 | * the streamed tree's subtree depth (full depth = this + 1 root | ||
| 194 | * level). */ | ||
| 195 |
1/2✓ Branch 0 taken 237 times.
✗ Branch 1 not taken.
|
237 | if (sub->nleaves > 0) |
| 196 | 237 | a->subtree_nlevels = sub->nlevels; | |
| 197 | 237 | a->nleaves_arr[child] = sub->nleaves; | |
| 198 | 474 | a->pages_arr[child] = (uint32_t) | |
| 199 | 237 | prism_compute_centroid_layout(sub, a->max_ent, 0, nfb); | |
| 200 | 237 | vs_free(nfb); | |
| 201 | |||
| 202 |
1/2✓ Branch 0 taken 237 times.
✗ Branch 1 not taken.
|
237 | if (a->leaf_sum != NULL) |
| 203 | { | ||
| 204 | 237 | const float *lc = hk_leaf_centroids(sub); | |
| 205 |
2/2✓ Branch 0 taken 1783 times.
✓ Branch 1 taken 237 times.
|
2020 | for (uint32_t l = 0; l < sub->nleaves; l++) |
| 206 |
2/2✓ Branch 0 taken 45726 times.
✓ Branch 1 taken 1783 times.
|
47509 | for (Dimension d = 0; d < a->dim; d++) |
| 207 | 45726 | a->leaf_sum[d] += lc[(size_t)l * a->dim + d]; | |
| 208 | } | ||
| 209 | |||
| 210 | 237 | prism_pbuild_blobstore_put(a->store, sub, sub->total_size); | |
| 211 | } | ||
| 212 | 102 | } | |
| 213 | |||
| 214 | /* What both routing-tree shapes hand back to do_parallel_build: the layout | ||
| 215 | * facts the leader publishes to the workers before phase 3. */ | ||
| 216 | typedef struct TreeLayout | ||
| 217 | { | ||
| 218 | BlockNumber first_posting; /* leaf c's head = first_posting + c */ | ||
| 219 | BlockNumber root_blk; /* tree root page */ | ||
| 220 | uint8_t nlevels; /* written depth */ | ||
| 221 | uint32_t nlist; /* actual leaf count */ | ||
| 222 | } TreeLayout; | ||
| 223 | |||
| 224 | /* | ||
| 225 | * Assemble and write the routing tree for the hierarchical shape | ||
| 226 | * (nlevels >= 2): each participant builds the subtree of every root child | ||
| 227 | * it owns; the leader records each batch's layout counts (PLAN), computes | ||
| 228 | * the block layout and the leaf-centroid mean, pre-extends the relation, | ||
| 229 | * replays the spilled subtree blobs to pages, and writes the root page | ||
| 230 | * above them. Fills global_mean and the layout the caller publishes. | ||
| 231 | */ | ||
| 232 | static bool | ||
| 233 | 50 | build_routing_tree_batched( | |
| 234 | PrismBuildShared *shared, | ||
| 235 | VsStorage *storage, | ||
| 236 | PrismBuildProgress *prog, | ||
| 237 | PrismDsmSamples *dsm_samples, | ||
| 238 | PrismDsmRootAssign *dsm_ra, | ||
| 239 | float *cents, | ||
| 240 | uint32_t km_k, | ||
| 241 | uint32_t nlist, | ||
| 242 | uint32_t fan_out, | ||
| 243 | Barrier *barrier, | ||
| 244 | RaBitQParams *rq_params, | ||
| 245 | uint32_t max_ent, | ||
| 246 | BlockNumber first_centroid, | ||
| 247 | float *global_mean, | ||
| 248 | PrismExactCentroidCollector *collector, | ||
| 249 | TreeLayout *out) | ||
| 250 | { | ||
| 251 | 50 | Dimension dim = shared->dim; | |
| 252 | 50 | int nparticipants = shared->nparticipants; | |
| 253 | 50 | PrismCentroidFormat fmt = shared->centroid_format; | |
| 254 | |||
| 255 | /* Size the ring slot from the actual per-child sample counts: a subtree | ||
| 256 | * can never hold more leaves than the samples routed to its root child, | ||
| 257 | * which caps the slot far below the analytic worst case (nlist here is | ||
| 258 | * the worst-case sizing bound, up to fan_out x the requested count). */ | ||
| 259 | 50 | uint32_t max_cc = 1; | |
| 260 | { | ||
| 261 | 50 | uint32_t *child_count = vs_alloc0((size_t)km_k * sizeof(uint32_t)); | |
| 262 | 50 | prism_pbuild_count_children( | |
| 263 | dsm_samples, dsm_ra, nparticipants, km_k, child_count); | ||
| 264 |
3/3✓ Branch 0 taken 126 times.
✓ Branch 1 taken 145 times.
✓ Branch 2 taken 20 times.
|
291 | for (uint32_t c = 0; c < km_k; c++) |
| 265 |
2/2✓ Branch 0 taken 76 times.
✓ Branch 1 taken 50 times.
|
241 | if (child_count[c] > max_cc) |
| 266 | 76 | max_cc = child_count[c]; | |
| 267 | 50 | vs_free(child_count); | |
| 268 | } | ||
| 269 | 50 | uint32_t nlist_c = (nlist + fan_out - 1) / fan_out; | |
| 270 | 20 | uint64_t slot_size = | |
| 271 | 50 | vs_hkmeans_max_blob_size_capped(nlist_c, fan_out, dim, max_cc); | |
| 272 | |||
| 273 | /* The slot must fit one allocation (the leader reads each spilled blob | ||
| 274 | * back through a palloc'd buffer), which also keeps the blob format's | ||
| 275 | * 32-bit interior offsets valid. */ | ||
| 276 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 50 times.
|
50 | if (slot_size > (uint64_t)MaxAllocSize) |
| 277 | ✗ | vs_error( | |
| 278 | "prism: subtree slot %llu MB exceeds the allocation limit " | ||
| 279 | "(nlist %u, fan_out %u); increase fan_out or decrease nlist", | ||
| 280 | (unsigned long long)(slot_size >> 20), | ||
| 281 | nlist, | ||
| 282 | fan_out); | ||
| 283 | /* The ring is the build's only region outside the sample budget; hold | ||
| 284 | * it to the same standard. work_mem_kb == 0 = unbudgeted back-end. */ | ||
| 285 |
2/2✓ Branch 0 taken 20 times.
✓ Branch 1 taken 30 times.
|
50 | if (shared->work_mem_kb > 0 && |
| 286 | 20 | (uint64_t)nparticipants * slot_size > | |
| 287 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 20 times.
|
20 | (uint64_t)shared->work_mem_kb * 1024) |
| 288 | ✗ | vs_error( | |
| 289 | "prism: subtree ring %llu MB exceeds maintenance_work_mem " | ||
| 290 | "(nlist %u, fan_out %u); increase maintenance_work_mem or " | ||
| 291 | "fan_out", | ||
| 292 | (unsigned long long)(((uint64_t)nparticipants * slot_size) >> | ||
| 293 | 20), | ||
| 294 | nlist, | ||
| 295 | fan_out); | ||
| 296 | |||
| 297 | /* Ring barrier: create + publish (handle, slot size) before arriving; | ||
| 298 | * the workers attach after. */ | ||
| 299 | 50 | void *ring_seg = NULL; | |
| 300 | 50 | char *subtrees_base = prism_pbuild_subtree_ring_create( | |
| 301 | shared, nparticipants, slot_size, &ring_seg); | ||
| 302 |
1/2✗ Branch 1 not taken.
✓ Branch 2 taken 20 times.
|
50 | vs_debug( |
| 303 | "prism: subtree ring %d x %llu KB (largest child %u samples)", | ||
| 304 | nparticipants, | ||
| 305 | (unsigned long long)(slot_size >> 10), | ||
| 306 | max_cc); | ||
| 307 | 50 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 308 | |||
| 309 | 50 | uint32_t *nleaves_arr = vs_alloc0((size_t)km_k * sizeof(uint32_t)); | |
| 310 | 50 | uint32_t *pages_arr = vs_alloc0((size_t)km_k * sizeof(uint32_t)); | |
| 311 | |||
| 312 | /* PLAN pass: build subtrees, discover leaf + page counts (and the | ||
| 313 | * leaf-centroid sum for global_mean). */ | ||
| 314 | 50 | prism_build_report_phase(prog, PRISM_BUILD_PHASE_SUBTREES); | |
| 315 | 147 | PlanCbArg planarg = { | |
| 316 | .nleaves_arr = nleaves_arr, | ||
| 317 | .pages_arr = pages_arr, | ||
| 318 | .max_ent = max_ent, | ||
| 319 | .subtree_nlevels = 1, | ||
| 320 | .dim = dim, | ||
| 321 | 49 | .leaf_sum = vs_alloc0((size_t)dim * sizeof(double)), | |
| 322 | 49 | .store = prism_pbuild_blobstore_begin(), | |
| 323 | }; | ||
| 324 | /* The batch schedule is largest-first; the blob store receives the | ||
| 325 | * subtrees in that order, so the replay below needs the same order | ||
| 326 | * to place each blob at its child's reserved block range. */ | ||
| 327 | 49 | uint32_t *child_order = vs_alloc((size_t)km_k * sizeof(uint32_t)); | |
| 328 | 49 | prism_pbuild_stream_subtrees( | |
| 329 | 0, | ||
| 330 | nparticipants, | ||
| 331 | dsm_samples, | ||
| 332 | dsm_ra, | ||
| 333 | cents, | ||
| 334 | km_k, | ||
| 335 | nlist, | ||
| 336 | fan_out, | ||
| 337 | dim, | ||
| 338 | shared->metric, | ||
| 339 | shared->km_max_iterations, | ||
| 340 | subtrees_base, | ||
| 341 | slot_size, | ||
| 342 | barrier, | ||
| 343 | plan_batch_cb, | ||
| 344 | &planarg, | ||
| 345 | child_order); | ||
| 346 | |||
| 347 | /* Every batch is consumed into the blob store; the ring is dead before | ||
| 348 | * the write pass starts. */ | ||
| 349 | 49 | prism_pbuild_subtree_ring_release(ring_seg); | |
| 350 | |||
| 351 | /* Root page(s) occupy the reserved block(s) at first_centroid; | ||
| 352 | * subtrees follow, so meta.first_centroid stays 1 (root written last, | ||
| 353 | * in place). */ | ||
| 354 | 49 | uint32_t root_pages = (km_k + max_ent - 1) / max_ent; | |
| 355 | 49 | uint32_t *leaf_off = vs_alloc((size_t)km_k * sizeof(uint32_t)); | |
| 356 | 49 | uint32_t *block_off = vs_alloc((size_t)km_k * sizeof(uint32_t)); | |
| 357 | 49 | uint32_t lo = 0; | |
| 358 | 49 | BlockNumber bo = 0; | |
| 359 |
3/3✓ Branch 0 taken 126 times.
✓ Branch 1 taken 141 times.
✓ Branch 2 taken 19 times.
|
286 | for (uint32_t c = 0; c < km_k; c++) |
| 360 | { | ||
| 361 | 237 | leaf_off[c] = lo; | |
| 362 | 237 | block_off[c] = (uint32_t)bo; | |
| 363 | 237 | lo += nleaves_arr[c]; | |
| 364 | 237 | bo += pages_arr[c]; | |
| 365 | } | ||
| 366 | 49 | uint32_t actual_nlist = lo; | |
| 367 | 49 | BlockNumber subtree_base = first_centroid + root_pages; | |
| 368 | 49 | BlockNumber first_posting = subtree_base + bo; | |
| 369 | |||
| 370 | /* The centroid-page region [first_centroid, first_posting) is now | ||
| 371 | * sized; the streaming pass below collects every internal node's | ||
| 372 | * exact centroids into it (the root included). */ | ||
| 373 |
2/2✓ Branch 0 taken 45 times.
✓ Branch 1 taken 4 times.
|
49 | if (collector != NULL) |
| 374 |
1/2✓ Branch 0 taken 19 times.
✗ Branch 1 not taken.
|
45 | prism_exact_centroid_collector_init( |
| 375 | collector, | ||
| 376 | dim, | ||
| 377 | fmt, | ||
| 378 | first_centroid, | ||
| 379 | 45 | (uint32_t)(first_posting - first_centroid), | |
| 380 |
1/2✓ Branch 0 taken 19 times.
✗ Branch 1 not taken.
|
45 | prism_exact_centroid_budget(shared->work_mem_kb), |
| 381 | prism_exact_centroid_expected_slots(actual_nlist, fan_out)); | ||
| 382 | |||
| 383 | /* Leaf-centroid mean from the PLAN pass -> the encoder centering. */ | ||
| 384 |
2/2✓ Branch 0 taken 1068 times.
✓ Branch 1 taken 49 times.
|
1117 | for (Dimension d = 0; d < dim; d++) |
| 385 | 1068 | global_mean[d] = actual_nlist > 0 ? (float)(planarg.leaf_sum[d] / | |
| 386 | 1068 | (double)actual_nlist) | |
| 387 |
1/2✓ Branch 0 taken 1068 times.
✗ Branch 1 not taken.
|
1068 | : 0.0f; |
| 388 |
2/2✓ Branch 0 taken 3 times.
✓ Branch 1 taken 46 times.
|
49 | if (shared->metric == DISTANCE_COSINE) |
| 389 | 3 | vs_l2_normalize(global_mean, dim); | |
| 390 | 49 | vs_free(planarg.leaf_sum); | |
| 391 | 49 | planarg.leaf_sum = NULL; | |
| 392 | |||
| 393 | /* Head blocks are formula-derived: leaf c's head is first_posting + c, | ||
| 394 | * a contiguous head region of actual_nlist pages. Pre-extend the | ||
| 395 | * relation to cover the centroid pages + the head region so the write | ||
| 396 | * pass can write both at reserved blocks; continuation pages are | ||
| 397 | * appended past it during prism_posting_build_lists. No O(nlist) reserve | ||
| 398 | * arrays. */ | ||
| 399 | 49 | prism_build_reserve_layout(storage, first_posting + actual_nlist); | |
| 400 | |||
| 401 | /* Streaming pass: read each subtree blob back from the store (child | ||
| 402 | * order matches the append order) and stream its centroid + head | ||
| 403 | * pages to the reserved block range, then the root page. Leader-only: | ||
| 404 | * the PLAN pass already produced every subtree, so the workers have | ||
| 405 | * nothing to contribute here and run no barriers for this phase. */ | ||
| 406 | 49 | prism_build_report_phase(prog, PRISM_BUILD_PHASE_CENTROID); | |
| 407 | 49 | BlockNumber *subtree_root_blk = vs_alloc( | |
| 408 | 30 | (size_t)km_k * sizeof(BlockNumber)); | |
| 409 | 19 | PrismHeadWriteCtx head; | |
| 410 | 68 | prism_head_write_ctx_init( | |
| 411 | 49 | &head, storage, rq_params, dim, shared->fastscan, first_posting); | |
| 412 | 49 | HKMeansResult *blob = vs_alloc(slot_size); | |
| 413 | 49 | prism_pbuild_blobstore_rewind(planarg.store); | |
| 414 |
3/3✓ Branch 0 taken 126 times.
✓ Branch 1 taken 141 times.
✓ Branch 2 taken 19 times.
|
286 | for (uint32_t i = 0; i < km_k; i++) |
| 415 | { | ||
| 416 | /* Blobs arrive in the largest-first batch order; c is the child | ||
| 417 | * whose reserved block range this blob belongs to. */ | ||
| 418 | 237 | uint32_t c = child_order[i]; | |
| 419 | 237 | (void)prism_pbuild_blobstore_get(planarg.store, blob, slot_size); | |
| 420 | 237 | BlockNumber base_blk = subtree_base + block_off[c]; | |
| 421 | 237 | subtree_root_blk[c] = prism_routing_subtree_write( | |
| 422 | storage, | ||
| 423 | blob, | ||
| 424 | dim, | ||
| 425 | shared->metric, | ||
| 426 | fan_out, | ||
| 427 | 1, /* subtrees hang off the level-0 root */ | ||
| 428 | fmt, | ||
| 429 | rq_params, | ||
| 430 | global_mean, | ||
| 431 | first_posting, | ||
| 432 | 237 | leaf_off[c], | |
| 433 | base_blk, | ||
| 434 | prism_write_leaf_head, | ||
| 435 | &head, | ||
| 436 | collector); | ||
| 437 | } | ||
| 438 | 49 | vs_free(blob); | |
| 439 | 49 | vs_free(child_order); | |
| 440 | 49 | prism_pbuild_blobstore_end(planarg.store); | |
| 441 | 49 | planarg.store = NULL; | |
| 442 | |||
| 443 | /* Root centroid page at the reserved first_centroid (children = | ||
| 444 | * subtree roots). Written last, but in place, so first_centroid | ||
| 445 | * stays 1. */ | ||
| 446 | 49 | BlockNumber root_blk = first_centroid; | |
| 447 | 49 | prism_centroid_write_node( | |
| 448 | storage, | ||
| 449 | dim, | ||
| 450 | cents, | ||
| 451 | km_k, | ||
| 452 | fmt, | ||
| 453 | 0, | ||
| 454 | 0, | ||
| 455 | 30 | (uint16_t)fan_out, | |
| 456 | rq_params, | ||
| 457 | global_mean, | ||
| 458 | subtree_root_blk, | ||
| 459 | NULL, | ||
| 460 | root_blk, | ||
| 461 | collector); | ||
| 462 | |||
| 463 | 49 | prism_head_write_ctx_cleanup(&head); | |
| 464 | 49 | vs_free(subtree_root_blk); | |
| 465 | 49 | vs_free(leaf_off); | |
| 466 | 49 | vs_free(block_off); | |
| 467 | 49 | vs_free(nleaves_arr); | |
| 468 | 49 | vs_free(pages_arr); | |
| 469 | |||
| 470 | 49 | out->first_posting = first_posting; | |
| 471 | 49 | out->root_blk = root_blk; | |
| 472 | 49 | out->nlevels = (uint8_t)(planarg.subtree_nlevels + 1); | |
| 473 | 49 | out->nlist = actual_nlist; | |
| 474 | 49 | return true; | |
| 475 | } | ||
| 476 | |||
| 477 | /* | ||
| 478 | * Assemble and write the routing tree for the flat shape (nlevels == 1): | ||
| 479 | * the root k-means centroids already ARE the leaves, so stream the | ||
| 480 | * one-level tree directly (root = leaf-parent page at first_centroid). | ||
| 481 | * Fills global_mean and the layout the caller publishes; false when the | ||
| 482 | * flat carrier cannot be built. | ||
| 483 | */ | ||
| 484 | static bool | ||
| 485 | 76 | build_routing_tree_flat( | |
| 486 | PrismBuildShared *shared, | ||
| 487 | VsStorage *storage, | ||
| 488 | PrismBuildProgress *prog, | ||
| 489 | float *cents, | ||
| 490 | uint32_t km_k, | ||
| 491 | uint32_t fan_out, | ||
| 492 | RaBitQParams *rq_params, | ||
| 493 | uint32_t max_ent, | ||
| 494 | BlockNumber first_centroid, | ||
| 495 | float *global_mean, | ||
| 496 | TreeLayout *out) | ||
| 497 | { | ||
| 498 | 76 | Dimension dim = shared->dim; | |
| 499 | 76 | PrismCentroidFormat fmt = shared->centroid_format; | |
| 500 | |||
| 501 | /* Flat (nlevels == 1): cents already holds every leaf centroid; stream | ||
| 502 | * the one-level tree directly (root = leaf-parent at first_centroid). | ||
| 503 | */ | ||
| 504 |
2/2✓ Branch 0 taken 4 times.
✓ Branch 1 taken 72 times.
|
76 | if (shared->metric == DISTANCE_COSINE) |
| 505 |
2/2✓ Branch 0 taken 41 times.
✓ Branch 1 taken 4 times.
|
45 | for (uint32_t c = 0; c < km_k; c++) |
| 506 | 41 | vs_l2_normalize(cents + (size_t)c * dim, dim); | |
| 507 | |||
| 508 | 76 | HKMeansResult *flat = vs_hkmeans_build_flat(cents, km_k, fan_out, dim); | |
| 509 |
2/2✓ Branch 0 taken 24 times.
✓ Branch 1 taken 52 times.
|
76 | if (flat == NULL) |
| 510 | ✗ | return false; | |
| 511 | |||
| 512 | 76 | uint32_t actual_nlist = flat->nleaves; | |
| 513 | 76 | BlockNumber *nfb = vs_alloc((size_t)flat->nnodes * sizeof(BlockNumber)); | |
| 514 | 24 | uint32_t centroid_pages = (uint32_t) | |
| 515 | 76 | prism_compute_centroid_layout(flat, max_ent, 0, nfb); | |
| 516 | 76 | vs_free(nfb); | |
| 517 | 76 | BlockNumber first_posting = first_centroid + centroid_pages; | |
| 518 | |||
| 519 | /* Leaf-centroid mean -> the encoder centering (flat tree in hand). */ | ||
| 520 | 76 | vec32_mean(hk_leaf_centroids(flat), actual_nlist, dim, global_mean); | |
| 521 |
2/2✓ Branch 0 taken 4 times.
✓ Branch 1 taken 72 times.
|
76 | if (shared->metric == DISTANCE_COSINE) |
| 522 | 4 | vs_l2_normalize(global_mean, dim); | |
| 523 | |||
| 524 | /* Head region: actual_nlist pages at first_posting (leaf c -> head | ||
| 525 | * first_posting + c). Pre-extend to cover centroid + head region; | ||
| 526 | * continuations append past it. No O(nlist) reserve. */ | ||
| 527 | 76 | prism_build_reserve_layout(storage, first_posting + actual_nlist); | |
| 528 | |||
| 529 | 76 | prism_build_report_phase(prog, PRISM_BUILD_PHASE_CENTROID); | |
| 530 | 24 | PrismHeadWriteCtx head; | |
| 531 | 100 | prism_head_write_ctx_init( | |
| 532 | 76 | &head, storage, rq_params, dim, shared->fastscan, first_posting); | |
| 533 | 76 | BlockNumber root_blk = prism_routing_subtree_write( | |
| 534 | storage, | ||
| 535 | flat, | ||
| 536 | dim, | ||
| 537 | shared->metric, | ||
| 538 | fan_out, | ||
| 539 | 0, /* the flat tree IS the root level */ | ||
| 540 | fmt, | ||
| 541 | rq_params, | ||
| 542 | global_mean, | ||
| 543 | first_posting, | ||
| 544 | 0, | ||
| 545 | first_centroid, | ||
| 546 | prism_write_leaf_head, | ||
| 547 | &head, | ||
| 548 | /* the flat tree's only level is the leaf level — nothing | ||
| 549 | * internal to collect */ | ||
| 550 | NULL); | ||
| 551 | 76 | prism_head_write_ctx_cleanup(&head); | |
| 552 | |||
| 553 | 76 | out->first_posting = first_posting; | |
| 554 | 76 | out->root_blk = root_blk; | |
| 555 | 76 | out->nlevels = (uint8_t)flat->nlevels; | |
| 556 | 76 | out->nlist = actual_nlist; | |
| 557 | 76 | vs_free(flat); | |
| 558 | 76 | return true; | |
| 559 | } | ||
| 560 | |||
| 561 | bool | ||
| 562 | 127 | do_parallel_build( | |
| 563 | Relation heap, | ||
| 564 | Relation index, | ||
| 565 | struct IndexInfo *index_info, | ||
| 566 | const PrismBuildConfig *config, | ||
| 567 | VsStorage *storage, | ||
| 568 | struct PrismBuildProgress *prog, | ||
| 569 | uint32_t *out_nlist, | ||
| 570 | uint8_t *out_tree_nlevels, | ||
| 571 | double *out_heap_tuples, | ||
| 572 | double *out_indtuples, | ||
| 573 | double *out_soar_dupes, | ||
| 574 | float **out_global_mean, | ||
| 575 | BlockNumber *out_first_posting) | ||
| 576 | { | ||
| 577 | 127 | int nworkers = index_info->ii_ParallelWorkers; | |
| 578 | |||
| 579 | 127 | *out_nlist = 0; | |
| 580 | 127 | *out_tree_nlevels = 0; | |
| 581 |
1/2✓ Branch 0 taken 127 times.
✗ Branch 1 not taken.
|
127 | if (out_first_posting) |
| 582 | 127 | *out_first_posting = InvalidBlockNumber; | |
| 583 |
1/2✓ Branch 0 taken 127 times.
✗ Branch 1 not taken.
|
127 | if (out_global_mean) |
| 584 | 127 | *out_global_mean = NULL; | |
| 585 | |||
| 586 | 45 | PrismPBuildLeader lead; | |
| 587 |
1/2✗ Branch 1 not taken.
✓ Branch 2 taken 127 times.
|
127 | if (!prism_pbuild_setup_shared(&lead, heap, index, config, nworkers)) |
| 588 | ✗ | return false; | |
| 589 | |||
| 590 | 127 | ParallelContext *pcxt = lead.pcxt; | |
| 591 | 127 | PrismBuildShared *shared = lead.shared; | |
| 592 | 127 | Barrier *barrier = lead.barrier; | |
| 593 | 127 | PrismDsmSamples *dsm_samples = lead.dsm_samples; | |
| 594 | 127 | char *centroids_base = lead.centroids_base; | |
| 595 | 127 | float *cents = lead.cents; | |
| 596 | 127 | char *km_workers_base = lead.km_workers_base; | |
| 597 | 127 | PrismDsmRootAssign *dsm_ra = lead.dsm_ra; | |
| 598 | 127 | WalUsage *walusage = lead.walusage; | |
| 599 | 127 | BufferUsage *bufferusage = lead.bufferusage; | |
| 600 | 127 | int nparticipants = lead.nparticipants; | |
| 601 | 127 | uint32_t km_k = lead.km_k; | |
| 602 | 127 | Dimension dim = lead.dim; | |
| 603 | 127 | uint32_t nlist = lead.nlist; | |
| 604 | 127 | uint64_t rabitq_seed = lead.rabitq_seed; | |
| 605 | 127 | uint32_t fan_out = lead.fan_out; | |
| 606 | |||
| 607 | /* | ||
| 608 | * Introspection: name the dataset-scaling allocations up front (always | ||
| 609 | * logged, regardless of the GUC) so an OOM in any of them is | ||
| 610 | * pre-explained, and record the committed DSM size so the per-phase memory | ||
| 611 | * lines can add it. nlist here is the worst-case upper bound (the tree is | ||
| 612 | * not built yet). | ||
| 613 | */ | ||
| 614 | 127 | prism_build_report_dsm_bytes(prog, (uint64_t)lead.dsm_total); | |
| 615 | 127 | prism_build_report_planned_alloc( | |
| 616 | prog, | ||
| 617 | 127 | (uint64_t)prism_dsm_samples_size( | |
| 618 | nparticipants, lead.max_per_worker, dim), | ||
| 619 | /* The tree is streamed to pages from a bounded ring of subtree | ||
| 620 | * slots (part of dsm_total); it is never held whole in memory. */ | ||
| 621 | 0, | ||
| 622 | 127 | (uint64_t)lead.dsm_total); | |
| 623 | |||
| 624 | 45 | instr_time t_launch_start; | |
| 625 | 127 | INSTR_TIME_SET_CURRENT(t_launch_start); | |
| 626 | |||
| 627 | /* | ||
| 628 | * The leader participates in every phase, so it joins the dynamic barrier | ||
| 629 | * as a party (workers attach the same way as they start). Attaching before | ||
| 630 | * launching workers guarantees the leader is counted before any worker can | ||
| 631 | * arrive, so the barrier never advances a phase without it. | ||
| 632 | */ | ||
| 633 | 127 | BarrierAttach(barrier); | |
| 634 | |||
| 635 | /* ---- Launch workers + wait until they've all attached ---- */ | ||
| 636 |
2/2✓ Branch 1 taken 1 times.
✓ Branch 2 taken 126 times.
|
127 | if (!prism_pbuild_launch(pcxt, barrier, shared)) |
| 637 | { | ||
| 638 | /* Teardown already ran; hand the sample segment back too so the | ||
| 639 | * serial fallback starts from a clean budget. */ | ||
| 640 | 1 | prism_pbuild_samples_release(dsm_samples, lead.sample_seg); | |
| 641 | 1 | return false; | |
| 642 | } | ||
| 643 | |||
| 644 | /* The launch may have narrowed the participant count to the party that | ||
| 645 | * actually attached; partition the phases below over that count. */ | ||
| 646 | 126 | nparticipants = shared->nparticipants; | |
| 647 | |||
| 648 | /* ==== Leader participates in all phases as worker_id=0 ==== */ | ||
| 649 | |||
| 650 | /* ---- Phase 1: sampling (leader runs the worker body as participant 0) | ||
| 651 | * ---- */ | ||
| 652 | 126 | prism_build_report_phase(prog, PRISM_BUILD_PHASE_SAMPLE); | |
| 653 | 126 | prism_pbuild_exec_sampling( | |
| 654 | 0, heap, index, index_info, shared, dsm_samples, barrier); | ||
| 655 | |||
| 656 | 44 | instr_time t_sample_end; | |
| 657 | 126 | INSTR_TIME_SET_CURRENT(t_sample_end); | |
| 658 | 126 | INSTR_TIME_SUBTRACT(t_sample_end, t_launch_start); | |
| 659 |
1/2✗ Branch 1 not taken.
✓ Branch 2 taken 44 times.
|
126 | vs_debug( |
| 660 | "prism: phase 1 (sampling) %.1fms", | ||
| 661 | INSTR_TIME_GET_MILLISEC(t_sample_end)); | ||
| 662 | |||
| 663 | 44 | instr_time t_km_start; | |
| 664 | 126 | INSTR_TIME_SET_CURRENT(t_km_start); | |
| 665 | |||
| 666 | /* ---- Phase 2: root k-means. The leader runs the worker body as | ||
| 667 | * participant 0; inside it additionally seeds the initial centroids and | ||
| 668 | * reduces the per-iteration accumulators (gated on participant 0). ---- */ | ||
| 669 | 126 | prism_build_report_phase(prog, PRISM_BUILD_PHASE_KMEANS); | |
| 670 | 126 | uint32_t km_iters = prism_pbuild_exec_kmeans( | |
| 671 | 0, shared, dsm_samples, centroids_base, km_workers_base, barrier); | ||
| 672 | |||
| 673 | { | ||
| 674 | 44 | instr_time t_km_elapsed; | |
| 675 | 126 | INSTR_TIME_SET_CURRENT(t_km_elapsed); | |
| 676 | 126 | INSTR_TIME_SUBTRACT(t_km_elapsed, t_km_start); | |
| 677 |
1/2✗ Branch 1 not taken.
✓ Branch 2 taken 44 times.
|
126 | vs_debug( |
| 678 | "prism: root kmeans %.1fms (%u iters, k=%u)", | ||
| 679 | INSTR_TIME_GET_MILLISEC(t_km_elapsed), | ||
| 680 | km_iters, | ||
| 681 | km_k); | ||
| 682 | } | ||
| 683 | |||
| 684 | /* | ||
| 685 | * Compute nlevels to decide flat vs hierarchical assembly. For | ||
| 686 | * nlevels == 1 (flat) the root k-means already produced every leaf, so | ||
| 687 | * the tree is built from the centroids directly. For nlevels >= 2 each | ||
| 688 | * participant builds the full subtree — to whatever depth nlist/fan_out | ||
| 689 | * needs — for the root children it owns, and the leader streams each to | ||
| 690 | * centroid pages a batch at a time (no in-RAM whole-tree assembly). | ||
| 691 | */ | ||
| 692 | 126 | uint32_t nlevels = vs_hkmeans_nlevels(nlist, fan_out); | |
| 693 | |||
| 694 | /* ---- Phase 2b: root assignment (leader as participant 0). ---- */ | ||
| 695 | 126 | prism_pbuild_exec_root_assign( | |
| 696 | 0, shared, dsm_samples, dsm_ra, centroids_base, barrier); | ||
| 697 | |||
| 698 | /* ---- Batched streaming tree build -------------------------------------- | ||
| 699 | * Workers build per-root-child subtrees into a bounded ring of slots; the | ||
| 700 | * leader records each batch's layout counts, keeps the blobs in a | ||
| 701 | * spillable store, and streams them to centroid pages. Peak subtree DSM is | ||
| 702 | * nparticipants slots, independent of nlist. | ||
| 703 | * ------------------------------ */ | ||
| 704 | 126 | prism_build_report_phase(prog, PRISM_BUILD_PHASE_SETUP); | |
| 705 | 126 | RaBitQParams *rq_params = vs_rabitq_create(dim, rabitq_seed); | |
| 706 | 126 | PrismCentroidFormat fmt = shared->centroid_format; | |
| 707 | 126 | uint32_t max_ent = prism_centroid_max_entries_fmt(dim, fmt); | |
| 708 | |||
| 709 | /* | ||
| 710 | * global_mean = mean of the LEAF centroids (the encoder centering) -- an | ||
| 711 | * unweighted per-cluster mean, not the per-vector sample mean. The | ||
| 712 | * quantization quality of every centroid and posting code depends on this | ||
| 713 | * anchor. Pages are written only after the PLAN pass has built every | ||
| 714 | * subtree, so | ||
| 715 | * the leaf-centroid sum is accumulated there (plan_batch_cb) and the mean | ||
| 716 | * is ready before any page is encoded. Filled per branch below. | ||
| 717 | */ | ||
| 718 | 126 | const size_t vec_nbytes = (size_t)dim * sizeof(float); | |
| 719 | 126 | float *global_mean = vs_alloc(vec_nbytes); | |
| 720 | |||
| 721 | 126 | BlockNumber first_centroid = PRISM_FIRST_CENTROID_BLKNO; | |
| 722 | 126 | BlockNumber first_posting = 0; | |
| 723 | 126 | BlockNumber root_blk = InvalidBlockNumber; | |
| 724 | |||
| 725 | /* Exact internal-centroid collection for the phase-2.5/3 build | ||
| 726 | * descent (see PrismExactCentroidCollector in index_build.h). */ | ||
| 727 | 126 | PrismExactCentroidCollector exact_centroids = {0}; | |
| 728 |
2/2✓ Branch 0 taken 20 times.
✓ Branch 1 taken 24 times.
|
126 | PrismExactCentroidCollector *collector = |
| 729 | 106 | prism_exact_centroid_enabled(nlevels, fmt) ? &exact_centroids | |
| 730 |
2/2✓ Branch 0 taken 26 times.
✓ Branch 1 taken 56 times.
|
82 | : NULL; |
| 731 | |||
| 732 | 44 | TreeLayout layout; | |
| 733 | 44 | bool tree_ok; | |
| 734 |
2/2✓ Branch 0 taken 50 times.
✓ Branch 1 taken 76 times.
|
126 | if (nlevels >= 2) |
| 735 | 50 | tree_ok = build_routing_tree_batched( | |
| 736 | shared, | ||
| 737 | storage, | ||
| 738 | prog, | ||
| 739 | dsm_samples, | ||
| 740 | dsm_ra, | ||
| 741 | cents, | ||
| 742 | km_k, | ||
| 743 | nlist, | ||
| 744 | fan_out, | ||
| 745 | barrier, | ||
| 746 | rq_params, | ||
| 747 | max_ent, | ||
| 748 | first_centroid, | ||
| 749 | global_mean, | ||
| 750 | collector, | ||
| 751 | &layout); | ||
| 752 | else | ||
| 753 | 76 | tree_ok = build_routing_tree_flat( | |
| 754 | shared, | ||
| 755 | storage, | ||
| 756 | prog, | ||
| 757 | cents, | ||
| 758 | km_k, | ||
| 759 | fan_out, | ||
| 760 | rq_params, | ||
| 761 | max_ent, | ||
| 762 | first_centroid, | ||
| 763 | global_mean, | ||
| 764 | &layout); | ||
| 765 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 125 times.
|
125 | if (!tree_ok) |
| 766 | { | ||
| 767 | ✗ | if (collector != NULL) | |
| 768 | ✗ | prism_exact_centroid_collector_cleanup(collector); | |
| 769 | ✗ | vs_free(global_mean); | |
| 770 | ✗ | prism_pbuild_samples_release(dsm_samples, lead.sample_seg); | |
| 771 | ✗ | WaitForParallelWorkersToFinish(pcxt); | |
| 772 | ✗ | prism_pbuild_teardown(pcxt); | |
| 773 | ✗ | return false; | |
| 774 | } | ||
| 775 | 125 | nlist = layout.nlist; | |
| 776 | 125 | first_posting = layout.first_posting; | |
| 777 | 125 | root_blk = layout.root_blk; | |
| 778 | 125 | const uint8_t out_nlevels = layout.nlevels; | |
| 779 | |||
| 780 | /* Publish the exact internal-node centroids the tree write collected, | ||
| 781 | * for the workers' phase-2.5/3 build descent (exact-centroid seam; must | ||
| 782 | * precede the tree-ready barrier). Without a collector the collection | ||
| 783 | * is an empty header and the workers' scoring hook stays inert. */ | ||
| 784 | 125 | void *exact_seg = NULL; | |
| 785 | 125 | char *exact_cents = prism_pbuild_exact_centroids_create( | |
| 786 | shared, | ||
| 787 | prism_exact_centroid_collection_size(collector), | ||
| 788 | &exact_seg); | ||
| 789 | 125 | prism_exact_centroid_collection_write(collector, exact_cents); | |
| 790 |
2/2✓ Branch 0 taken 45 times.
✓ Branch 1 taken 80 times.
|
125 | if (collector != NULL) |
| 791 | 45 | prism_exact_centroid_collector_cleanup(collector); | |
| 792 | |||
| 793 | /* Publish the routing state the workers read in phase 3: nlist, the tree | ||
| 794 | * root block + depth, the posting-head base (leaf c's head = first_posting | ||
| 795 | * + c), the global mean, and the page store. */ | ||
| 796 | 125 | shared->nlist = nlist; | |
| 797 | 125 | shared->first_centroid = root_blk; | |
| 798 | 125 | shared->nlevels = out_nlevels; | |
| 799 | 125 | shared->first_posting = first_posting; | |
| 800 |
1/2✓ Branch 0 taken 125 times.
✗ Branch 1 not taken.
|
125 | if (out_first_posting) |
| 801 | 125 | *out_first_posting = first_posting; | |
| 802 | { | ||
| 803 | 43 | float *dsm_gmean = | |
| 804 | 125 | shm_toc_lookup(pcxt->toc, PRISM_DSM_KEY_GLOBAL_MEAN, false); | |
| 805 | 125 | memcpy(dsm_gmean, global_mean, vec_nbytes); | |
| 806 | } | ||
| 807 | 125 | prism_pbuild_publish_storage(shared, storage); | |
| 808 | |||
| 809 |
1/2✓ Branch 0 taken 125 times.
✗ Branch 1 not taken.
|
125 | if (out_global_mean) |
| 810 | 125 | *out_global_mean = global_mean; /* caller owns it (metadata write) */ | |
| 811 | else | ||
| 812 | ✗ | vs_free(global_mean); | |
| 813 | |||
| 814 | /* No in-RAM tree exists; the caller's metadata write needs only the | ||
| 815 | * shape. */ | ||
| 816 | 125 | *out_nlist = nlist; | |
| 817 | 125 | *out_tree_nlevels = out_nlevels; | |
| 818 | |||
| 819 | /* The refine decision, now that the actual leaf count is known: refine | ||
| 820 | * only when the sample was bounded below the table AND the leaves are | ||
| 821 | * sample-thin -- a leaf's encode reference is a sample mean whose error | ||
| 822 | * shrinks with its sample count, so past the threshold the full-table | ||
| 823 | * scan recomputes what the sample already got right. Published before | ||
| 824 | * the ready barrier; workers gate the refine phase (and its barriers) | ||
| 825 | * on it after that barrier. */ | ||
| 826 | { | ||
| 827 | 125 | uint64_t collected = 0; | |
| 828 | 125 | uint64_t seen = 0; | |
| 829 |
2/2✓ Branch 0 taken 315 times.
✓ Branch 1 taken 125 times.
|
440 | for (int t = 0; t < nparticipants; t++) |
| 830 | { | ||
| 831 | 315 | collected += prism_dsm_sample_counts(dsm_samples)[t]; | |
| 832 | 315 | seen += prism_dsm_sample_seen(dsm_samples)[t]; | |
| 833 | } | ||
| 834 | /* kept < seen means the sampling scans skipped real rows, so the | ||
| 835 | * sample is a strict subset of the table (kept == seen means the | ||
| 836 | * sample IS the table -- nothing to refine from). */ | ||
| 837 |
1/2✓ Branch 0 taken 3 times.
✗ Branch 1 not taken.
|
3 | shared->refine = shared->refine_threshold > 0 && collected < seen && |
| 838 |
4/6✓ Branch 0 taken 3 times.
✓ Branch 1 taken 122 times.
✓ Branch 2 taken 3 times.
✗ Branch 3 not taken.
✓ Branch 4 taken 3 times.
✗ Branch 5 not taken.
|
128 | shared->refine_tile_cap > 0 && nlist > 0 && |
| 839 |
2/2✓ Branch 0 taken 1 times.
✓ Branch 1 taken 2 times.
|
3 | collected / nlist < |
| 840 | ✗ | (uint64_t)shared->refine_threshold; | |
| 841 |
2/5✗ Branch 0 not taken.
✓ Branch 1 taken 82 times.
✓ Branch 2 taken 43 times.
✗ Branch 3 not taken.
✗ Branch 4 not taken.
|
125 | vs_debug( |
| 842 | "prism: refine gate: kept=%" PRIu64 " seen=%" PRIu64 | ||
| 843 | " nlist=%u threshold=%u -> %s", | ||
| 844 | collected, | ||
| 845 | seen, | ||
| 846 | nlist, | ||
| 847 | shared->refine_threshold, | ||
| 848 | shared->refine ? "refine" : "skip"); | ||
| 849 | } | ||
| 850 | |||
| 851 | /* The samples are dead (their last readers were the subtree builders); | ||
| 852 | * lay the refine accumulator over them before the ready barrier so the | ||
| 853 | * workers -- who pass that barrier ahead of the refine phase -- see an | ||
| 854 | * initialized header. */ | ||
| 855 | 125 | PrismDsmRefineAccum *refine_accum = NULL; | |
| 856 |
2/2✓ Branch 0 taken 2 times.
✓ Branch 1 taken 123 times.
|
125 | if (shared->refine) |
| 857 | { | ||
| 858 | 2 | refine_accum = prism_pbuild_refine_overlay(dsm_samples); | |
| 859 | 2 | refine_accum->nleaves = shared->refine_tile_cap; | |
| 860 | 2 | refine_accum->dim = dim; | |
| 861 | } | ||
| 862 | |||
| 863 | /* Initialize the shared cluster sorter for the launched-worker count | ||
| 864 | * BEFORE the ready barrier, so it is ready when workers attach in phase 3. | ||
| 865 | * The leader merges only (it does not sort a share). */ | ||
| 866 | 43 | void *sortshared = | |
| 867 | 125 | shm_toc_lookup(pcxt->toc, PRISM_DSM_KEY_SORTSHARED, false); | |
| 868 | 125 | prism_pbuild_sort_shared_init( | |
| 869 | 125 | sortshared, pcxt->nworkers_launched, pcxt->seg); | |
| 870 | |||
| 871 | /* Barrier: centroid/head pages written + routing state published + sorter | ||
| 872 | * ready; workers build their page-backed router next. */ | ||
| 873 | 125 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 874 | |||
| 875 | /* ---- Phase 2.5: page-backed full-table refine (only when subsampled) | ||
| 876 | * ---- The workers route + accumulate; the leader (participant 0) clears | ||
| 877 | * the tiled accumulator, resets the scan per tile, and rewrites each | ||
| 878 | * leaf's head-page pt_centroid to the full-table mean. Gated on | ||
| 879 | * shared->refine, matching the workers, so the internal barriers stay | ||
| 880 | * in lockstep. */ | ||
| 881 |
2/2✓ Branch 0 taken 2 times.
✓ Branch 1 taken 123 times.
|
125 | if (shared->refine) |
| 882 | { | ||
| 883 | 2 | prism_build_report_phase(prog, PRISM_BUILD_PHASE_REFINE); | |
| 884 | 2 | PrismDsmRefineAccum *accum = refine_accum; | |
| 885 | 2 | PrismHeadWriteCtx rhead; | |
| 886 | 4 | prism_head_write_ctx_init( | |
| 887 | &rhead, | ||
| 888 | storage, | ||
| 889 | rq_params, | ||
| 890 | dim, | ||
| 891 | 2 | shared->fastscan, | |
| 892 | first_posting); | ||
| 893 | 2 | prism_pbuild_exec_refine_paged( | |
| 894 | 0, | ||
| 895 | heap, | ||
| 896 | index, | ||
| 897 | index_info, | ||
| 898 | shared, | ||
| 899 | NULL, | ||
| 900 | first_posting, | ||
| 901 | accum, | ||
| 902 | barrier, | ||
| 903 | prism_write_leaf_head, | ||
| 904 | &rhead); | ||
| 905 | 2 | prism_head_write_ctx_cleanup(&rhead); | |
| 906 | } | ||
| 907 | |||
| 908 | /* Phase 3: workers scan + route page-backed + encode + sort; the leader | ||
| 909 | * merges. Reported exactly once -- the seam fires the "prism-build-load" | ||
| 910 | * test hook, and a 'wait' attached there must pause the build a single | ||
| 911 | * time -- and before the scan-reset barrier below releases the workers, | ||
| 912 | * so progress reflects the whole (multi-hour at scale) scan. */ | ||
| 913 | 125 | prism_build_report_phase(prog, PRISM_BUILD_PHASE_SCAN_PARALLEL); | |
| 914 | |||
| 915 | /* The samples (and the refine overlay riding in them) are dead; hand the | ||
| 916 | * segment back before the posting sort claims its own memory budget. */ | ||
| 917 | 125 | prism_pbuild_samples_release(dsm_samples, lead.sample_seg); | |
| 918 | 125 | dsm_samples = NULL; | |
| 919 | |||
| 920 | /* Re-init the scan for the posting phase (the refine passes above consumed | ||
| 921 | * it). Guarded by the barrier below so no worker scans before the reset. | ||
| 922 | */ | ||
| 923 | 125 | prism_pbuild_rescan(heap, shared); | |
| 924 | |||
| 925 | /* Barrier: scan reset for the posting phase; workers start the page-backed | ||
| 926 | * posting scan+sort. */ | ||
| 927 | 125 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 928 | |||
| 929 | 43 | instr_time t_scan_start; | |
| 930 | 125 | INSTR_TIME_SET_CURRENT(t_scan_start); | |
| 931 | |||
| 932 | /* Barrier: wait until every worker has finished sorting its run, then | ||
| 933 | * merge and build. The leader does not scan; the workers cover the heap. | ||
| 934 | */ | ||
| 935 | 125 | BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | |
| 936 | 125 | BarrierDetach(barrier); | |
| 937 | |||
| 938 | 125 | WaitForParallelWorkersToFinish(pcxt); | |
| 939 | |||
| 940 | 43 | instr_time t_scan_end; | |
| 941 | 125 | INSTR_TIME_SET_CURRENT(t_scan_end); | |
| 942 | 125 | INSTR_TIME_SUBTRACT(t_scan_end, t_scan_start); | |
| 943 | |||
| 944 |
2/2✓ Branch 0 taken 190 times.
✓ Branch 1 taken 125 times.
|
315 | for (int i = 0; i < pcxt->nworkers_launched; i++) |
| 945 | 84 | InstrAccumParallelQuery(&bufferusage[i], &walusage[i]); | |
| 946 | |||
| 947 | 125 | *out_heap_tuples = shared->reltuples; | |
| 948 | 125 | *out_indtuples = shared->indtuples; | |
| 949 | 125 | *out_soar_dupes = shared->soar_dupes; | |
| 950 | |||
| 951 | /* The workers cover the heap cooperatively while the leader blocks on the | ||
| 952 | * scan barrier above, so there is no leader-side loop to advance the % mid | ||
| 953 | * scan; publish the final scanned count now that the workers have | ||
| 954 | * finished. | ||
| 955 | */ | ||
| 956 | 125 | prism_build_report_progress(prog, shared->indtuples); | |
| 957 | |||
| 958 | 43 | instr_time t_merge_start; | |
| 959 | 125 | INSTR_TIME_SET_CURRENT(t_merge_start); | |
| 960 | |||
| 961 | /* | ||
| 962 | * Merge the workers' sorted runs and build each cluster's posting list | ||
| 963 | * with a single resident page builder, in cluster order (same loop as the | ||
| 964 | * serial path). Entries arrive grouped by cluster, so there is no | ||
| 965 | * partial-page fold and no chain to splice. | ||
| 966 | */ | ||
| 967 | 125 | prism_build_report_phase(prog, PRISM_BUILD_PHASE_POSTING); | |
| 968 | 293 | PrismSorter *sorter = prism_pbuild_sort_begin( | |
| 969 | sortshared, | ||
| 970 | 125 | pcxt->seg, | |
| 971 | 0, | ||
| 972 | pcxt->nworkers_launched, | ||
| 973 | true, | ||
| 974 | 125 | (uint32_t)prism_posting_entry_size(dim), | |
| 975 | shared->work_mem_kb); | ||
| 976 | 125 | prism_posting_build_lists( | |
| 977 | sorter, | ||
| 978 | storage, | ||
| 979 | nlist, | ||
| 980 | dim, | ||
| 981 | 125 | shared->fastscan, | |
| 982 | rq_params, | ||
| 983 | first_posting); | ||
| 984 | |||
| 985 | 125 | uint32_t total_pages = RelationGetNumberOfBlocks(index) - first_posting; | |
| 986 | |||
| 987 | 43 | instr_time t_merge_end; | |
| 988 | 125 | INSTR_TIME_SET_CURRENT(t_merge_end); | |
| 989 | 125 | INSTR_TIME_SUBTRACT(t_merge_end, t_merge_start); | |
| 990 | |||
| 991 |
1/2✗ Branch 1 not taken.
✓ Branch 2 taken 43 times.
|
125 | vs_debug( |
| 992 | "prism: parallel streaming build with %d workers, " | ||
| 993 | "%u clusters, %u pages, " | ||
| 994 | "scan+drain %.1fms, finalize %.1fms", | ||
| 995 | pcxt->nworkers_launched, | ||
| 996 | nlist, | ||
| 997 | total_pages, | ||
| 998 | INSTR_TIME_GET_MILLISEC(t_scan_end), | ||
| 999 | INSTR_TIME_GET_MILLISEC(t_merge_end)); | ||
| 1000 | |||
| 1001 | /* Workers detached from the exact-centroid collection when routing ended; | ||
| 1002 | * the leader's release is the last and frees it. */ | ||
| 1003 | 125 | prism_pbuild_exact_centroids_release(exact_seg); | |
| 1004 | |||
| 1005 | 125 | prism_pbuild_teardown(pcxt); | |
| 1006 | |||
| 1007 | 125 | return true; | |
| 1008 | } | ||
| 1009 |