| 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_backend.c - PostgreSQL implementations of the build's back-end | ||
| 6 | * seams. | ||
| 7 | * | ||
| 8 | * The build driver is shared between the PG extension and the standalone | ||
| 9 | * engine; the parts that genuinely differ by back-end are expressed as | ||
| 10 | * same-named seam functions, with the PG versions here and the thread-based | ||
| 11 | * versions under src/standalone. It holds the PG implementations of all the | ||
| 12 | * build seams: DSM setup over shm_toc, ParallelContext launch/teardown, worker | ||
| 13 | * attach/detach, the heap parallel scan, and per-worker count reduction. | ||
| 14 | */ | ||
| 15 | |||
| 16 | #include <postgres.h> | ||
| 17 | |||
| 18 | #include "vs_config.h" | ||
| 19 | |||
| 20 | #include <access/parallel.h> | ||
| 21 | #include <access/table.h> | ||
| 22 | #include <access/tableam.h> | ||
| 23 | #include <catalog/index.h> | ||
| 24 | #include <catalog/pg_operator_d.h> | ||
| 25 | #include <catalog/pg_type_d.h> | ||
| 26 | #include <executor/instrument.h> | ||
| 27 | #include <executor/tuptable.h> | ||
| 28 | #include <miscadmin.h> | ||
| 29 | #include <optimizer/plancat.h> | ||
| 30 | #include <pgstat.h> | ||
| 31 | #include <storage/buffile.h> | ||
| 32 | #include <storage/latch.h> | ||
| 33 | #include <storage/proc.h> | ||
| 34 | #include <storage/spin.h> | ||
| 35 | #include <tcop/tcopprot.h> | ||
| 36 | #include <utils/rel.h> | ||
| 37 | #include <utils/snapmgr.h> | ||
| 38 | #include <utils/tuplesort.h> | ||
| 39 | #include <utils/wait_event.h> | ||
| 40 | |||
| 41 | #include "build.h" | ||
| 42 | #include "index/index_build.h" | ||
| 43 | #include "index/parallel_build.h" | ||
| 44 | #include "pg/bufstorage.h" | ||
| 45 | #include "quant/matrix.h" | ||
| 46 | #include "quant/rabitq.h" | ||
| 47 | #include "support_pg.h" | ||
| 48 | #include "typeinfo.h" | ||
| 49 | #include "types/vec32.h" | ||
| 50 | |||
| 51 | /* | ||
| 52 | * PG-specific shared build state: the neutral PrismBuildShared plus the | ||
| 53 | * relation identity, query id, and the spinlock guarding its counters. A | ||
| 54 | * ParallelTableScanDesc is appended after it in the DSM segment. The base is | ||
| 55 | * the first member, so the PrismBuildShared * the workers look up out of the | ||
| 56 | * toc is recovered here as a PrismBuildSharedPg *. | ||
| 57 | */ | ||
| 58 | typedef struct PrismBuildSharedPg | ||
| 59 | { | ||
| 60 | PrismBuildShared base; | ||
| 61 | Oid heaprelid; | ||
| 62 | Oid indexrelid; | ||
| 63 | int64 queryid; | ||
| 64 | /* DSM handle of the dedicated sample segment (sample-region seam). */ | ||
| 65 | dsm_handle sample_handle; | ||
| 66 | /* DSM handle of the per-child subtree ring, created by the leader after | ||
| 67 | * root assignment (when the per-child sample counts that bound the slot | ||
| 68 | * size are known) and published before the ring barrier. */ | ||
| 69 | dsm_handle subtree_ring_handle; | ||
| 70 | /* DSM handle of the exact centroid collection, created by the leader | ||
| 71 | * after the streaming tree write and published before the tree-ready | ||
| 72 | * barrier (exact-centroid seam). */ | ||
| 73 | dsm_handle exact_centroids_handle; | ||
| 74 | slock_t mutex; | ||
| 75 | /* Striped locks guarding the shared leaf-refinement accumulator. */ | ||
| 76 | slock_t accum_locks[PRISM_REFINE_LOCK_STRIPES]; | ||
| 77 | } PrismBuildSharedPg; | ||
| 78 | |||
| 79 | #define ParallelTableScanFromVsShared(shared) \ | ||
| 80 | ((ParallelTableScanDesc)((char *)(shared) + \ | ||
| 81 | BUFFERALIGN(sizeof(PrismBuildSharedPg)))) | ||
| 82 | |||
| 83 | /* | ||
| 84 | * Bridges PostgreSQL's heap-tuple callback to the back-end-neutral scan | ||
| 85 | * callback: skip nulls and hand the vector's data pointer (into the varlena, | ||
| 86 | * no copy) plus its tid to the shared logic. | ||
| 87 | * | ||
| 88 | * A halfvec column cannot be handed over no-copy -- the shared logic is | ||
| 89 | * float32 -- so it is widened into buf, reused across tuples. The callback | ||
| 90 | * consumes the pointer before returning, so one buffer per participant is | ||
| 91 | * enough. | ||
| 92 | */ | ||
| 93 | typedef struct PrismPgScanAdapter | ||
| 94 | { | ||
| 95 | PrismBuildScanCb cb; | ||
| 96 | void *state; | ||
| 97 | Vec32Access input; | ||
| 98 | } PrismPgScanAdapter; | ||
| 99 | |||
| 100 | static void | ||
| 101 | 230560 | prism_pg_scan_adapter( | |
| 102 | Relation index, | ||
| 103 | ItemPointer tid, | ||
| 104 | Datum *values, | ||
| 105 | bool *isnull, | ||
| 106 | bool tuple_is_alive, | ||
| 107 | void *adapter_state) | ||
| 108 | { | ||
| 109 | 230560 | PrismPgScanAdapter *a = (PrismPgScanAdapter *)adapter_state; | |
| 110 | |||
| 111 | 230560 | (void)index; | |
| 112 | 230560 | (void)tuple_is_alive; | |
| 113 | |||
| 114 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 230560 times.
|
230560 | if (isnull[0]) |
| 115 | ✗ | return; | |
| 116 | |||
| 117 | 230560 | Vec32Ref vref = vec32_read(&a->input, values[0]); | |
| 118 | 230560 | a->cb(a->state, *tid, vref.data); | |
| 119 | } | ||
| 120 | |||
| 121 | /* | ||
| 122 | * Scan every vector via the shared parallel table scan, invoking cb per live | ||
| 123 | * tuple. The standalone back-end provides a same-named function that iterates | ||
| 124 | * its in-memory vector array (work-stealing) instead. allow_sync/progress | ||
| 125 | * are PostgreSQL table_index_build_scan flags, ignored in standalone; progress | ||
| 126 | * gates pg_stat_progress_create_index reporting, so only the leader sets it | ||
| 127 | * (otherwise every participant would inflate the tuples-scanned counter). | ||
| 128 | */ | ||
| 129 | double | ||
| 130 | 218 | prism_build_scan( | |
| 131 | Relation heap, | ||
| 132 | Relation index, | ||
| 133 | struct IndexInfo *indexInfo, | ||
| 134 | PrismBuildShared *shared, | ||
| 135 | bool allow_sync, | ||
| 136 | bool progress, | ||
| 137 | PrismBuildScanCb cb, | ||
| 138 | void *state) | ||
| 139 | { | ||
| 140 | 218 | Dimension dim = (Dimension)TupleDescAttr(index->rd_att, 0)->atttypmod; | |
| 141 | /* Build-time path: resolved directly, since the per-backend cache reads a | ||
| 142 | * metadata page this build has not written yet. */ | ||
| 143 | 218 | PrismPgScanAdapter actx = { | |
| 144 | .cb = cb, | ||
| 145 | .state = state, | ||
| 146 | 218 | .input = vec32_access( | |
| 147 | prism_index_type_info(index), dim, CurrentMemoryContext), | ||
| 148 | }; | ||
| 149 | |||
| 150 | #if PG_VERSION_NUM >= 190000 | ||
| 151 | TableScanDesc scan = table_beginscan_parallel( | ||
| 152 | heap, ParallelTableScanFromVsShared(shared), SO_NONE); | ||
| 153 | #else | ||
| 154 | 218 | TableScanDesc scan = table_beginscan_parallel( | |
| 155 | heap, ParallelTableScanFromVsShared(shared)); | ||
| 156 | #endif | ||
| 157 | |||
| 158 | 218 | return table_index_build_scan( | |
| 159 | heap, | ||
| 160 | index, | ||
| 161 | indexInfo, | ||
| 162 | allow_sync, | ||
| 163 | progress, | ||
| 164 | prism_pg_scan_adapter, | ||
| 165 | &actx, | ||
| 166 | scan); | ||
| 167 | } | ||
| 168 | |||
| 169 | /* | ||
| 170 | * Relation lock modes a parallel worker uses for the heap and index. A | ||
| 171 | * concurrent build (CREATE INDEX CONCURRENTLY) must take weak locks: the | ||
| 172 | * leader holds only ShareUpdateExclusive on the heap, and acquiring a strong | ||
| 173 | * index lock in a worker logs it for hot-standby and assigns an XID, which is | ||
| 174 | * illegal in a parallel worker. Matches PostgreSQL's btree parallel build. The | ||
| 175 | * same pair is used to open the relations on attach and to close them on | ||
| 176 | * detach. | ||
| 177 | */ | ||
| 178 | static void | ||
| 179 | 170 | worker_lockmodes(bool concurrent, LOCKMODE *heapmode, LOCKMODE *indexmode) | |
| 180 | { | ||
| 181 | 506 | *heapmode = concurrent ? ShareUpdateExclusiveLock : ShareLock; | |
| 182 | 166 | *indexmode = concurrent ? RowExclusiveLock : AccessExclusiveLock; | |
| 183 | } | ||
| 184 | |||
| 185 | /* Per-worker build memory context (see prism_pbuild_worker_attach). */ | ||
| 186 | static MemoryContext vs_pbuild_worker_ctx = NULL; | ||
| 187 | static MemoryContext vs_pbuild_worker_oldctx = NULL; | ||
| 188 | |||
| 189 | /* | ||
| 190 | * Join the parallel build: look up the shared state, open the heap and index, | ||
| 191 | * start per-worker instrumentation, and attach to the phase barrier. The | ||
| 192 | * standalone back-end provides a same-named function that takes the shared | ||
| 193 | * state and vectors directly and joins a thread barrier. | ||
| 194 | */ | ||
| 195 | void | ||
| 196 | 86 | prism_pbuild_worker_attach(shm_toc *toc, PrismPBuildWorker *w) | |
| 197 | { | ||
| 198 | 86 | PrismBuildShared *shared = | |
| 199 | 86 | shm_toc_lookup(toc, PRISM_DSM_KEY_SHARED, false); | |
| 200 | 86 | PrismBuildSharedPg *pg = (PrismBuildSharedPg *)shared; | |
| 201 | 86 | Barrier *barrier = shm_toc_lookup(toc, PRISM_DSM_KEY_BARRIER, false); | |
| 202 | |||
| 203 | 86 | char *sharedquery = shm_toc_lookup(toc, PRISM_DSM_KEY_QUERY_TEXT, true); | |
| 204 | 86 | debug_query_string = sharedquery; | |
| 205 | 86 | pgstat_report_activity(STATE_RUNNING, debug_query_string); | |
| 206 | 86 | pgstat_report_query_id(pg->queryid, false); | |
| 207 | |||
| 208 | 86 | LOCKMODE heapmode, indexmode; | |
| 209 |
2/2✓ Branch 0 taken 84 times.
✓ Branch 1 taken 2 times.
|
86 | worker_lockmodes(shared->concurrent, &heapmode, &indexmode); |
| 210 | |||
| 211 | 86 | w->shared = shared; | |
| 212 | 86 | w->barrier = barrier; | |
| 213 | 86 | w->heapRel = table_open(pg->heaprelid, heapmode); | |
| 214 | 86 | w->indexRel = index_open(pg->indexrelid, indexmode); | |
| 215 | 86 | w->worker_id = ParallelWorkerNumber + 1; | |
| 216 | 86 | w->dim = shared->dim; | |
| 217 | |||
| 218 | /* | ||
| 219 | * All of this worker's build allocations (routing state, per-child | ||
| 220 | * subtree blobs, batch buffers) go into a named context so | ||
| 221 | * pg_backend_memory_contexts attributes them to the build and they are | ||
| 222 | * reclaimed together at detach. A file-static is safe for the same | ||
| 223 | * reason as vs_pbuild_snapshot: one build per worker, non-reentrant. | ||
| 224 | */ | ||
| 225 | 86 | vs_pbuild_worker_ctx = AllocSetContextCreate( | |
| 226 | CurrentMemoryContext, "vs worker build", ALLOCSET_DEFAULT_SIZES); | ||
| 227 | 86 | vs_pbuild_worker_oldctx = MemoryContextSwitchTo(vs_pbuild_worker_ctx); | |
| 228 | |||
| 229 | 86 | InstrStartParallelQuery(); | |
| 230 | |||
| 231 | /* | ||
| 232 | * Attach to the dynamic phase barrier (before the first phase). The leader | ||
| 233 | * attaches before launching workers, so the party tracks all participants | ||
| 234 | * that actually start and stays in lockstep regardless of how many workers | ||
| 235 | * PostgreSQL launched. | ||
| 236 | */ | ||
| 237 | 86 | BarrierAttach(barrier); | |
| 238 | 86 | } | |
| 239 | |||
| 240 | /* | ||
| 241 | * Leave the parallel build: report this worker's buffer/WAL usage back to the | ||
| 242 | * leader and close the relations. The standalone back-end's same-named | ||
| 243 | * function joins the thread and is otherwise a no-op. | ||
| 244 | */ | ||
| 245 | void | ||
| 246 | 84 | prism_pbuild_worker_detach(shm_toc *toc, PrismPBuildWorker *w) | |
| 247 | { | ||
| 248 | 84 | MemoryContextSwitchTo(vs_pbuild_worker_oldctx); | |
| 249 | 84 | MemoryContextDelete(vs_pbuild_worker_ctx); | |
| 250 | 84 | vs_pbuild_worker_ctx = NULL; | |
| 251 | 84 | vs_pbuild_worker_oldctx = NULL; | |
| 252 | |||
| 253 | 84 | BufferUsage *bufferusage = | |
| 254 | 84 | shm_toc_lookup(toc, PRISM_DSM_KEY_BUFFER_USAGE, false); | |
| 255 | 84 | WalUsage *walusage = shm_toc_lookup(toc, PRISM_DSM_KEY_WAL_USAGE, false); | |
| 256 | 84 | InstrEndParallelQuery( | |
| 257 | 84 | &bufferusage[ParallelWorkerNumber], | |
| 258 | 84 | &walusage[ParallelWorkerNumber]); | |
| 259 | |||
| 260 | 84 | LOCKMODE heapmode, indexmode; | |
| 261 |
2/2✓ Branch 0 taken 82 times.
✓ Branch 1 taken 2 times.
|
84 | worker_lockmodes(w->shared->concurrent, &heapmode, &indexmode); |
| 262 | 84 | index_close(w->indexRel, indexmode); | |
| 263 | 84 | table_close(w->heapRel, heapmode); | |
| 264 | 84 | } | |
| 265 | |||
| 266 | /* | ||
| 267 | * Page-backed routing storage seam (see parallel_build.h). PG workers are | ||
| 268 | * separate processes, so each opens its own VsStorage on the worker's index | ||
| 269 | * relation; the leader's storage pointer cannot cross the process boundary, so | ||
| 270 | * publish is a no-op here. | ||
| 271 | */ | ||
| 272 | void | ||
| 273 | 43 | prism_pbuild_publish_storage(PrismBuildShared *shared, VsStorage *s) | |
| 274 | { | ||
| 275 | 43 | (void)shared; | |
| 276 | 43 | (void)s; | |
| 277 | 43 | } | |
| 278 | |||
| 279 | VsStorage * | ||
| 280 | 84 | prism_pbuild_worker_storage(PrismPBuildWorker *w) | |
| 281 | { | ||
| 282 | /* No table relation needed (routing reads index pages only, no rerank). */ | ||
| 283 | 84 | VsPgStorage *s = palloc(sizeof(VsPgStorage)); | |
| 284 | 84 | vs_pg_storage_init(s, w->indexRel, NULL, w->shared->metric); | |
| 285 | 84 | s->build_mode = true; /* reads only; matches the leader's build storage */ | |
| 286 | 84 | return &s->base; /* base is the first member */ | |
| 287 | } | ||
| 288 | |||
| 289 | void | ||
| 290 | 84 | prism_pbuild_worker_storage_release(VsStorage *s) | |
| 291 | { | ||
| 292 | /* The route helper releases every page it reads, so no buffer stays | ||
| 293 | * pinned; just free the wrapper (allocated in the worker's memory | ||
| 294 | * context). */ | ||
| 295 | 84 | pfree(s); | |
| 296 | 84 | } | |
| 297 | |||
| 298 | /* | ||
| 299 | * Snapshot used to initialize the parallel heap scan for CREATE INDEX | ||
| 300 | * CONCURRENTLY (an MVCC snapshot; SnapshotAny/NULL for a normal build). It | ||
| 301 | * must stay registered for the whole parallel operation and is released at | ||
| 302 | * teardown. A file-static is safe: an index build is single-threaded and | ||
| 303 | * non-reentrant in the leader backend, and every parallel-build exit path runs | ||
| 304 | * prism_pbuild_teardown. | ||
| 305 | */ | ||
| 306 | static Snapshot vs_pbuild_snapshot = NULL; | ||
| 307 | |||
| 308 | /* | ||
| 309 | * Tear the parallel context down and leave parallel mode. The standalone | ||
| 310 | * back-end provides a same-named function that joins its worker threads and | ||
| 311 | * frees the shared arena instead. | ||
| 312 | */ | ||
| 313 | void | ||
| 314 | 44 | prism_pbuild_teardown(ParallelContext *pcxt) | |
| 315 | { | ||
| 316 |
2/2✓ Branch 0 taken 1 times.
✓ Branch 1 taken 43 times.
|
44 | if (vs_pbuild_snapshot != NULL) |
| 317 | { | ||
| 318 | 1 | UnregisterSnapshot(vs_pbuild_snapshot); | |
| 319 | 1 | vs_pbuild_snapshot = NULL; | |
| 320 | } | ||
| 321 | 44 | DestroyParallelContext(pcxt); | |
| 322 | 44 | ExitParallelMode(); | |
| 323 | 44 | } | |
| 324 | |||
| 325 | /* | ||
| 326 | * Leader-only subtree blob store (see parallel_build.h): a BufFile temp | ||
| 327 | * file. Small blob sets never leave the kernel page cache; large ones spill | ||
| 328 | * to pgsql_tmp automatically, so the store adds no unbounded memory. Blobs | ||
| 329 | * are length-prefixed and read back strictly in append order. | ||
| 330 | */ | ||
| 331 | struct PrismBlobStore | ||
| 332 | { | ||
| 333 | BufFile *file; | ||
| 334 | }; | ||
| 335 | |||
| 336 | PrismBlobStore * | ||
| 337 | 170 | prism_pbuild_blobstore_begin(void) | |
| 338 | { | ||
| 339 | 170 | PrismBlobStore *bs = palloc(sizeof(PrismBlobStore)); | |
| 340 | 170 | bs->file = BufFileCreateTemp(false); | |
| 341 | 170 | return bs; | |
| 342 | } | ||
| 343 | |||
| 344 | void | ||
| 345 | 3771 | prism_pbuild_blobstore_put(PrismBlobStore *bs, const void *blob, uint64_t size) | |
| 346 | { | ||
| 347 | 3771 | BufFileWrite(bs->file, &size, sizeof(size)); | |
| 348 | 3771 | BufFileWrite(bs->file, blob, (size_t)size); | |
| 349 | 3771 | } | |
| 350 | |||
| 351 | void | ||
| 352 | 170 | prism_pbuild_blobstore_rewind(PrismBlobStore *bs) | |
| 353 | { | ||
| 354 |
1/2✗ Branch 1 not taken.
✓ Branch 2 taken 170 times.
|
170 | if (BufFileSeek(bs->file, 0, 0, SEEK_SET) != 0) |
| 355 | ✗ | ereport(ERROR, | |
| 356 | (errcode_for_file_access(), | ||
| 357 | errmsg("could not rewind subtree blob store"))); | ||
| 358 | 170 | } | |
| 359 | |||
| 360 | uint64_t | ||
| 361 | 3771 | prism_pbuild_blobstore_get(PrismBlobStore *bs, void *buf, uint64_t max_size) | |
| 362 | { | ||
| 363 | 3771 | uint64_t size; | |
| 364 | 3771 | BufFileReadExact(bs->file, &size, sizeof(size)); | |
| 365 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 3771 times.
|
3771 | if (size > max_size) |
| 366 | ✗ | ereport(ERROR, | |
| 367 | (errcode(ERRCODE_INTERNAL_ERROR), | ||
| 368 | errmsg("subtree blob larger than its slot (%llu > %llu)", | ||
| 369 | (unsigned long long)size, | ||
| 370 | (unsigned long long)max_size))); | ||
| 371 | 3771 | BufFileReadExact(bs->file, buf, (size_t)size); | |
| 372 | 3771 | return size; | |
| 373 | } | ||
| 374 | |||
| 375 | void | ||
| 376 | 170 | prism_pbuild_blobstore_end(PrismBlobStore *bs) | |
| 377 | { | ||
| 378 | 170 | BufFileClose(bs->file); | |
| 379 | 170 | pfree(bs); | |
| 380 | 170 | } | |
| 381 | |||
| 382 | /* | ||
| 383 | * Launch the worker participants and wait until they have all attached to the | ||
| 384 | * barrier (so the dynamic party reaches launched+1 before the leader advances | ||
| 385 | * the first phase). Returns false — after tearing the context down — if no | ||
| 386 | * workers started, so the caller falls back to a serial build. The standalone | ||
| 387 | * back-end provides a same-named function that spawns threads and joins them | ||
| 388 | * at the barrier instead. | ||
| 389 | */ | ||
| 390 | bool | ||
| 391 | 45 | prism_pbuild_launch( | |
| 392 | ParallelContext *pcxt, Barrier *barrier, PrismBuildShared *shared) | ||
| 393 | { | ||
| 394 | 45 | LaunchParallelWorkers(pcxt); | |
| 395 | |||
| 396 |
2/2✓ Branch 0 taken 1 times.
✓ Branch 1 taken 44 times.
|
45 | if (pcxt->nworkers_launched == 0) |
| 397 | { | ||
| 398 | 1 | WaitForParallelWorkersToFinish(pcxt); | |
| 399 | 1 | prism_pbuild_teardown(pcxt); | |
| 400 | 1 | return false; | |
| 401 | } | ||
| 402 | |||
| 403 | /* | ||
| 404 | * Workers attach to the barrier dynamically, so the party (1 leader + N | ||
| 405 | * launched) is not final until they all have; if the leader arrived first | ||
| 406 | * it could advance the phase alone and strand late workers. | ||
| 407 | * WaitForParallelWorkersToAttach surfaces a startup failure as an error | ||
| 408 | * rather than a hang; then poll until the live participant count is whole. | ||
| 409 | */ | ||
| 410 | 44 | WaitForParallelWorkersToAttach(pcxt); | |
| 411 |
2/2✓ Branch 2 taken 300 times.
✓ Branch 3 taken 44 times.
|
344 | while (BarrierParticipants(barrier) < pcxt->nworkers_launched + 1) |
| 412 | { | ||
| 413 |
2/2✓ Branch 0 taken 4 times.
✓ Branch 1 taken 296 times.
|
300 | CHECK_FOR_INTERRUPTS(); |
| 414 | 300 | (void)WaitLatch( | |
| 415 | MyLatch, | ||
| 416 | WL_LATCH_SET | WL_TIMEOUT | WL_EXIT_ON_PM_DEATH, | ||
| 417 | 1L, | ||
| 418 | WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN); | ||
| 419 | 300 | ResetLatch(MyLatch); | |
| 420 | } | ||
| 421 | |||
| 422 | /* The launch can fall short of the plan (the parallel-worker pool under | ||
| 423 | * max_parallel_workers is shared with concurrent queries). Narrow the | ||
| 424 | * participant count to the party that attached so the phases partition | ||
| 425 | * their work over participants that exist; the per-participant DSM | ||
| 426 | * regions keep their planned size and leave the tail slots unused. */ | ||
| 427 | 44 | shared->nparticipants = pcxt->nworkers_launched + 1; | |
| 428 | 44 | return true; | |
| 429 | } | ||
| 430 | |||
| 431 | /* | ||
| 432 | * Allocate and populate the parallel build's shared state: the DSM segment and | ||
| 433 | * its regions (shared header, barrier, sample/centroid/assignment slots, the | ||
| 434 | * tree blob, the per-worker page queues, usage counters). Does not launch | ||
| 435 | * workers; the caller does. Coarse PG block — the standalone back-end provides | ||
| 436 | * a same-named function over a heap arena. Returns false (after tearing the | ||
| 437 | * parallel context down) if the DSM segment could not be created. | ||
| 438 | */ | ||
| 439 | bool | ||
| 440 | 45 | prism_pbuild_setup_shared( | |
| 441 | PrismPBuildLeader *lead, | ||
| 442 | Relation heap, | ||
| 443 | Relation index, | ||
| 444 | const PrismBuildConfig *config, | ||
| 445 | int nworkers) | ||
| 446 | { | ||
| 447 | 45 | Dimension dim = config->dim; | |
| 448 | 45 | uint32_t nlist = config->nlist; | |
| 449 | 45 | int nparticipants = nworkers + 1; | |
| 450 | 45 | uint64_t rabitq_seed = VS_RABITQ_BUILD_SEED; | |
| 451 | 90 | uint32_t fan_out = config->fan_out > 0 ? config->fan_out | |
| 452 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 45 times.
|
45 | : prism_auto_fan_out(0, nlist, 0); |
| 453 | 45 | uint32_t km_k = fan_out < nlist ? fan_out : nlist; | |
| 454 | |||
| 455 | /* | ||
| 456 | * Size the k-means sample set. The samples live in one shared-memory | ||
| 457 | * region (total_samples * dim floats) that must stay resident for the | ||
| 458 | * whole tree build -- root k-means and every subtree -- so its size is the | ||
| 459 | * build's dominant memory cost. The ideal is ~256 samples per list, but at | ||
| 460 | * fine nlist that can dwarf available RAM (nlist=480k -> 123M samples -> | ||
| 461 | * ~360 GB), far beyond what a DSM segment can hold. | ||
| 462 | * | ||
| 463 | * Bound it by maintenance_work_mem: that is the build's memory budget and | ||
| 464 | * the knob operators already raise for large index builds. When the budget | ||
| 465 | * is smaller than the ideal the stride sampler simply draws a coarser (but | ||
| 466 | * still uniform) subsample to fit. estimate_rel_size() caps it to the rows | ||
| 467 | * that actually exist so small tables don't over-allocate; it is only a | ||
| 468 | * hint -- the budget is the hard bound, so an inaccurate estimate cannot | ||
| 469 | * over-commit shared memory. | ||
| 470 | */ | ||
| 471 | 45 | uint64_t want_samples = (uint64_t)nlist * 256; | |
| 472 | |||
| 473 | /* maintenance_work_mem is in kB; reserve it for the sample region. */ | ||
| 474 | 45 | uint64_t mem_bytes = (uint64_t)maintenance_work_mem * UINT64CONST(1024); | |
| 475 | 45 | uint64_t vec_nbytes = (uint64_t)dim * sizeof(float); | |
| 476 |
1/2✓ Branch 0 taken 45 times.
✗ Branch 1 not taken.
|
45 | uint64_t budget = vec_nbytes > 0 ? mem_bytes / vec_nbytes : want_samples; |
| 477 | 45 | if (budget < 10000) | |
| 478 | budget = 10000; /* k-means needs a workable minimum */ | ||
| 479 | |||
| 480 | 45 | uint64_t total64 = Min(want_samples, budget); | |
| 481 | |||
| 482 | 45 | BlockNumber est_pages; | |
| 483 | 45 | double est_tuples; | |
| 484 | 45 | double allvisfrac; | |
| 485 | 45 | estimate_rel_size(heap, NULL, &est_pages, &est_tuples, &allvisfrac); | |
| 486 |
3/4✓ Branch 0 taken 45 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 38 times.
✓ Branch 3 taken 7 times.
|
45 | if (est_tuples > 0.0 && (double)total64 > est_tuples) |
| 487 | 38 | total64 = (uint64_t)est_tuples; | |
| 488 | |||
| 489 | /* total64 <= want_samples = nlist*256 <= 512M (nlist reloption max 2M), so | ||
| 490 | * it always fits a uint32. */ | ||
| 491 | 45 | uint32_t total_samples = (uint32_t)Max(total64, UINT64CONST(1)); | |
| 492 | |||
| 493 |
3/4✓ Branch 0 taken 2 times.
✓ Branch 1 taken 43 times.
✓ Branch 2 taken 2 times.
✗ Branch 3 not taken.
|
45 | if (want_samples > budget && est_tuples > (double)budget) |
| 494 |
1/2✓ Branch 1 taken 2 times.
✗ Branch 2 not taken.
|
2 | elog(LOG, |
| 495 | VS_EXTENSION_NAME | ||
| 496 | ": k-means sample set limited to %u of the " | ||
| 497 | "ideal " UINT64_FORMAT " vectors by maintenance_work_mem " | ||
| 498 | "(%d kB); raise maintenance_work_mem for finer centroid " | ||
| 499 | "training on large tables", | ||
| 500 | total_samples, | ||
| 501 | want_samples, | ||
| 502 | maintenance_work_mem); | ||
| 503 | |||
| 504 | 45 | uint32_t max_per_worker = (total_samples + nparticipants - 1) / | |
| 505 | nparticipants; | ||
| 506 | |||
| 507 | /* | ||
| 508 | * The refine decision itself is the leader's, made after clustering when | ||
| 509 | * the actual leaf count is known (see PrismBuildShared.refine); setup only | ||
| 510 | * publishes the gate inputs. The pass is single-shot by design: refine | ||
| 511 | * routes every row over the centroid PAGES, which it never rewrites (it | ||
| 512 | * updates the heads' encode references, which routing does not read), so | ||
| 513 | * the row-to-leaf assignment is a fixed point -- a second pass would | ||
| 514 | * rescan the whole table to recompute the exact same means. | ||
| 515 | */ | ||
| 516 | |||
| 517 | 45 | EnterParallelMode(); | |
| 518 | |||
| 519 | 45 | ParallelContext *pcxt = CreateParallelContext( | |
| 520 | VS_MODULE_NAME, "prism_parallel_build_main", nworkers); | ||
| 521 | |||
| 522 | /* | ||
| 523 | * The heap scan's snapshot. A normal build sees all tuples (SnapshotAny); | ||
| 524 | * CREATE INDEX CONCURRENTLY must use an MVCC snapshot so it indexes only | ||
| 525 | * tuples visible to it (heapam asserts SnapshotAny <-> a valid OldestXmin, | ||
| 526 | * so the concurrent path must not pass SnapshotAny). Register it for the | ||
| 527 | * duration — its serialized size also affects the DSM estimate below — and | ||
| 528 | * release it in prism_pbuild_teardown. Mirrors PostgreSQL's nbtsort.c. | ||
| 529 | */ | ||
| 530 | 90 | Snapshot snapshot = config->concurrent | |
| 531 | 1 | ? RegisterSnapshot(GetTransactionSnapshot()) | |
| 532 |
2/2✓ Branch 0 taken 1 times.
✓ Branch 1 taken 44 times.
|
45 | : SnapshotAny; |
| 533 |
2/2✓ Branch 0 taken 44 times.
✓ Branch 1 taken 1 times.
|
45 | vs_pbuild_snapshot = (snapshot != SnapshotAny) ? snapshot : NULL; |
| 534 | 45 | Size est_shared = add_size( | |
| 535 | BUFFERALIGN(sizeof(PrismBuildSharedPg)), | ||
| 536 | table_parallelscan_estimate(heap, snapshot)); | ||
| 537 | |||
| 538 | 45 | shm_toc_estimate_chunk(&pcxt->estimator, est_shared); | |
| 539 | 45 | shm_toc_estimate_chunk(&pcxt->estimator, sizeof(Barrier)); | |
| 540 | /* K-means shared centroids + norms (root level, k=km_k) */ | ||
| 541 | 45 | shm_toc_estimate_chunk( | |
| 542 | &pcxt->estimator, prism_dsm_centroids_size(km_k, dim)); | ||
| 543 | /* K-means per-worker accumulators */ | ||
| 544 | 45 | shm_toc_estimate_chunk( | |
| 545 | &pcxt->estimator, | ||
| 546 | prism_dsm_km_workers_size(nparticipants, km_k, dim)); | ||
| 547 | /* Root assignments: per-worker uint32_t[max_per_worker] */ | ||
| 548 | 45 | shm_toc_estimate_chunk( | |
| 549 | &pcxt->estimator, | ||
| 550 | prism_dsm_root_assign_size(nparticipants, max_per_worker)); | ||
| 551 | /* Shared coordinator for the cluster-keyed posting sort (sort seam). Sized | ||
| 552 | * for the planned participant count (upper bound on launched workers). */ | ||
| 553 | 45 | shm_toc_estimate_chunk( | |
| 554 | &pcxt->estimator, prism_pbuild_sort_shared_size(nparticipants)); | ||
| 555 | |||
| 556 | /* Page-backed routing: the global mean, published by the leader before the | ||
| 557 | * tree-ready barrier. The posting-head base (leaf c's head = first_posting | ||
| 558 | * + c) is a scalar in PrismBuildShared, so no O(nlist) head array is | ||
| 559 | * shared. | ||
| 560 | */ | ||
| 561 | 45 | shm_toc_estimate_chunk(&pcxt->estimator, (Size)vec_nbytes); | |
| 562 | |||
| 563 | /* The leaf-refinement accumulator overlays the sample segment once the | ||
| 564 | * samples are dead (see the sample-region seam), so it needs no DSM chunk | ||
| 565 | * of its own; its tile capacity is bounded by the memory budget, by | ||
| 566 | * MaxAllocSize, and by the region it overlays. */ | ||
| 567 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 45 times.
|
45 | Size samp_sz = prism_dsm_samples_size(nparticipants, max_per_worker, dim); |
| 568 | 45 | uint64_t refine_cap_bytes = | |
| 569 | 45 | Min((uint64_t)maintenance_work_mem * 1024, (uint64_t)MaxAllocSize); | |
| 570 | 45 | if (refine_cap_bytes > | |
| 571 | 45 | (uint64_t)samp_sz - offsetof(PrismDsmRefineAccum, sums)) | |
| 572 | refine_cap_bytes = (uint64_t)samp_sz - | ||
| 573 | offsetof(PrismDsmRefineAccum, sums); | ||
| 574 | 45 | uint32_t refine_tile_cap = | |
| 575 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 45 times.
|
45 | prism_refine_tile_leaves(nlist, dim, refine_cap_bytes); |
| 576 | |||
| 577 | 45 | shm_toc_estimate_chunk( | |
| 578 | &pcxt->estimator, mul_size(sizeof(WalUsage), pcxt->nworkers)); | ||
| 579 | 45 | shm_toc_estimate_chunk( | |
| 580 | &pcxt->estimator, mul_size(sizeof(BufferUsage), pcxt->nworkers)); | ||
| 581 | |||
| 582 | 45 | int querylen = 0; | |
| 583 |
1/2✓ Branch 0 taken 45 times.
✗ Branch 1 not taken.
|
45 | if (debug_query_string) |
| 584 | { | ||
| 585 | 45 | querylen = strlen(debug_query_string); | |
| 586 | 45 | shm_toc_estimate_chunk(&pcxt->estimator, querylen + 1); | |
| 587 | } | ||
| 588 | |||
| 589 | /* nkeys: shared, barrier, centroids, km_workers, root_assign, | ||
| 590 | * sortshared, wal, buffer, global_mean + optionally query_text (samples | ||
| 591 | * and the subtree ring live in their own segments; the refine | ||
| 592 | * accumulator overlays the samples) */ | ||
| 593 | 45 | int nkeys = 9; | |
| 594 |
1/2✓ Branch 0 taken 45 times.
✗ Branch 1 not taken.
|
45 | if (debug_query_string) |
| 595 | 45 | nkeys++; | |
| 596 | 45 | shm_toc_estimate_keys(&pcxt->estimator, nkeys); | |
| 597 | |||
| 598 | /* Total bytes the leader is about to commit to the DSM segment (sum of all | ||
| 599 | * estimated chunks), for the planned-allocation introspection line. */ | ||
| 600 | 45 | Size dsm_total = pcxt->estimator.space_for_chunks; | |
| 601 | |||
| 602 | 45 | InitializeParallelDSM(pcxt); | |
| 603 | |||
| 604 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 45 times.
|
45 | if (pcxt->seg == NULL) |
| 605 | { | ||
| 606 | ✗ | prism_pbuild_teardown(pcxt); | |
| 607 | ✗ | return false; | |
| 608 | } | ||
| 609 | |||
| 610 | /* ---- Populate shared state ---- */ | ||
| 611 | 45 | PrismBuildSharedPg *pg = shm_toc_allocate(pcxt->toc, est_shared); | |
| 612 | 45 | PrismBuildShared *shared = &pg->base; | |
| 613 | 45 | pg->heaprelid = RelationGetRelid(heap); | |
| 614 | 45 | pg->indexrelid = RelationGetRelid(index); | |
| 615 | 45 | pg->queryid = pgstat_get_my_query_id(); | |
| 616 | 45 | shared->concurrent = config->concurrent; | |
| 617 | 45 | shared->dim = dim; | |
| 618 | 45 | shared->metric = config->metric; | |
| 619 | 45 | shared->nlist = nlist; | |
| 620 | 45 | shared->fan_out = fan_out; /* resolved (auto if config 0) */ | |
| 621 | 45 | shared->subtree_slot_size = 0; /* leader sizes the ring post-assign */ | |
| 622 | 45 | shared->soar_lambda = config->soar_lambda; | |
| 623 | 45 | shared->boundary_epsilon = config->boundary_epsilon; | |
| 624 | 45 | shared->fastscan = config->fastscan; | |
| 625 | 45 | shared->centroid_format = config->centroid_format; | |
| 626 | 45 | shared->rabitq_seed = rabitq_seed; | |
| 627 | 45 | shared->nparticipants = nparticipants; | |
| 628 | 45 | shared->work_mem_kb = maintenance_work_mem; | |
| 629 | 45 | shared->max_samples_per_worker = max_per_worker; | |
| 630 | 45 | shared->km_max_iterations = 20; | |
| 631 | 45 | shared->km_tolerance = 1e-4f; | |
| 632 | 45 | shared->km_k = km_k; | |
| 633 | 45 | shared->km_converged = false; | |
| 634 | 45 | shared->refine_threshold = (uint32_t)prism_leaf_refine_threshold; | |
| 635 | 45 | shared->refine = false; /* leader decides post-clustering */ | |
| 636 | 45 | shared->refine_tile_cap = refine_tile_cap; | |
| 637 | /* Build routes for accuracy, not query speed (see PRISM_BUILD_CENTROID_* | ||
| 638 | * in posting_build.h): decouple from the query-tuned GUCs. */ | ||
| 639 | 45 | shared->centroid_error_scale = PRISM_BUILD_CENTROID_ERROR_SCALE; | |
| 640 | 45 | shared->centroid_beam_scale = PRISM_BUILD_CENTROID_BEAM_SCALE; | |
| 641 | 45 | shared->fastscan_bits = prism_fastscan_bits; | |
| 642 | 45 | SpinLockInit(&pg->mutex); | |
| 643 |
2/2✓ Branch 0 taken 11520 times.
✓ Branch 1 taken 45 times.
|
11565 | for (int i = 0; i < PRISM_REFINE_LOCK_STRIPES; i++) |
| 644 | 11520 | SpinLockInit(&pg->accum_locks[i]); | |
| 645 | 45 | shared->reltuples = 0.0; | |
| 646 | 45 | shared->indtuples = 0.0; | |
| 647 | 45 | shared->soar_dupes = 0.0; | |
| 648 | 45 | table_parallelscan_initialize( | |
| 649 | heap, ParallelTableScanFromVsShared(shared), snapshot); | ||
| 650 | 45 | shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_SHARED, shared); | |
| 651 | |||
| 652 | /* | ||
| 653 | * Barrier for phase synchronization. It is a *dynamic* barrier (init with | ||
| 654 | * 0 parties): the leader and every launched worker join via BarrierAttach, | ||
| 655 | * so the party tracks the participants that actually start. A non-zero | ||
| 656 | * (static) party is wrong here — BarrierAttach forbids it | ||
| 657 | * (Assert(!static_party)), and a fixed planned count would deadlock when | ||
| 658 | * PostgreSQL launches fewer workers than planned. The leader attaches in | ||
| 659 | * do_parallel_build before launching workers. | ||
| 660 | */ | ||
| 661 | 45 | Barrier *barrier = shm_toc_allocate(pcxt->toc, sizeof(Barrier)); | |
| 662 | 45 | BarrierInit(barrier, 0); | |
| 663 | 45 | shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_BARRIER, barrier); | |
| 664 | |||
| 665 | /* Sample slots, in a dedicated DSM segment so the build's largest | ||
| 666 | * working set can be handed back (last detach destroys it) before the | ||
| 667 | * posting sort claims its own budget -- see the sample-region seam. | ||
| 668 | * Workers attach via the handle published in the shared state. */ | ||
| 669 | 45 | dsm_segment *sample_seg = dsm_create(samp_sz, 0); | |
| 670 | 45 | PrismDsmSamples *dsm_samples = dsm_segment_address(sample_seg); | |
| 671 | /* Only the header + per-participant counts are read before being written; | ||
| 672 | * the sample data is filled by the sampling pass and read back bounded by | ||
| 673 | * those counts, so zeroing the (multi-GB) data region is wasted work. */ | ||
| 674 | 45 | memset(dsm_samples, | |
| 675 | 0, | ||
| 676 | 45 | MAXALIGN(sizeof(PrismDsmSamples)) + | |
| 677 | (size_t)nparticipants * sizeof(uint32_t)); | ||
| 678 | 45 | dsm_samples->nparticipants = nparticipants; | |
| 679 | 45 | dsm_samples->max_per_worker = max_per_worker; | |
| 680 | 45 | dsm_samples->dim = dim; | |
| 681 | 45 | pg->sample_handle = dsm_segment_handle(sample_seg); | |
| 682 | |||
| 683 | /* Shared centroids + norms (root k-means, k=km_k) */ | ||
| 684 | 45 | Size cent_sz = prism_dsm_centroids_size(km_k, dim); | |
| 685 | 45 | char *centroids_base = shm_toc_allocate(pcxt->toc, cent_sz); | |
| 686 | 45 | memset(centroids_base, 0, cent_sz); | |
| 687 | 45 | shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_CENTROIDS, centroids_base); | |
| 688 | 45 | float *cents = prism_dsm_centroids(centroids_base); | |
| 689 | |||
| 690 | /* Per-worker k-means accumulators */ | ||
| 691 | 45 | Size km_sz = prism_dsm_km_workers_size(nparticipants, km_k, dim); | |
| 692 | 45 | char *km_workers_base = shm_toc_allocate(pcxt->toc, km_sz); | |
| 693 | 45 | memset(km_workers_base, 0, km_sz); | |
| 694 | 45 | shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_KM_WORKERS, km_workers_base); | |
| 695 | |||
| 696 | /* Root assignment slots */ | ||
| 697 | 45 | Size ra_sz = prism_dsm_root_assign_size(nparticipants, max_per_worker); | |
| 698 | 45 | PrismDsmRootAssign *dsm_ra = shm_toc_allocate(pcxt->toc, ra_sz); | |
| 699 | 45 | memset(dsm_ra, 0, ra_sz); | |
| 700 | 45 | dsm_ra->nparticipants = nparticipants; | |
| 701 | 45 | dsm_ra->max_per_worker = max_per_worker; | |
| 702 | 45 | shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_ROOT_ASSIGN, dsm_ra); | |
| 703 | |||
| 704 | /* Shared coordinator for the cluster-keyed posting sort (sort seam). Sized | ||
| 705 | * for the planned participant count; the leader calls | ||
| 706 | * prism_pbuild_sort_shared_init with the actual launched count after | ||
| 707 | * launch. | ||
| 708 | */ | ||
| 709 | 45 | Size sort_sz = prism_pbuild_sort_shared_size(nparticipants); | |
| 710 | 45 | void *sortshared = shm_toc_allocate(pcxt->toc, sort_sz); | |
| 711 | 45 | memset(sortshared, 0, sort_sz); | |
| 712 | 45 | shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_SORTSHARED, sortshared); | |
| 713 | |||
| 714 | /* Page-backed routing region (filled by the leader before the tree-ready | ||
| 715 | * barrier): the global mean. The posting-head base is a scalar in | ||
| 716 | * PrismBuildShared, so there is no O(nlist) head array here. */ | ||
| 717 | 45 | float *dsm_gmean = shm_toc_allocate(pcxt->toc, (Size)vec_nbytes); | |
| 718 | 45 | memset(dsm_gmean, 0, (Size)vec_nbytes); | |
| 719 | 45 | shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_GLOBAL_MEAN, dsm_gmean); | |
| 720 | |||
| 721 | 45 | WalUsage *walusage = shm_toc_allocate( | |
| 722 | 45 | pcxt->toc, mul_size(sizeof(WalUsage), pcxt->nworkers)); | |
| 723 | 45 | memset(walusage, 0, mul_size(sizeof(WalUsage), pcxt->nworkers)); | |
| 724 | 45 | shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_WAL_USAGE, walusage); | |
| 725 | |||
| 726 | 45 | BufferUsage *bufferusage = shm_toc_allocate( | |
| 727 | 45 | pcxt->toc, mul_size(sizeof(BufferUsage), pcxt->nworkers)); | |
| 728 | 45 | memset(bufferusage, 0, mul_size(sizeof(BufferUsage), pcxt->nworkers)); | |
| 729 | 45 | shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_BUFFER_USAGE, bufferusage); | |
| 730 | |||
| 731 |
1/2✓ Branch 0 taken 45 times.
✗ Branch 1 not taken.
|
45 | if (debug_query_string) |
| 732 | { | ||
| 733 | 45 | char *sq = shm_toc_allocate(pcxt->toc, querylen + 1); | |
| 734 | 45 | memcpy(sq, debug_query_string, querylen + 1); | |
| 735 | 45 | shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_QUERY_TEXT, sq); | |
| 736 | } | ||
| 737 | |||
| 738 | 45 | lead->pcxt = pcxt; | |
| 739 | 45 | lead->shared = shared; | |
| 740 | 45 | lead->barrier = barrier; | |
| 741 | 45 | lead->dsm_samples = dsm_samples; | |
| 742 | 45 | lead->sample_seg = sample_seg; | |
| 743 | 45 | lead->centroids_base = centroids_base; | |
| 744 | 45 | lead->cents = cents; | |
| 745 | 45 | lead->km_workers_base = km_workers_base; | |
| 746 | 45 | lead->dsm_ra = dsm_ra; | |
| 747 | 45 | lead->walusage = walusage; | |
| 748 | 45 | lead->bufferusage = bufferusage; | |
| 749 | 45 | lead->nparticipants = nparticipants; | |
| 750 | 45 | lead->km_k = km_k; | |
| 751 | 45 | lead->max_per_worker = max_per_worker; | |
| 752 | 45 | lead->dim = dim; | |
| 753 | 45 | lead->nlist = nlist; | |
| 754 | 45 | lead->rabitq_seed = rabitq_seed; | |
| 755 | 45 | lead->fan_out = fan_out; | |
| 756 | 45 | lead->dsm_total = dsm_total; | |
| 757 | 45 | return true; | |
| 758 | } | ||
| 759 | |||
| 760 | /* | ||
| 761 | * Accumulate one worker's tuple counts into the shared state under the lock. | ||
| 762 | * Back-end seam: only the spinlock is PG-specific (it lives in the derived | ||
| 763 | * struct); the standalone version locks a pthread mutex around the same adds. | ||
| 764 | */ | ||
| 765 | void | ||
| 766 | 84 | prism_pbuild_worker_add_counts( | |
| 767 | PrismBuildShared *shared, | ||
| 768 | double indtuples, | ||
| 769 | double soar_dupes, | ||
| 770 | double heap_tuples) | ||
| 771 | { | ||
| 772 | 84 | PrismBuildSharedPg *pg = (PrismBuildSharedPg *)shared; | |
| 773 | |||
| 774 |
1/2✗ Branch 1 not taken.
✓ Branch 2 taken 84 times.
|
84 | SpinLockAcquire(&pg->mutex); |
| 775 | 84 | shared->indtuples += indtuples; | |
| 776 | 84 | shared->soar_dupes += soar_dupes; | |
| 777 | 84 | shared->reltuples += heap_tuples; | |
| 778 | 84 | SpinLockRelease(&pg->mutex); | |
| 779 | 84 | } | |
| 780 | |||
| 781 | /* | ||
| 782 | * Sample-region seam (see parallel_build.h): the samples live in their own | ||
| 783 | * DSM segment; workers attach by the handle the leader published, and every | ||
| 784 | * participant detaches independently once the refine pass is done -- the | ||
| 785 | * last detach destroys the segment and returns the build's largest working | ||
| 786 | * set before the posting sort claims its budget. | ||
| 787 | */ | ||
| 788 | PrismDsmSamples * | ||
| 789 | 86 | prism_pbuild_samples_attach( | |
| 790 | shm_toc *toc, PrismBuildShared *shared, void **seg_out) | ||
| 791 | { | ||
| 792 | 86 | PrismBuildSharedPg *pg = (PrismBuildSharedPg *)shared; | |
| 793 | |||
| 794 | 86 | (void)toc; | |
| 795 | 86 | dsm_segment *seg = dsm_attach(pg->sample_handle); | |
| 796 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 86 times.
|
86 | if (seg == NULL) |
| 797 | ✗ | ereport(ERROR, | |
| 798 | (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), | ||
| 799 | errmsg("could not attach to prism sample segment"))); | ||
| 800 | 86 | *seg_out = seg; | |
| 801 | 86 | return (PrismDsmSamples *)dsm_segment_address(seg); | |
| 802 | } | ||
| 803 | |||
| 804 | void | ||
| 805 | 128 | prism_pbuild_samples_release(PrismDsmSamples *samples, void *seg) | |
| 806 | { | ||
| 807 | 128 | (void)samples; | |
| 808 |
1/2✓ Branch 0 taken 128 times.
✗ Branch 1 not taken.
|
128 | if (seg != NULL) |
| 809 | 128 | dsm_detach((dsm_segment *)seg); | |
| 810 | 128 | } | |
| 811 | |||
| 812 | /* | ||
| 813 | * Subtree-ring seam: the ring of nparticipants subtree slots lives in its | ||
| 814 | * own DSM segment, created by the leader only after root assignment -- the | ||
| 815 | * per-child sample counts known then bound the slot far below the analytic | ||
| 816 | * worst case. The handle travels in shared; every participant detaches when | ||
| 817 | * the streaming pass is done, and the last detach frees the memory. | ||
| 818 | */ | ||
| 819 | char * | ||
| 820 | 20 | prism_pbuild_subtree_ring_create( | |
| 821 | PrismBuildShared *shared, | ||
| 822 | int nparticipants, | ||
| 823 | uint64_t slot_size, | ||
| 824 | void **seg_out) | ||
| 825 | { | ||
| 826 | 20 | PrismBuildSharedPg *pg = (PrismBuildSharedPg *)shared; | |
| 827 | |||
| 828 | 20 | dsm_segment *seg = dsm_create( | |
| 829 | prism_dsm_child_subtrees_size(nparticipants, slot_size), 0); | ||
| 830 | 20 | pg->subtree_ring_handle = dsm_segment_handle(seg); | |
| 831 | 20 | shared->subtree_slot_size = slot_size; | |
| 832 | 20 | *seg_out = seg; | |
| 833 | 20 | return (char *)dsm_segment_address(seg); | |
| 834 | } | ||
| 835 | |||
| 836 | char * | ||
| 837 | 38 | prism_pbuild_subtree_ring_attach(PrismBuildShared *shared, void **seg_out) | |
| 838 | { | ||
| 839 | 38 | PrismBuildSharedPg *pg = (PrismBuildSharedPg *)shared; | |
| 840 | |||
| 841 | 38 | dsm_segment *seg = dsm_attach(pg->subtree_ring_handle); | |
| 842 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 38 times.
|
38 | if (seg == NULL) |
| 843 | ✗ | ereport(ERROR, | |
| 844 | (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), | ||
| 845 | errmsg("could not attach to prism subtree segment"))); | ||
| 846 | 38 | *seg_out = seg; | |
| 847 | 38 | return (char *)dsm_segment_address(seg); | |
| 848 | } | ||
| 849 | |||
| 850 | void | ||
| 851 | 55 | prism_pbuild_subtree_ring_release(void *seg) | |
| 852 | { | ||
| 853 |
1/2✓ Branch 0 taken 55 times.
✗ Branch 1 not taken.
|
55 | if (seg != NULL) |
| 854 | 55 | dsm_detach((dsm_segment *)seg); | |
| 855 | 55 | } | |
| 856 | |||
| 857 | /* | ||
| 858 | * Exact-centroid seam: the exact centroid collection lives in its own DSM | ||
| 859 | * segment, created by the leader only after the streaming tree write (its | ||
| 860 | * size — a few MB — is known then). The handle travels in shared; every | ||
| 861 | * participant detaches when its page-backed routing is done, and the last | ||
| 862 | * detach frees the memory. | ||
| 863 | */ | ||
| 864 | char * | ||
| 865 | 43 | prism_pbuild_exact_centroids_create( | |
| 866 | PrismBuildShared *shared, uint64_t nbytes, void **seg_out) | ||
| 867 | { | ||
| 868 | 43 | PrismBuildSharedPg *pg = (PrismBuildSharedPg *)shared; | |
| 869 | |||
| 870 | 43 | dsm_segment *seg = dsm_create(nbytes, 0); | |
| 871 | 43 | pg->exact_centroids_handle = dsm_segment_handle(seg); | |
| 872 | 43 | *seg_out = seg; | |
| 873 | 43 | return (char *)dsm_segment_address(seg); | |
| 874 | } | ||
| 875 | |||
| 876 | char * | ||
| 877 | 84 | prism_pbuild_exact_centroids_attach(PrismBuildShared *shared, void **seg_out) | |
| 878 | { | ||
| 879 | 84 | PrismBuildSharedPg *pg = (PrismBuildSharedPg *)shared; | |
| 880 | |||
| 881 | 84 | dsm_segment *seg = dsm_attach(pg->exact_centroids_handle); | |
| 882 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 84 times.
|
84 | if (seg == NULL) |
| 883 | ✗ | ereport(ERROR, | |
| 884 | (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), | ||
| 885 | errmsg("could not attach to prism exact-centroid " | ||
| 886 | "segment"))); | ||
| 887 | 84 | *seg_out = seg; | |
| 888 | 84 | return (char *)dsm_segment_address(seg); | |
| 889 | } | ||
| 890 | |||
| 891 | void | ||
| 892 | 127 | prism_pbuild_exact_centroids_release(void *seg) | |
| 893 | { | ||
| 894 |
1/2✓ Branch 0 taken 127 times.
✗ Branch 1 not taken.
|
127 | if (seg != NULL) |
| 895 | 127 | dsm_detach((dsm_segment *)seg); | |
| 896 | 127 | } | |
| 897 | |||
| 898 | /* | ||
| 899 | * Re-initialize the parallel scan for the posting pass; the sampling pass | ||
| 900 | * consumed the first one. Back-end seam: the standalone version resets its | ||
| 901 | * work-stealing cursor. | ||
| 902 | */ | ||
| 903 | void | ||
| 904 | 45 | prism_pbuild_rescan(Relation heap, PrismBuildShared *shared) | |
| 905 | { | ||
| 906 | 45 | table_parallelscan_reinitialize( | |
| 907 | heap, ParallelTableScanFromVsShared(shared)); | ||
| 908 | 45 | } | |
| 909 | |||
| 910 | /* | ||
| 911 | * Striped-lock seam for the shared leaf-refinement accumulator. Contention is | ||
| 912 | * low because rows spread across nleaves leaves into PRISM_REFINE_LOCK_STRIPES | ||
| 913 | * stripes. | ||
| 914 | */ | ||
| 915 | void | ||
| 916 | 12600 | prism_pbuild_accum_lock(PrismBuildShared *shared, uint32_t stripe) | |
| 917 | { | ||
| 918 | 12600 | PrismBuildSharedPg *pg = (PrismBuildSharedPg *)shared; | |
| 919 |
2/2✓ Branch 1 taken 70 times.
✓ Branch 2 taken 12530 times.
|
12600 | SpinLockAcquire(&pg->accum_locks[stripe]); |
| 920 | 12600 | } | |
| 921 | |||
| 922 | void | ||
| 923 | 12600 | prism_pbuild_accum_unlock(PrismBuildShared *shared, uint32_t stripe) | |
| 924 | { | ||
| 925 | 12600 | PrismBuildSharedPg *pg = (PrismBuildSharedPg *)shared; | |
| 926 | 12600 | SpinLockRelease(&pg->accum_locks[stripe]); | |
| 927 | 12600 | } | |
| 928 | |||
| 929 | /* ---------------------------------------------------------------- | ||
| 930 | * Posting sort seam (PG back-end) — parallel tuplesort | ||
| 931 | * | ||
| 932 | * Each entry is sorted by a uint32 cluster id; the payload is a fixed-size | ||
| 933 | * opaque blob carried in a bytea column. Workers feed partial runs into the | ||
| 934 | * shared Sharedsort (in the build DSM); the leader merges and reads them back | ||
| 935 | * grouped by cluster. Memory is bounded by maintenance_work_mem (tuplesort | ||
| 936 | * spills past it). The standalone back-end provides the same-named seam over | ||
| 937 | * in-memory arrays. | ||
| 938 | * ---------------------------------------------------------------- */ | ||
| 939 | struct PrismSorter | ||
| 940 | { | ||
| 941 | Tuplesortstate *ts; | ||
| 942 | TupleDesc tupdesc; | ||
| 943 | TupleTableSlot *slot; | ||
| 944 | SortCoordinate coord; | ||
| 945 | uint32_t entry_size; | ||
| 946 | bytea *payload; /* reusable scratch: VARHDRSZ + entry_size */ | ||
| 947 | bool is_leader; | ||
| 948 | }; | ||
| 949 | |||
| 950 | Size | ||
| 951 | 90 | prism_pbuild_sort_shared_size(int nparticipants) | |
| 952 | { | ||
| 953 | 90 | return tuplesort_estimate_shared(nparticipants); | |
| 954 | } | ||
| 955 | |||
| 956 | void | ||
| 957 | 43 | prism_pbuild_sort_shared_init(void *region, int nparticipants, void *seg) | |
| 958 | { | ||
| 959 | 43 | tuplesort_initialize_shared( | |
| 960 | (Sharedsort *)region, nparticipants, (dsm_segment *)seg); | ||
| 961 | 43 | } | |
| 962 | |||
| 963 | PrismSorter * | ||
| 964 | 278 | prism_pbuild_sort_begin( | |
| 965 | void *region, | ||
| 966 | void *seg, | ||
| 967 | int participant, | ||
| 968 | int nparticipants, | ||
| 969 | bool is_leader, | ||
| 970 | uint32_t entry_size, | ||
| 971 | int work_mem_kb) | ||
| 972 | { | ||
| 973 | 278 | (void)participant; | |
| 974 | 278 | PrismSorter *s = palloc0(sizeof(PrismSorter)); | |
| 975 | 278 | s->entry_size = entry_size; | |
| 976 | 278 | s->is_leader = is_leader; | |
| 977 | |||
| 978 | 278 | s->tupdesc = CreateTemplateTupleDesc(2); | |
| 979 | 278 | TupleDescInitEntry(s->tupdesc, 1, "cluster", INT4OID, -1, 0); | |
| 980 | 278 | TupleDescInitEntry(s->tupdesc, 2, "payload", BYTEAOID, -1, 0); | |
| 981 | #if PG_VERSION_NUM >= 190000 | ||
| 982 | TupleDescFinalize(s->tupdesc); | ||
| 983 | #endif | ||
| 984 | |||
| 985 | /* region == NULL: a plain, non-parallel sort (the serial build) — no | ||
| 986 | * coordinate, no attach. Otherwise a parallel participant: a leader | ||
| 987 | * merging runs, or a worker producing one (which attaches to the shared | ||
| 988 | * fileset). */ | ||
| 989 |
2/2✓ Branch 0 taken 127 times.
✓ Branch 1 taken 151 times.
|
278 | if (region != NULL) |
| 990 | { | ||
| 991 | 127 | s->coord = palloc0(sizeof(SortCoordinateData)); | |
| 992 | 127 | s->coord->sharedsort = (Sharedsort *)region; | |
| 993 |
2/2✓ Branch 0 taken 43 times.
✓ Branch 1 taken 84 times.
|
127 | if (is_leader) |
| 994 | { | ||
| 995 | 43 | s->coord->isWorker = false; | |
| 996 | 43 | s->coord->nParticipants = nparticipants; | |
| 997 | } | ||
| 998 | else | ||
| 999 | { | ||
| 1000 | 84 | s->coord->isWorker = true; | |
| 1001 | 84 | s->coord->nParticipants = -1; | |
| 1002 | } | ||
| 1003 | } | ||
| 1004 | |||
| 1005 | 278 | AttrNumber attNums[1] = {1}; | |
| 1006 | 278 | Oid sortOps[1] = {Int4LessOperator}; | |
| 1007 | 278 | Oid sortColls[1] = {InvalidOid}; | |
| 1008 | 278 | bool nullsF[1] = {false}; | |
| 1009 | 278 | s->ts = tuplesort_begin_heap( | |
| 1010 | s->tupdesc, | ||
| 1011 | 1, | ||
| 1012 | attNums, | ||
| 1013 | sortOps, | ||
| 1014 | sortColls, | ||
| 1015 | nullsF, | ||
| 1016 | work_mem_kb, | ||
| 1017 | s->coord, | ||
| 1018 | TUPLESORT_NONE); | ||
| 1019 | |||
| 1020 | /* Workers attach to the shared fileset; the leader holds it via the | ||
| 1021 | * backend's dsm reference and must not attach (per tuplesort.h). */ | ||
| 1022 |
2/2✓ Branch 0 taken 84 times.
✓ Branch 1 taken 194 times.
|
278 | if (region != NULL && !is_leader) |
| 1023 | 84 | tuplesort_attach_shared((Sharedsort *)region, (dsm_segment *)seg); | |
| 1024 | |||
| 1025 | 278 | s->slot = MakeSingleTupleTableSlot(s->tupdesc, &TTSOpsMinimalTuple); | |
| 1026 | 278 | s->payload = (bytea *)palloc(VARHDRSZ + entry_size); | |
| 1027 | 278 | SET_VARSIZE(s->payload, VARHDRSZ + entry_size); | |
| 1028 | 278 | return s; | |
| 1029 | } | ||
| 1030 | |||
| 1031 | void | ||
| 1032 | 394144 | prism_pbuild_sort_put(PrismSorter *s, uint32_t cluster, const void *entry) | |
| 1033 | { | ||
| 1034 | 394144 | memcpy(VARDATA(s->payload), entry, s->entry_size); | |
| 1035 | 394144 | ExecClearTuple(s->slot); | |
| 1036 | 394144 | s->slot->tts_values[0] = Int32GetDatum((int32)cluster); | |
| 1037 | 394144 | s->slot->tts_isnull[0] = false; | |
| 1038 | 394144 | s->slot->tts_values[1] = PointerGetDatum(s->payload); | |
| 1039 | 394144 | s->slot->tts_isnull[1] = false; | |
| 1040 | 394144 | ExecStoreVirtualTuple(s->slot); | |
| 1041 | 394144 | tuplesort_puttupleslot(s->ts, s->slot); | |
| 1042 | 394144 | } | |
| 1043 | |||
| 1044 | void | ||
| 1045 | 278 | prism_pbuild_sort_performsort(PrismSorter *s) | |
| 1046 | { | ||
| 1047 | 278 | tuplesort_performsort(s->ts); | |
| 1048 | 278 | } | |
| 1049 | |||
| 1050 | bool | ||
| 1051 | 394338 | prism_pbuild_sort_getnext( | |
| 1052 | PrismSorter *s, uint32_t *cluster, const void **entry) | ||
| 1053 | { | ||
| 1054 | 394338 | bool isnull; | |
| 1055 |
2/2✓ Branch 1 taken 394144 times.
✓ Branch 2 taken 194 times.
|
394338 | if (!tuplesort_gettupleslot(s->ts, true, false, s->slot, NULL)) |
| 1056 | return false; | ||
| 1057 | 394144 | *cluster = (uint32_t)DatumGetInt32(slot_getattr(s->slot, 1, &isnull)); | |
| 1058 | 394144 | bytea *pl = DatumGetByteaPP(slot_getattr(s->slot, 2, &isnull)); | |
| 1059 |
2/2✓ Branch 0 taken 392317 times.
✓ Branch 1 taken 1827 times.
|
394144 | *entry = VARDATA_ANY(pl); |
| 1060 | 394144 | return true; | |
| 1061 | } | ||
| 1062 | |||
| 1063 | void | ||
| 1064 | 278 | prism_pbuild_sort_end(PrismSorter *s) | |
| 1065 | { | ||
| 1066 | 278 | tuplesort_end(s->ts); | |
| 1067 | 278 | ExecDropSingleTupleTableSlot(s->slot); | |
| 1068 | 278 | FreeTupleDesc(s->tupdesc); | |
| 1069 |
2/2✓ Branch 0 taken 127 times.
✓ Branch 1 taken 151 times.
|
278 | if (s->coord != NULL) /* NULL for a non-parallel (serial) sort */ |
| 1070 | 127 | pfree(s->coord); | |
| 1071 | 278 | pfree(s->payload); | |
| 1072 | 278 | pfree(s); | |
| 1073 | 278 | } | |
| 1074 |