| Line | Branch | Exec | Source |
|---|---|---|---|
| 1 | /* | ||
| 2 | * Copyright (c) 2026 Tiger Data, Inc. | ||
| 3 | * Licensed under the PostgreSQL License. See LICENSE for details. | ||
| 4 | * | ||
| 5 | * build.c - Index build for prism | ||
| 6 | * | ||
| 7 | * Serial build phases (do_serial_build; the parallel shape lives in | ||
| 8 | * parallel_build_leader.c and falls back here when no workers launch): | ||
| 9 | * 1. Resolve dimension, metric, centroid format, nlist, fan_out | ||
| 10 | * 2. Sample vectors for clustering (BlockSampler + reservoir) | ||
| 11 | * 3. Routing-tree PLAN pass: cluster the sample once, recording each | ||
| 12 | * node into a spillable blob store, and size the page layout | ||
| 13 | * 4. Pre-extend the relation for the centroid pages + posting heads | ||
| 14 | * 5. Routing-tree WRITE pass: replay the recorded nodes and stream | ||
| 15 | * every centroid page and posting-head page (pt_centroid resident) | ||
| 16 | * 6. Optional leaf refinement (prism.leaf_refine_threshold) re-centers | ||
| 17 | * the head encode references from the full table | ||
| 18 | * 7. Single heap scan: route each row page-backed (the same | ||
| 19 | * prism_query_route the query/insert use) into the cluster-keyed sort, | ||
| 20 | * then build each posting list | ||
| 21 | * 8. Write the metadata page (final tuple count) and WAL-log | ||
| 22 | * | ||
| 23 | * Block layout: | ||
| 24 | * Block 0: Metadata page | ||
| 25 | * Blocks 1..C: Centroid pages (streamed before the scan) | ||
| 26 | * Blocks C+1..H: Posting-list head pages (leaf c at C+1+c) | ||
| 27 | * Blocks H+1..N: Posting continuation pages (appended during the scan) | ||
| 28 | * | ||
| 29 | * Memory layout: | ||
| 30 | * build_ctx — all build-phase allocations; deleted in one shot | ||
| 31 | * tmp_ctx — per-tuple scratch; reset after each callback | ||
| 32 | */ | ||
| 33 | |||
| 34 | #include <postgres.h> | ||
| 35 | |||
| 36 | #include <access/heaptoast.h> | ||
| 37 | #include <access/htup_details.h> | ||
| 38 | #include <access/table.h> | ||
| 39 | #include <access/tableam.h> | ||
| 40 | #include <access/xloginsert.h> | ||
| 41 | #include <catalog/index.h> | ||
| 42 | #include <common/pg_prng.h> | ||
| 43 | #include <math.h> | ||
| 44 | #include <miscadmin.h> | ||
| 45 | #include <portability/instr_time.h> | ||
| 46 | #include <utils/memutils.h> | ||
| 47 | #include <utils/rel.h> | ||
| 48 | #include <utils/sampling.h> | ||
| 49 | |||
| 50 | #include "algo/distance.h" | ||
| 51 | #include "algo/hkmeans.h" | ||
| 52 | #include "algo/kmeans.h" | ||
| 53 | #include "algo/vecops.h" | ||
| 54 | #include "build.h" | ||
| 55 | #include "core/log.h" | ||
| 56 | #include "index/centroid_build.h" | ||
| 57 | #include "index/centroid_page.h" | ||
| 58 | #include "index/index_base.h" | ||
| 59 | #include "index/index_build.h" | ||
| 60 | #include "index/parallel_build.h" | ||
| 61 | #include "index/posting_build.h" | ||
| 62 | #include "index/posting_page.h" | ||
| 63 | #include "index/query_scan.h" | ||
| 64 | #include "meta.h" | ||
| 65 | #include "pg/bufstorage.h" | ||
| 66 | #include "quant/rabitq.h" | ||
| 67 | #include "support_pg.h" | ||
| 68 | #include "typeinfo.h" | ||
| 69 | #include "types/vec16.h" | ||
| 70 | #include "types/vec32.h" | ||
| 71 | |||
| 72 | /* ---------------------------------------------------------------- | ||
| 73 | * Build state (PrismBuildParams is in build.h, shared with the | ||
| 74 | * parallel leader) | ||
| 75 | * ---------------------------------------------------------------- */ | ||
| 76 | |||
| 77 | typedef struct PrismBuildState | ||
| 78 | { | ||
| 79 | PrismBuildParams params; | ||
| 80 | |||
| 81 | double indtuples; /* total count */ | ||
| 82 | double soar_dupes; /* replicated SOAR vectors */ | ||
| 83 | |||
| 84 | /* | ||
| 85 | * Posting entries are streamed into a cluster-keyed tuplesort during the | ||
| 86 | * heap scan, then read back grouped by cluster so the build holds only ONE | ||
| 87 | * posting-page builder at a time. This bounds build memory to | ||
| 88 | * maintenance_work_mem (the tuplesort stays in RAM until it exceeds it, | ||
| 89 | * then spills); per-cluster resident builders would cost O(nlist) memory. | ||
| 90 | */ | ||
| 91 | /* Posting entries are RaBitQ-encoded during the scan (relative to the | ||
| 92 | * assigned cluster centroid) and fed to a cluster-keyed sorter — the same | ||
| 93 | * PrismSorter seam the parallel build uses, here in non-parallel mode (no | ||
| 94 | * coordinate). The shared build loop (prism_posting_build_lists) reads | ||
| 95 | * them back grouped by cluster and writes one resident page builder at a | ||
| 96 | * time. | ||
| 97 | */ | ||
| 98 | PrismSorter *sorter; | ||
| 99 | RaBitQParams *rq_params; | ||
| 100 | |||
| 101 | /* Sampling */ | ||
| 102 | float *samples; /* [max_samples * dim] row-major */ | ||
| 103 | int nsamples; /* current sample count */ | ||
| 104 | int max_samples; /* target sample count */ | ||
| 105 | double rowstoskip; /* reservoir sampling state */ | ||
| 106 | |||
| 107 | ReservoirStateData rstate; | ||
| 108 | |||
| 109 | /* PG context */ | ||
| 110 | Relation heap; | ||
| 111 | Relation index; | ||
| 112 | struct IndexInfo *index_info; | ||
| 113 | MemoryContext build_ctx; /* all build allocations */ | ||
| 114 | MemoryContext tmp_ctx; /* per-tuple scratch */ | ||
| 115 | |||
| 116 | PrismBuildProgress *prog; /* phase/progress reporting seam (serial path) */ | ||
| 117 | |||
| 118 | /* | ||
| 119 | * Page-backed assignment (routes each row the same way the query/insert | ||
| 120 | * do, so the in-RAM tree is not needed for the scan). qs holds the routing | ||
| 121 | * state; route is the shared route+encode+emit context (also used by the | ||
| 122 | * parallel posting workers), which references qs and the sorter. | ||
| 123 | */ | ||
| 124 | PrismQueryState qs; | ||
| 125 | PrismBuildRouteCtx route; | ||
| 126 | |||
| 127 | /* | ||
| 128 | * Column type bound to this index's dimension, with the conversion buffer | ||
| 129 | * it needs. Held on the state because both callbacks run under a tmp_ctx | ||
| 130 | * that is reset after every tuple. | ||
| 131 | */ | ||
| 132 | Vec32Access input; | ||
| 133 | } PrismBuildState; | ||
| 134 | |||
| 135 | /* ---------------------------------------------------------------- | ||
| 136 | * Helpers | ||
| 137 | * ---------------------------------------------------------------- */ | ||
| 138 | |||
| 139 | /* ---------------------------------------------------------------- | ||
| 140 | * Sampling | ||
| 141 | * ---------------------------------------------------------------- */ | ||
| 142 | |||
| 143 | static void | ||
| 144 | 251996 | sample_callback( | |
| 145 | Relation index, | ||
| 146 | ItemPointer tid, | ||
| 147 | Datum *values, | ||
| 148 | bool *isnull, | ||
| 149 | bool tuple_is_alive, | ||
| 150 | void *state) | ||
| 151 | { | ||
| 152 | 251996 | PrismBuildState *bs = (PrismBuildState *)state; | |
| 153 | |||
| 154 | 251996 | (void)index; | |
| 155 | 251996 | (void)tid; | |
| 156 | 251996 | (void)tuple_is_alive; | |
| 157 | |||
| 158 |
2/2✓ Branch 0 taken 20 times.
✓ Branch 1 taken 251976 times.
|
251996 | if (isnull[0]) |
| 159 | 20 | return; | |
| 160 | |||
| 161 | 251976 | MemoryContext old_ctx = MemoryContextSwitchTo(bs->tmp_ctx); | |
| 162 | |||
| 163 | 251976 | Dimension dim = bs->params.dim; | |
| 164 | 251976 | Vec32Ref vref = vec32_read(&bs->input, values[0]); | |
| 165 | 251976 | const float *src = vref.data; | |
| 166 | 251976 | const size_t vec_nbytes = (size_t)dim * sizeof(float); | |
| 167 | |||
| 168 |
2/2✓ Branch 0 taken 245976 times.
✓ Branch 1 taken 6000 times.
|
251976 | if (bs->nsamples < bs->max_samples) |
| 169 | { | ||
| 170 | 245976 | memcpy(bs->samples + (size_t)bs->nsamples * dim, src, vec_nbytes); | |
| 171 | 245976 | bs->nsamples++; | |
| 172 | } | ||
| 173 | else | ||
| 174 | { | ||
| 175 |
2/2✓ Branch 0 taken 5453 times.
✓ Branch 1 taken 547 times.
|
6000 | if (bs->rowstoskip < 0) |
| 176 | 5453 | bs->rowstoskip = reservoir_get_next_S( | |
| 177 | &bs->rstate, bs->nsamples, bs->max_samples); | ||
| 178 | |||
| 179 |
2/2✓ Branch 0 taken 5453 times.
✓ Branch 1 taken 547 times.
|
6000 | if (bs->rowstoskip <= 0) |
| 180 | { | ||
| 181 | 10906 | int k = (int)(bs->max_samples * | |
| 182 | 5453 | sampler_random_fract(&bs->rstate.randstate)); | |
| 183 |
2/4✓ Branch 0 taken 5453 times.
✗ Branch 1 not taken.
✗ Branch 2 not taken.
✓ Branch 3 taken 5453 times.
|
5453 | Assert(k >= 0 && k < bs->max_samples); |
| 184 | 5453 | memcpy(bs->samples + (size_t)k * dim, src, vec_nbytes); | |
| 185 | } | ||
| 186 | 6000 | bs->rowstoskip -= 1; | |
| 187 | 6000 | bs->nsamples++; | |
| 188 | } | ||
| 189 | |||
| 190 | 251976 | MemoryContextSwitchTo(old_ctx); | |
| 191 | 251976 | MemoryContextReset(bs->tmp_ctx); | |
| 192 | } | ||
| 193 | |||
| 194 | static void | ||
| 195 | 151 | sample_rows(PrismBuildState *bs) | |
| 196 | { | ||
| 197 | 151 | BlockNumber totalblocks = RelationGetNumberOfBlocks(bs->heap); | |
| 198 | 151 | BlockSamplerData bsampler; | |
| 199 | |||
| 200 | 151 | bs->rowstoskip = -1; | |
| 201 | |||
| 202 | 151 | BlockSampler_Init( | |
| 203 | &bsampler, | ||
| 204 | totalblocks, | ||
| 205 | bs->max_samples, | ||
| 206 | pg_prng_uint32(&pg_global_prng_state)); | ||
| 207 | 151 | reservoir_init_selection_state(&bs->rstate, bs->max_samples); | |
| 208 | |||
| 209 |
2/2✓ Branch 2 taken 6600 times.
✓ Branch 3 taken 151 times.
|
6751 | while (BlockSampler_HasMore(&bsampler)) |
| 210 | { | ||
| 211 | 6600 | BlockNumber targblock = BlockSampler_Next(&bsampler); | |
| 212 | |||
| 213 | 6600 | table_index_build_range_scan( | |
| 214 | bs->heap, | ||
| 215 | bs->index, | ||
| 216 | bs->index_info, | ||
| 217 | false, | ||
| 218 | true, | ||
| 219 | false, | ||
| 220 | targblock, | ||
| 221 | 1, | ||
| 222 | sample_callback, | ||
| 223 | (void *)bs, | ||
| 224 | NULL); | ||
| 225 | } | ||
| 226 | 151 | } | |
| 227 | |||
| 228 | /* ---------------------------------------------------------------- | ||
| 229 | * Build callback — single-pass: route each row page-backed, | ||
| 230 | * stream encoded entries into the cluster-keyed sorter | ||
| 231 | * ---------------------------------------------------------------- */ | ||
| 232 | |||
| 233 | static void | ||
| 234 | 251996 | build_callback( | |
| 235 | Relation index, | ||
| 236 | ItemPointer tid, | ||
| 237 | Datum *values, | ||
| 238 | bool *isnull, | ||
| 239 | bool tuple_is_alive, | ||
| 240 | void *state) | ||
| 241 | { | ||
| 242 | 251996 | PrismBuildState *bs = (PrismBuildState *)state; | |
| 243 | |||
| 244 | 251996 | (void)index; | |
| 245 | 251996 | (void)tuple_is_alive; | |
| 246 | |||
| 247 |
2/2✓ Branch 0 taken 20 times.
✓ Branch 1 taken 251976 times.
|
251996 | if (isnull[0]) |
| 248 | 20 | return; | |
| 249 | |||
| 250 | 251976 | MemoryContext old_ctx = MemoryContextSwitchTo(bs->tmp_ctx); | |
| 251 | |||
| 252 | 251976 | Vec32Ref vref = vec32_read(&bs->input, values[0]); | |
| 253 | |||
| 254 | /* Route + encode + stream via the shared page-backed helper (the parallel | ||
| 255 | * posting workers use the very same call). tmp_ctx is reset after every | ||
| 256 | * tuple, so nothing reachable from this call may allocate memory that | ||
| 257 | * outlives the callback: the route context's buffers (candidates, batch, | ||
| 258 | * encode scratch) are all preallocated for exactly this reason. */ | ||
| 259 | 251976 | prism_build_route_emit(&bs->route, vref.data, *tid); | |
| 260 | |||
| 261 |
2/2✓ Branch 0 taken 15 times.
✓ Branch 1 taken 251961 times.
|
251976 | if (((uint64_t)bs->route.indtuples % 10000) == 0) |
| 262 | 15 | prism_build_report_progress(bs->prog, bs->route.indtuples); | |
| 263 | |||
| 264 | 251976 | MemoryContextSwitchTo(old_ctx); | |
| 265 | 251976 | MemoryContextReset(bs->tmp_ctx); | |
| 266 | } | ||
| 267 | |||
| 268 | /* ---------------------------------------------------------------- | ||
| 269 | * Write metadata page | ||
| 270 | * ---------------------------------------------------------------- */ | ||
| 271 | |||
| 272 | static void | ||
| 273 | 194 | write_meta_page( | |
| 274 | VsStorage *storage, | ||
| 275 | Dimension dim, | ||
| 276 | uint8_t nlevels, | ||
| 277 | uint8_t fan_out, | ||
| 278 | BlockNumber first_centroid, | ||
| 279 | BlockNumber first_posting, | ||
| 280 | uint32_t ncentroid_pages, | ||
| 281 | uint32_t nlist, | ||
| 282 | PrismCentroidFormat centroid_format, | ||
| 283 | DistanceMetric metric, | ||
| 284 | uint64_t rabitq_seed, | ||
| 285 | bool fastscan, | ||
| 286 | const float *global_mean) | ||
| 287 | { | ||
| 288 | /* Block 0 must already exist (extended or new_page'd by caller). | ||
| 289 | * Use write_page to write in-place. */ | ||
| 290 | 194 | Page page = vs_storage_write_page(storage, 0); | |
| 291 | |||
| 292 | 194 | PageInit(page, BLCKSZ, PRISM_META_SIZE(dim)); | |
| 293 | |||
| 294 | 194 | PrismMetaPage *meta = (PrismMetaPage *)PageGetSpecialPointer(page); | |
| 295 | 194 | meta->magic = PRISM_META_MAGIC; | |
| 296 | 194 | meta->dim = dim; | |
| 297 | 194 | meta->nlevels = nlevels; | |
| 298 | 194 | meta->centroid_format = (uint8_t)centroid_format; | |
| 299 | 194 | meta->first_centroid = first_centroid; | |
| 300 | 194 | meta->first_posting = first_posting; | |
| 301 | 194 | meta->ncentroid_pages = ncentroid_pages; | |
| 302 | 194 | meta->nlist = nlist; | |
| 303 | 194 | meta->metric = (uint8_t)metric; | |
| 304 | 194 | meta->fan_out = (uint8_t)fan_out; | |
| 305 | 194 | meta->flags = fastscan ? PRISM_META_FLAG_FASTSCAN : 0; | |
| 306 | 194 | meta->reserved = 0; | |
| 307 | 194 | meta->rabitq_seed = rabitq_seed; | |
| 308 | |||
| 309 | 194 | memcpy(prism_meta_global_mean(meta), global_mean, dim * sizeof(float)); | |
| 310 | |||
| 311 | 194 | vs_storage_commit_page(storage, 0); | |
| 312 | 194 | } | |
| 313 | |||
| 314 | /* ---------------------------------------------------------------- | ||
| 315 | * Helpers: resolve metric and centroid format from opclass/relopts | ||
| 316 | * ---------------------------------------------------------------- */ | ||
| 317 | |||
| 318 | static DistanceMetric | ||
| 319 | 196 | prism_get_metric(Relation index) | |
| 320 | { | ||
| 321 | 196 | FmgrInfo *procinfo = index_getprocinfo(index, 1, PRISM_METRIC_PROC); | |
| 322 | 196 | return (DistanceMetric)DatumGetInt32( | |
| 323 | FunctionCall1Coll(procinfo, InvalidOid, (Datum)0)); | ||
| 324 | } | ||
| 325 | |||
| 326 | static PrismCentroidFormat | ||
| 327 | 196 | prism_resolve_format(Relation index, DistanceMetric metric, Dimension dim) | |
| 328 | { | ||
| 329 | 196 | PrismOptions *opts = (PrismOptions *)index->rd_options; | |
| 330 | 551 | int cc = (opts != NULL) ? opts->centroid_compression | |
| 331 |
2/2✓ Branch 0 taken 159 times.
✓ Branch 1 taken 37 times.
|
196 | : PRISM_CENTROID_COMPRESSION_AUTO; |
| 332 | 355 | int cfs = (opts != NULL) ? opts->centroid_fastscan | |
| 333 | 159 | : PRISM_FASTSCAN_MODE_AUTO; | |
| 334 | |||
| 335 | /* | ||
| 336 | * RaBitQ centroids estimate L2 distance, which routes correctly for | ||
| 337 | * L2 and (normalized) cosine but not inner product. ON forces it and | ||
| 338 | * errors for inner product; AUTO compresses everything except inner | ||
| 339 | * product; OFF disables it. | ||
| 340 | */ | ||
| 341 | 196 | bool compressed; | |
| 342 |
3/3✓ Branch 0 taken 59 times.
✓ Branch 1 taken 125 times.
✓ Branch 2 taken 12 times.
|
196 | switch (cc) |
| 343 | { | ||
| 344 | 59 | case PRISM_CENTROID_COMPRESSION_ON: | |
| 345 |
2/2✓ Branch 0 taken 1 times.
✓ Branch 1 taken 58 times.
|
59 | if (metric == DISTANCE_INNER_PRODUCT) |
| 346 |
1/2✓ Branch 1 taken 1 times.
✗ Branch 2 not taken.
|
1 | ereport(ERROR, |
| 347 | (errcode(ERRCODE_INVALID_PARAMETER_VALUE), | ||
| 348 | errmsg("centroid_compression=on is not supported " | ||
| 349 | "with vec32_ip_ops"))); | ||
| 350 | compressed = true; | ||
| 351 | break; | ||
| 352 | case PRISM_CENTROID_COMPRESSION_OFF: | ||
| 353 | compressed = false; | ||
| 354 | break; | ||
| 355 | 125 | default: /* AUTO */ | |
| 356 | 125 | compressed = (metric != DISTANCE_INNER_PRODUCT); | |
| 357 | 125 | break; | |
| 358 | } | ||
| 359 | |||
| 360 | /* | ||
| 361 | * FASTSCAN centroids are a packed layout over the RaBitQ-compressed | ||
| 362 | * representation, so they require compression (and, like RaBitQ | ||
| 363 | * centroids, don't apply to inner product). The layout stores fixed | ||
| 364 | * 32-candidate groups, so past the dimension where a group no | ||
| 365 | * longer fits a page it cannot be used at all. AUTO uses it exactly | ||
| 366 | * where it is available; ON errors where it is not. | ||
| 367 | */ | ||
| 368 |
2/2✓ Branch 0 taken 7 times.
✓ Branch 1 taken 188 times.
|
195 | bool cfs_fits = prism_centroid_fastscan_max_groups(dim) > 0; |
| 369 |
2/2✓ Branch 0 taken 7 times.
✓ Branch 1 taken 188 times.
|
195 | if (cfs == PRISM_FASTSCAN_MODE_ON) |
| 370 | { | ||
| 371 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 7 times.
|
7 | if (metric == DISTANCE_INNER_PRODUCT) |
| 372 | ✗ | ereport(ERROR, | |
| 373 | (errcode(ERRCODE_INVALID_PARAMETER_VALUE), | ||
| 374 | errmsg("centroid_fastscan=on is not supported " | ||
| 375 | "with vec32_ip_ops"))); | ||
| 376 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 7 times.
|
7 | if (!compressed) |
| 377 | ✗ | ereport(ERROR, | |
| 378 | (errcode(ERRCODE_INVALID_PARAMETER_VALUE), | ||
| 379 | errmsg("centroid_fastscan=on requires " | ||
| 380 | "centroid_compression"))); | ||
| 381 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 7 times.
|
7 | if (!cfs_fits) |
| 382 | ✗ | ereport(ERROR, | |
| 383 | (errcode(ERRCODE_PROGRAM_LIMIT_EXCEEDED), | ||
| 384 | errmsg("centroid_fastscan=on does not support " | ||
| 385 | "%d dimensions", | ||
| 386 | dim))); | ||
| 387 | } | ||
| 388 |
4/4✓ Branch 0 taken 145 times.
✓ Branch 1 taken 50 times.
✓ Branch 2 taken 1 times.
✓ Branch 3 taken 144 times.
|
195 | if (compressed && cfs != PRISM_FASTSCAN_MODE_OFF && cfs_fits) |
| 389 | return PRISM_CENTROID_FMT_FASTSCAN; | ||
| 390 | |||
| 391 |
2/2✓ Branch 0 taken 14 times.
✓ Branch 1 taken 37 times.
|
51 | if (compressed) |
| 392 | return PRISM_CENTROID_FMT_RABITQ; | ||
| 393 | |||
| 394 | /* Uncompressed centroids follow the column type. Resolved directly rather | ||
| 395 | * than through the per-backend cache: there is no metadata page to | ||
| 396 | * populate that cache from until this build writes one. */ | ||
| 397 | 14 | return prism_index_type_info(index)->centroid_format; | |
| 398 | } | ||
| 399 | |||
| 400 | /* | ||
| 401 | * Resolve the posting-page format. Like the centroid layout, a | ||
| 402 | * fastscan posting group has a dimension-dependent fixed size; the | ||
| 403 | * binding constraint is a cluster's first page, which also carries | ||
| 404 | * the full-precision centroid reference. AUTO uses fastscan exactly | ||
| 405 | * where a group fits; ON errors where it does not. | ||
| 406 | */ | ||
| 407 | static bool | ||
| 408 | 195 | prism_resolve_fastscan(Relation index, Dimension dim) | |
| 409 | { | ||
| 410 | 195 | PrismOptions *opts = (PrismOptions *)index->rd_options; | |
| 411 |
2/2✓ Branch 0 taken 158 times.
✓ Branch 1 taken 37 times.
|
195 | int fs = (opts != NULL) ? opts->fastscan : PRISM_FASTSCAN_MODE_AUTO; |
| 412 |
3/3✓ Branch 0 taken 10 times.
✓ Branch 1 taken 161 times.
✓ Branch 2 taken 24 times.
|
195 | bool fits = prism_fastscan_max_groups(dim, true) > 0; |
| 413 | |||
| 414 |
3/3✓ Branch 0 taken 10 times.
✓ Branch 1 taken 161 times.
✓ Branch 2 taken 24 times.
|
195 | switch (fs) |
| 415 | { | ||
| 416 | 10 | case PRISM_FASTSCAN_MODE_ON: | |
| 417 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 10 times.
|
10 | if (!fits) |
| 418 | ✗ | ereport(ERROR, | |
| 419 | (errcode(ERRCODE_PROGRAM_LIMIT_EXCEEDED), | ||
| 420 | errmsg("fastscan=on does not support %d dimensions", | ||
| 421 | dim))); | ||
| 422 | return true; | ||
| 423 | case PRISM_FASTSCAN_MODE_OFF: | ||
| 424 | return false; | ||
| 425 | 161 | default: /* AUTO */ | |
| 426 | 161 | return fits; | |
| 427 | } | ||
| 428 | } | ||
| 429 | |||
| 430 | static uint32_t | ||
| 431 | 195 | prism_get_fan_out(Relation index) | |
| 432 | { | ||
| 433 | 195 | PrismOptions *opts = (PrismOptions *)index->rd_options; | |
| 434 |
1/2✓ Branch 0 taken 158 times.
✗ Branch 1 not taken.
|
158 | if (opts != NULL && opts->fan_out >= PRISM_MIN_FAN_OUT) |
| 435 | 158 | return (uint32_t)opts->fan_out; | |
| 436 | return PRISM_DEFAULT_FAN_OUT; | ||
| 437 | } | ||
| 438 | |||
| 439 | static void | ||
| 440 | 197 | resolve_build_params( | |
| 441 | Relation heap, Relation index, PrismBuildParams *p, double *est_rows) | ||
| 442 | { | ||
| 443 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 197 times.
|
197 | Dimension dim = (Dimension)TupleDescAttr(index->rd_att, 0)->atttypmod; |
| 444 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 197 times.
|
197 | if (dim == 0) |
| 445 | ✗ | ereport(ERROR, | |
| 446 | (errcode(ERRCODE_INVALID_PARAMETER_VALUE), | ||
| 447 | errmsg("column does not have dimensions"))); | ||
| 448 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 197 times.
|
197 | if (dim > VEC32_MAX_DIM) |
| 449 | ✗ | ereport(ERROR, | |
| 450 | (errcode(ERRCODE_PROGRAM_LIMIT_EXCEEDED), | ||
| 451 | errmsg("column cannot have more than %d dimensions", | ||
| 452 | VEC32_MAX_DIM))); | ||
| 453 |
2/2✓ Branch 0 taken 1 times.
✓ Branch 1 taken 196 times.
|
197 | if (dim > PRISM_INDEX_MAX_DIM) |
| 454 |
1/2✓ Branch 1 taken 1 times.
✗ Branch 2 not taken.
|
1 | ereport(ERROR, |
| 455 | (errcode(ERRCODE_PROGRAM_LIMIT_EXCEEDED), | ||
| 456 | errmsg("prism indexes support at most %d dimensions", | ||
| 457 | PRISM_INDEX_MAX_DIM))); | ||
| 458 | |||
| 459 | 196 | p->dim = dim; | |
| 460 | 196 | p->metric = prism_get_metric(index); | |
| 461 | 196 | p->centroid_format = prism_resolve_format(index, p->metric, dim); | |
| 462 |
2/2✓ Branch 0 taken 158 times.
✓ Branch 1 taken 37 times.
|
195 | p->fan_out = prism_get_fan_out(index); |
| 463 | |||
| 464 | /* nlist: use relopt if set, otherwise auto from sqrt(reltuples). | ||
| 465 | * The row estimate is computed once and surfaced to the caller (the | ||
| 466 | * progress-reporting total needs it too): the never-analyzed | ||
| 467 | * fallback samples heap pages, and sampling twice would double that | ||
| 468 | * cost and could even disagree with itself on a growing heap. */ | ||
| 469 | 195 | PrismOptions *opts = (PrismOptions *)index->rd_options; | |
| 470 |
2/2✓ Branch 0 taken 158 times.
✓ Branch 1 taken 37 times.
|
195 | uint32_t nlist_opt = (opts != NULL) ? (uint32_t)opts->nlist : 0; |
| 471 | 158 | uint32_t tpages_opt = (opts != NULL) ? (uint32_t)opts->target_pages : 0; | |
| 472 | |||
| 473 | 195 | *est_rows = prism_estimate_heap_tuples(heap); | |
| 474 | |||
| 475 |
2/2✓ Branch 0 taken 102 times.
✓ Branch 1 taken 93 times.
|
195 | if (nlist_opt > 0) |
| 476 | { | ||
| 477 | 102 | p->nlist = nlist_opt; | |
| 478 | } | ||
| 479 | else | ||
| 480 | { | ||
| 481 | 93 | p->nlist = prism_auto_nlist(*est_rows, p->dim, tpages_opt); | |
| 482 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 93 times.
|
93 | if (p->nlist > PRISM_MAX_NLIST) |
| 483 | ✗ | p->nlist = PRISM_MAX_NLIST; | |
| 484 | } | ||
| 485 | |||
| 486 | 390 | p->fan_out = | |
| 487 | 195 | prism_auto_fan_out(p->fan_out, p->nlist, PRISM_DEFAULT_FAN_OUT); | |
| 488 | |||
| 489 |
1/2✓ Branch 0 taken 158 times.
✗ Branch 1 not taken.
|
158 | p->kmeans_nredo = (opts != NULL && opts->kmeans_nredo > 0) |
| 490 | ? (uint32_t)opts->kmeans_nredo | ||
| 491 |
2/2✓ Branch 0 taken 158 times.
✓ Branch 1 taken 37 times.
|
353 | : 1; |
| 492 | |||
| 493 | 390 | p->soar_lambda = (opts != NULL) ? opts->soar_lambda | |
| 494 |
2/2✓ Branch 0 taken 158 times.
✓ Branch 1 taken 37 times.
|
195 | : PRISM_DEFAULT_SOAR_LAMBDA; |
| 495 | 390 | p->boundary_epsilon = (opts != NULL) ? opts->boundary_epsilon | |
| 496 |
2/2✓ Branch 0 taken 158 times.
✓ Branch 1 taken 37 times.
|
195 | : PRISM_DEFAULT_BOUNDARY_EPSILON; |
| 497 | 195 | p->fastscan = prism_resolve_fastscan(index, dim); | |
| 498 | |||
| 499 | /* The exact-centroid collection size is known from the resolved | ||
| 500 | * shape alone, so an under-budgeted maintenance_work_mem can be | ||
| 501 | * reported now — at build start, before the expensive phases — | ||
| 502 | * with the setting that would fit. The collector still enforces | ||
| 503 | * the budget against the real size at collection time. */ | ||
| 504 |
2/2✓ Branch 0 taken 181 times.
✓ Branch 1 taken 14 times.
|
195 | if (p->centroid_format == PRISM_CENTROID_FMT_RABITQ || |
| 505 | p->centroid_format == PRISM_CENTROID_FMT_FASTSCAN) | ||
| 506 | { | ||
| 507 | 181 | uint64_t expected = | |
| 508 |
2/2✓ Branch 0 taken 149 times.
✓ Branch 1 taken 32 times.
|
181 | prism_exact_centroid_expected_bytes(p->nlist, p->fan_out, dim); |
| 509 |
1/2✓ Branch 0 taken 181 times.
✗ Branch 1 not taken.
|
181 | uint64_t budget = prism_exact_centroid_budget( |
| 510 | (uint64_t)maintenance_work_mem); | ||
| 511 |
2/2✓ Branch 0 taken 1 times.
✓ Branch 1 taken 180 times.
|
181 | if (expected > budget) |
| 512 |
1/2✓ Branch 1 taken 1 times.
✗ Branch 2 not taken.
|
1 | vs_warn("maintenance_work_mem is likely too small for exact " |
| 513 | "centroid collection (about " UINT64_FORMAT | ||
| 514 | " kB needed for nlist=%u, fan_out=%u; the budget is " | ||
| 515 | "one eighth of maintenance_work_mem): the build will " | ||
| 516 | "fall back to estimated internal scoring — raise " | ||
| 517 | "maintenance_work_mem to at least " UINT64_FORMAT " kB", | ||
| 518 | expected / 1024 + 1, | ||
| 519 | p->nlist, | ||
| 520 | p->fan_out, | ||
| 521 | expected * 8 / 1024 + 1); | ||
| 522 | } | ||
| 523 | 195 | } | |
| 524 | |||
| 525 | /* | ||
| 526 | * Estimated heap row count. Uses the planner's reltuples when available, | ||
| 527 | * else samples real pages and counts normal line pointers (no | ||
| 528 | * visibility checks -- dead-but-unpruned rows count, like reltuples | ||
| 529 | * right after a bulk delete). Feeds the automatic partition count and | ||
| 530 | * the pg_stat_progress_create_index total. | ||
| 531 | * | ||
| 532 | * The density is measured rather than derived from the vector width | ||
| 533 | * because any width formula must assume a storage strategy, and a wrong | ||
| 534 | * assumption is catastrophic in both directions: a TOASTed column's | ||
| 535 | * main-fork rows are ~64 bytes where the dimension predicts kilobytes | ||
| 536 | * (undercounting rows two orders of magnitude), while a PLAIN-stored | ||
| 537 | * vector column holds ~2 rows per page where the TOAST assumption | ||
| 538 | * predicts ~128 (overcounting 64x -- and the automatic partition count | ||
| 539 | * inherits the error, sharding the index into starved clusters). | ||
| 540 | */ | ||
| 541 | double | ||
| 542 | 481 | prism_estimate_heap_tuples(Relation heap) | |
| 543 | { | ||
| 544 |
2/2✓ Branch 0 taken 362 times.
✓ Branch 1 taken 119 times.
|
481 | if (heap->rd_rel->reltuples > 0) |
| 545 | 362 | return heap->rd_rel->reltuples; | |
| 546 | |||
| 547 | 119 | BlockNumber nblocks = RelationGetNumberOfBlocks(heap); | |
| 548 |
2/2✓ Branch 0 taken 110 times.
✓ Branch 1 taken 9 times.
|
119 | if (nblocks == 0) |
| 549 | return 0; | ||
| 550 | |||
| 551 | 110 | uint32_t samples = Min(nblocks, 32); | |
| 552 | 110 | double rows = 0; | |
| 553 | |||
| 554 |
2/2✓ Branch 0 taken 801 times.
✓ Branch 1 taken 110 times.
|
911 | for (uint32_t i = 0; i < samples; i++) |
| 555 | { | ||
| 556 | 801 | BlockNumber blkno = (BlockNumber)(((uint64_t)nblocks * | |
| 557 | 801 | (2 * (uint64_t)i + 1)) / | |
| 558 | 801 | (2 * samples)); | |
| 559 | 801 | Buffer buf = ReadBuffer(heap, blkno); | |
| 560 | |||
| 561 | 801 | LockBuffer(buf, BUFFER_LOCK_SHARE); | |
| 562 | 801 | Page page = BufferGetPage(buf); | |
| 563 | |||
| 564 | /* An extended-but-unwritten page has pd_lower == 0, where | ||
| 565 | * PageGetMaxOffsetNumber underflows; such pages hold no rows. */ | ||
| 566 |
1/2✓ Branch 0 taken 801 times.
✗ Branch 1 not taken.
|
801 | if (!PageIsNew(page)) |
| 567 | { | ||
| 568 |
1/2✓ Branch 0 taken 801 times.
✗ Branch 1 not taken.
|
801 | OffsetNumber max = PageGetMaxOffsetNumber(page); |
| 569 |
2/2✓ Branch 0 taken 55624 times.
✓ Branch 1 taken 801 times.
|
56425 | for (OffsetNumber off = FirstOffsetNumber; off <= max; off++) |
| 570 | { | ||
| 571 |
1/2✓ Branch 0 taken 55624 times.
✗ Branch 1 not taken.
|
55624 | if (ItemIdIsNormal(PageGetItemId(page, off))) |
| 572 | 55624 | rows += 1; | |
| 573 | } | ||
| 574 | } | ||
| 575 | 801 | UnlockReleaseBuffer(buf); | |
| 576 | } | ||
| 577 | 110 | return rows / samples * nblocks; | |
| 578 | } | ||
| 579 | |||
| 580 | /* ---------------------------------------------------------------- | ||
| 581 | * Sample and cluster vectors | ||
| 582 | * ---------------------------------------------------------------- */ | ||
| 583 | |||
| 584 | /* | ||
| 585 | * Draw the maintenance_work_mem-bounded k-means sample into bs->samples (left | ||
| 586 | * resident for the caller to cluster + free), normalize it for cosine, and | ||
| 587 | * resolve the leaf target against the sample size. A heap with no indexable | ||
| 588 | * rows yields one synthetic sample and a single-cluster target, so the | ||
| 589 | * build always proceeds. | ||
| 590 | */ | ||
| 591 | static void | ||
| 592 | 151 | sample_for_build( | |
| 593 | PrismBuildState *bs, uint32_t *out_nlist, bool *out_subsampled) | ||
| 594 | { | ||
| 595 | 151 | Dimension dim = bs->params.dim; | |
| 596 | 151 | uint32_t nlist = bs->params.nlist; | |
| 597 | |||
| 598 | /* Bound the sample buffer by maintenance_work_mem (and MaxAllocSize). The | ||
| 599 | * tree is trained on this sample; leaf centroids are refined on the full | ||
| 600 | * table by a later (page-backed) refine pass when subsampling loses | ||
| 601 | * quality. */ | ||
| 602 | 151 | const uint64_t vec_nbytes = (uint64_t)dim * sizeof(float); | |
| 603 | |||
| 604 | 151 | uint64_t ideal_samples = Max((uint64_t)10000, (uint64_t)nlist * 256); | |
| 605 | 151 | uint64_t budget = (uint64_t)maintenance_work_mem * 1024 / vec_nbytes; | |
| 606 | 151 | uint64_t alloc_cap = MaxAllocSize / vec_nbytes; | |
| 607 | 151 | uint64_t cap = Min(budget, alloc_cap); | |
| 608 | 151 | if (cap < 10000) | |
| 609 | cap = 10000; | ||
| 610 | 151 | *out_subsampled = ideal_samples > cap; | |
| 611 |
2/2✓ Branch 0 taken 144 times.
✓ Branch 1 taken 7 times.
|
151 | bs->max_samples = (int)Min(ideal_samples, cap); |
| 612 | |||
| 613 | 151 | bs->nsamples = 0; | |
| 614 | 151 | bs->samples = palloc((size_t)bs->max_samples * vec_nbytes); | |
| 615 | |||
| 616 | 151 | sample_rows(bs); | |
| 617 | |||
| 618 |
2/2✓ Branch 0 taken 3 times.
✓ Branch 1 taken 148 times.
|
151 | if (bs->nsamples > bs->max_samples) |
| 619 | 3 | bs->nsamples = bs->max_samples; | |
| 620 | |||
| 621 |
2/2✓ Branch 0 taken 10 times.
✓ Branch 1 taken 141 times.
|
151 | if (bs->params.metric == DISTANCE_COSINE) |
| 622 | { | ||
| 623 | /* | ||
| 624 | * Normalize the sample onto the unit sphere, dropping zero-norm | ||
| 625 | * vectors: they carry no direction, so under cosine they sit at | ||
| 626 | * the origin -- far from every real point -- and k-means drags | ||
| 627 | * centroids toward them, warping the encode references (and so | ||
| 628 | * the distance estimates) for every vector in the affected | ||
| 629 | * subtree. They stay in the posting lists like any other row; | ||
| 630 | * they just don't get a vote on the clustering. | ||
| 631 | * | ||
| 632 | * Dropped here, by compacting the retained sample, rather than | ||
| 633 | * at collection time: the reservoir offers every live heap row, | ||
| 634 | * so a collection-time check would compute an O(dim) norm per | ||
| 635 | * table row to filter out a tiny fraction, while this pass | ||
| 636 | * touches only the <= max_samples retained rows (and copies | ||
| 637 | * nothing at all when no zeros were sampled). The cost is a | ||
| 638 | * sample budget short by the sampled-zero count -- in expectation | ||
| 639 | * the table's zero fraction, far below k-means run-to-run | ||
| 640 | * variance -- and slots spent on zeros bought nothing before | ||
| 641 | * this pass existed either. nlist is clamped to the compacted | ||
| 642 | * count below, so even a zero-heavy table degrades to fewer | ||
| 643 | * leaves rather than starved ones (all-zero degenerates to the | ||
| 644 | * synthetic single-cluster build). | ||
| 645 | */ | ||
| 646 | int kept = 0; | ||
| 647 |
2/2✓ Branch 0 taken 6842 times.
✓ Branch 1 taken 10 times.
|
6852 | for (int i = 0; i < bs->nsamples; i++) |
| 648 | { | ||
| 649 | 6842 | float *v = bs->samples + (size_t)i * dim; | |
| 650 | 6842 | float norm_sq = vs_l2_norm_squared(v, dim); | |
| 651 | |||
| 652 |
2/2✓ Branch 0 taken 44 times.
✓ Branch 1 taken 6798 times.
|
6842 | if (norm_sq == 0.0f) |
| 653 | 44 | continue; | |
| 654 | 6798 | vec32_scale(v, 1.0f / sqrtf(norm_sq), v, dim); | |
| 655 | /* The buffer is a dense row-major matrix consumed directly | ||
| 656 | * by k-means, so a dropped row leaves a hole that must be | ||
| 657 | * closed: shift each kept row down over it. Until the first | ||
| 658 | * drop, kept == i and rows stay normalized in place. */ | ||
| 659 |
2/2✓ Branch 0 taken 2998 times.
✓ Branch 1 taken 3800 times.
|
6798 | if (kept != i) |
| 660 | 2998 | memcpy(bs->samples + (size_t)kept * dim, | |
| 661 | v, | ||
| 662 | (size_t)dim * sizeof(float)); | ||
| 663 | 6798 | kept++; | |
| 664 | } | ||
| 665 | 10 | bs->nsamples = kept; | |
| 666 | } | ||
| 667 | |||
| 668 |
2/2✓ Branch 0 taken 8 times.
✓ Branch 1 taken 143 times.
|
151 | if (bs->nsamples == 0) |
| 669 | { | ||
| 670 | /* | ||
| 671 | * No indexable rows (empty table, or every row dead or NULL). | ||
| 672 | * Synthesize one unit-basis sample so the normal machinery | ||
| 673 | * emits a valid single-cluster index: the sample only shapes | ||
| 674 | * the leaf centroid -- the heap scan that fills posting lists | ||
| 675 | * still contributes nothing -- and later inserts route to that | ||
| 676 | * leaf like any single-cluster index (empty posting heads are | ||
| 677 | * a supported shape; every cluster gets one). A unit vector | ||
| 678 | * rather than zeros keeps the centroid safe for cosine | ||
| 679 | * normalization and the RaBitQ norm factors. Clustering | ||
| 680 | * quality after bulk loading comes from REINDEX, exactly as | ||
| 681 | * for any index built far below its final row count. | ||
| 682 | */ | ||
| 683 | 8 | bs->samples = repalloc(bs->samples, vec_nbytes); | |
| 684 | 8 | bs->max_samples = 1; | |
| 685 | 8 | memset(bs->samples, 0, vec_nbytes); | |
| 686 | 8 | bs->samples[0] = 1.0f; | |
| 687 | 8 | bs->nsamples = 1; | |
| 688 | } | ||
| 689 | |||
| 690 | 151 | if ((uint32_t)bs->nsamples < nlist) | |
| 691 | nlist = (uint32_t)bs->nsamples; | ||
| 692 | |||
| 693 | 151 | *out_nlist = nlist; | |
| 694 | 151 | } | |
| 695 | |||
| 696 | /* ---------------------------------------------------------------- | ||
| 697 | * Serial build | ||
| 698 | * ---------------------------------------------------------------- */ | ||
| 699 | |||
| 700 | /* ---------------------------------------------------------------- | ||
| 701 | * Page-backed leaf refinement (streaming build) | ||
| 702 | * | ||
| 703 | * The streaming write trained the centroids on the | ||
| 704 | * maintenance_work_mem-bounded sample. When that subsampled, refine the | ||
| 705 | * per-leaf encode reference on the full table: route every row page-backed to | ||
| 706 | * its leaf (the same routing the scan + query use), accumulate per-leaf means, | ||
| 707 | * and rewrite each leaf's head-page pt_centroid to the full-table mean. This | ||
| 708 | * tightens the RaBitQ residuals (the dominant recall factor) for every vector | ||
| 709 | * in the leaf. The accumulator is tiled to a bounded ceiling (like the posting | ||
| 710 | * reserve), so memory stays O(maintenance_work_mem) regardless of nlist; | ||
| 711 | * nleaves above the tile just means more (re-scanned) tiles. Routing stays on | ||
| 712 | * the sample-trained centroid pages, so a single pass reaches the fixed point | ||
| 713 | * (assignments do not shift). | ||
| 714 | * ---------------------------------------------------------------- */ | ||
| 715 | |||
| 716 | typedef struct RefineHeadState | ||
| 717 | { | ||
| 718 | PrismQueryState *qs; | ||
| 719 | BlockNumber first_posting; /* leaf c's head = first_posting + c */ | ||
| 720 | uint32_t nlist; | ||
| 721 | Dimension dim; | ||
| 722 | bool cosine; | ||
| 723 | MemoryContext tmp_ctx; | ||
| 724 | double *sums; /* [tile * dim], indexed by leaf - tile_lo */ | ||
| 725 | uint64_t *cnts; /* [tile] */ | ||
| 726 | float *scratch; /* [dim] normalized copy for cosine */ | ||
| 727 | uint32_t tile_lo; | ||
| 728 | uint32_t tile_hi; | ||
| 729 | } RefineHeadState; | ||
| 730 | |||
| 731 | static void | ||
| 732 | 60000 | refine_head_cb( | |
| 733 | Relation index, | ||
| 734 | ItemPointer tid, | ||
| 735 | Datum *values, | ||
| 736 | bool *isnull, | ||
| 737 | bool tuple_is_alive, | ||
| 738 | void *state) | ||
| 739 | { | ||
| 740 | 60000 | RefineHeadState *rs = (RefineHeadState *)state; | |
| 741 | |||
| 742 | 60000 | (void)index; | |
| 743 | 60000 | (void)tid; | |
| 744 | 60000 | (void)tuple_is_alive; | |
| 745 | |||
| 746 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 60000 times.
|
60000 | if (isnull[0]) |
| 747 | ✗ | return; | |
| 748 | |||
| 749 | 60000 | MemoryContext old_ctx = MemoryContextSwitchTo(rs->tmp_ctx); | |
| 750 | |||
| 751 | 60000 | const float *vin = Vec32ToRef(DatumGetVec32(values[0])).data; | |
| 752 | 60000 | uint32_t idx; | |
| 753 | 120000 | const float *v = prism_refine_route_row( | |
| 754 | 60000 | rs->qs, | |
| 755 | rs->first_posting, | ||
| 756 | vin, | ||
| 757 | 60000 | rs->dim, | |
| 758 | 60000 | rs->cosine, | |
| 759 | rs->scratch, | ||
| 760 | rs->tile_lo, | ||
| 761 | rs->tile_hi, | ||
| 762 | &idx); | ||
| 763 |
2/2✓ Branch 0 taken 36000 times.
✓ Branch 1 taken 24000 times.
|
60000 | if (v != NULL) |
| 764 | { | ||
| 765 | 36000 | double *sum = rs->sums + (size_t)idx * rs->dim; | |
| 766 |
2/2✓ Branch 0 taken 2304000 times.
✓ Branch 1 taken 36000 times.
|
2340000 | for (Dimension j = 0; j < rs->dim; j++) |
| 767 | 2304000 | sum[j] += v[j]; | |
| 768 | 36000 | rs->cnts[idx]++; | |
| 769 | } | ||
| 770 | |||
| 771 | 60000 | MemoryContextSwitchTo(old_ctx); | |
| 772 | 60000 | MemoryContextReset(rs->tmp_ctx); | |
| 773 | } | ||
| 774 | |||
| 775 | static void | ||
| 776 | 3 | serial_refine_heads( | |
| 777 | PrismBuildState *bs, | ||
| 778 | PrismHeadWriteCtx *headctx, | ||
| 779 | PrismQueryState *qs, | ||
| 780 | BlockNumber first_posting, | ||
| 781 | uint32_t nlist) | ||
| 782 | { | ||
| 783 | 3 | Dimension dim = bs->params.dim; | |
| 784 | |||
| 785 | 3 | uint64_t cap_bytes = | |
| 786 | 3 | Min((uint64_t)maintenance_work_mem * 1024, (uint64_t)MaxAllocSize); | |
| 787 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 3 times.
|
3 | uint32_t tile = prism_refine_tile_leaves(nlist, dim, cap_bytes); |
| 788 | |||
| 789 | 12 | RefineHeadState rs = { | |
| 790 | .qs = qs, | ||
| 791 | .first_posting = first_posting, | ||
| 792 | .dim = dim, | ||
| 793 | 3 | .cosine = (bs->params.metric == DISTANCE_COSINE), | |
| 794 | 3 | .tmp_ctx = bs->tmp_ctx, | |
| 795 | 3 | .sums = palloc((size_t)tile * dim * sizeof(double)), | |
| 796 | 3 | .cnts = palloc((size_t)tile * sizeof(uint64_t)), | |
| 797 | 3 | .scratch = palloc((size_t)dim * sizeof(float)), | |
| 798 | }; | ||
| 799 | |||
| 800 |
2/2✓ Branch 0 taken 5 times.
✓ Branch 1 taken 3 times.
|
8 | for (uint32_t lo = 0; lo < nlist; lo += tile) |
| 801 | { | ||
| 802 | 5 | uint32_t hi = Min(lo + tile, nlist); | |
| 803 | 5 | rs.tile_lo = lo; | |
| 804 | 5 | rs.tile_hi = hi; | |
| 805 | 5 | memset(rs.sums, 0, (size_t)(hi - lo) * dim * sizeof(double)); | |
| 806 | 5 | memset(rs.cnts, 0, (size_t)(hi - lo) * sizeof(uint64_t)); | |
| 807 | |||
| 808 | 5 | table_index_build_scan( | |
| 809 | bs->heap, | ||
| 810 | bs->index, | ||
| 811 | bs->index_info, | ||
| 812 | true, | ||
| 813 | false, | ||
| 814 | refine_head_cb, | ||
| 815 | (void *)&rs, | ||
| 816 | NULL); | ||
| 817 | |||
| 818 | 5 | prism_refine_write_means( | |
| 819 | 5 | rs.sums, | |
| 820 | 5 | rs.cnts, | |
| 821 | lo, | ||
| 822 | hi, | ||
| 823 | dim, | ||
| 824 | rs.scratch, | ||
| 825 | prism_write_leaf_head, | ||
| 826 | headctx); | ||
| 827 | } | ||
| 828 | |||
| 829 | 3 | pfree(rs.sums); | |
| 830 | 3 | pfree(rs.cnts); | |
| 831 | 3 | pfree(rs.scratch); | |
| 832 | 3 | } | |
| 833 | |||
| 834 | /* | ||
| 835 | * Serial build fallback, mirroring do_parallel_build's role for the | ||
| 836 | * non-parallel path. Streams the centroid tree straight to pages (no in-RAM | ||
| 837 | * tree, no O(nlist*dim) blob), then scans the heap once through build_callback | ||
| 838 | * (page-backed routing) to fill the posting lists. | ||
| 839 | * | ||
| 840 | * No in-RAM tree is materialized; the streamed tree's shape (leaf count + | ||
| 841 | * depth) comes back through out_nlist/out_tree_nlevels. Posting-list heads | ||
| 842 | * are formula-derived (first_posting + leaf), so no head array is returned. | ||
| 843 | * A heap with no indexable rows still builds: the sampler substitutes one | ||
| 844 | * synthetic sample and the result is a valid single-cluster index with an | ||
| 845 | * empty posting head. The serial path writes its own metadata page; the | ||
| 846 | * caller must not finalize again. | ||
| 847 | */ | ||
| 848 | static void | ||
| 849 | 151 | do_serial_build( | |
| 850 | PrismBuildState *bs, | ||
| 851 | VsStorage *storage, | ||
| 852 | uint64_t rabitq_seed, | ||
| 853 | uint32_t *out_nlist, | ||
| 854 | uint8_t *out_tree_nlevels, | ||
| 855 | float **out_global_mean, | ||
| 856 | double *out_heap_tuples, | ||
| 857 | double *out_indtuples, | ||
| 858 | double *out_soar_dupes) | ||
| 859 | { | ||
| 860 | 151 | const PrismBuildParams *p = &bs->params; | |
| 861 | 151 | Dimension dim = p->dim; | |
| 862 | |||
| 863 | 151 | *out_nlist = 0; | |
| 864 | 151 | *out_tree_nlevels = 0; | |
| 865 | |||
| 866 | 151 | prism_build_report_phase(bs->prog, PRISM_BUILD_PHASE_SAMPLE); | |
| 867 | |||
| 868 | 151 | uint32_t target_nlist = 0; | |
| 869 | 151 | bool subsampled = false; | |
| 870 | 151 | sample_for_build(bs, &target_nlist, &subsampled); | |
| 871 | |||
| 872 | 151 | KMeansOptions km_opts = VS_KMEANS_OPTIONS_DEFAULT; | |
| 873 | 151 | km_opts.algorithm = KMEANS_ALGO_LLOYD; | |
| 874 | 151 | km_opts.nredo = p->kmeans_nredo; | |
| 875 | |||
| 876 | /* | ||
| 877 | * Plan pass: cluster the sample once and discover the tree shape (leaf | ||
| 878 | * count, depth, centroid page count, leaf-centroid mean) without | ||
| 879 | * writing. Each node's clustering is recorded in a spillable blob store | ||
| 880 | * (BufFile-backed, so the resident cost stays one node) and the write | ||
| 881 | * pass below replays it instead of running k-means again -- the same | ||
| 882 | * single-clustering shape as the parallel build's subtree store. | ||
| 883 | * target_nlist (not the resolved leaf count) drives both passes. | ||
| 884 | */ | ||
| 885 | 151 | prism_build_report_phase(bs->prog, PRISM_BUILD_PHASE_KMEANS); | |
| 886 | 151 | PrismBlobStore *node_store = prism_pbuild_blobstore_begin(); | |
| 887 | 151 | PrismStreamTreePlan plan; | |
| 888 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 151 times.
|
151 | if (!prism_routing_tree_plan( |
| 889 | 151 | bs->samples, | |
| 890 | 151 | (uint32_t)bs->nsamples, | |
| 891 | dim, | ||
| 892 | target_nlist, | ||
| 893 | 151 | p->fan_out, | |
| 894 | 151 | p->metric, | |
| 895 | 151 | p->centroid_format, | |
| 896 | &km_opts, | ||
| 897 | node_store, | ||
| 898 | &plan)) | ||
| 899 | ✗ | ereport(ERROR, | |
| 900 | (errcode(ERRCODE_INTERNAL_ERROR), | ||
| 901 | errmsg("hierarchical k-means failed"))); | ||
| 902 | 151 | prism_pbuild_blobstore_rewind(node_store); | |
| 903 | |||
| 904 | 151 | uint32_t nlist = plan.nleaves; | |
| 905 | 151 | bs->params.nlist = nlist; | |
| 906 | |||
| 907 | /* | ||
| 908 | * Re-center the encoder on the leaf-centroid mean the plan pass reported | ||
| 909 | * (the write pass reproduces the identical tree). The in-RAM-tree build | ||
| 910 | * centers on the mean of the leaf centroids, not the per-vector sample | ||
| 911 | * mean, and the quantization quality of every centroid and posting code | ||
| 912 | * depends on this anchor. | ||
| 913 | */ | ||
| 914 | /* The streamed pages encode against the leaf-centroid mean the plan | ||
| 915 | * pass accumulated. */ | ||
| 916 | 151 | const size_t vec_nbytes = (size_t)dim * sizeof(float); | |
| 917 | |||
| 918 | 151 | float *global_mean = palloc(vec_nbytes); | |
| 919 |
2/2✓ Branch 0 taken 10 times.
✓ Branch 1 taken 141 times.
|
151 | memcpy(global_mean, plan.leaf_mean, vec_nbytes); |
| 920 |
2/2✓ Branch 0 taken 10 times.
✓ Branch 1 taken 141 times.
|
151 | if (p->metric == DISTANCE_COSINE) |
| 921 | 10 | vs_l2_normalize(global_mean, dim); | |
| 922 | 151 | vs_free(plan.leaf_mean); | |
| 923 | 151 | plan.leaf_mean = NULL; | |
| 924 | |||
| 925 | 151 | prism_build_report_phase(bs->prog, PRISM_BUILD_PHASE_SETUP); | |
| 926 | 151 | RaBitQParams *rq_params = vs_rabitq_create(dim, rabitq_seed); | |
| 927 | |||
| 928 | /* Centroid area is [first_centroid, first_posting); block 0 is metadata. | ||
| 929 | */ | ||
| 930 | 151 | BlockNumber first_centroid = PRISM_FIRST_CENTROID_BLKNO; | |
| 931 | 151 | BlockNumber first_posting = first_centroid + plan.centroid_pages; | |
| 932 | |||
| 933 | /* | ||
| 934 | * Head blocks are formula-derived: leaf c's head is first_posting + c, a | ||
| 935 | * contiguous head region of nlist pages. Pre-extend the relation to cover | ||
| 936 | * metadata + centroid + head region so the streaming write can place | ||
| 937 | * centroids at reserved blocks and each leaf's head page already exists | ||
| 938 | * when the write pass emits it; continuation pages are appended past the | ||
| 939 | * head region during prism_posting_build_lists. No O(nlist) reserve | ||
| 940 | * arrays. | ||
| 941 | */ | ||
| 942 | 151 | prism_build_reserve_layout(storage, first_posting + nlist); | |
| 943 | |||
| 944 | /* | ||
| 945 | * Write pass: stream the centroid pages (reserved blocks, post-order, root | ||
| 946 | * last) and, per leaf, its head page carrying pt_centroid. Both build and | ||
| 947 | * query then route page-backed over these centroid pages; the metadata | ||
| 948 | * page is written after the scan, when the tuple count is final. | ||
| 949 | */ | ||
| 950 | 151 | prism_build_report_phase(bs->prog, PRISM_BUILD_PHASE_CENTROID); | |
| 951 | |||
| 952 | /* Exact internal-centroid collection for the build descent (see | ||
| 953 | * PrismExactCentroidCollector in index_build.h). */ | ||
| 954 | 151 | PrismExactCentroidCollector exact_centroids = {0}; | |
| 955 | 334 | PrismExactCentroidCollector *collector = | |
| 956 |
2/2✓ Branch 0 taken 36 times.
✓ Branch 1 taken 115 times.
|
151 | prism_exact_centroid_enabled(plan.nlevels, p->centroid_format) |
| 957 | ? &exact_centroids | ||
| 958 | : NULL; | ||
| 959 | 32 | if (collector != NULL) | |
| 960 |
1/2✓ Branch 0 taken 32 times.
✗ Branch 1 not taken.
|
64 | prism_exact_centroid_collector_init( |
| 961 | collector, | ||
| 962 | dim, | ||
| 963 | p->centroid_format, | ||
| 964 | first_centroid, | ||
| 965 | plan.centroid_pages, | ||
| 966 | prism_exact_centroid_budget((uint64_t)maintenance_work_mem), | ||
| 967 |
1/2✓ Branch 0 taken 32 times.
✗ Branch 1 not taken.
|
32 | prism_exact_centroid_expected_slots(p->nlist, p->fan_out)); |
| 968 | |||
| 969 | 151 | PrismHeadWriteCtx headctx; | |
| 970 | 302 | prism_head_write_ctx_init( | |
| 971 | 151 | &headctx, storage, rq_params, dim, p->fastscan, first_posting); | |
| 972 | 302 | BlockNumber root = prism_routing_tree_write( | |
| 973 | storage, | ||
| 974 | 151 | (uint32_t)bs->nsamples, | |
| 975 | dim, | ||
| 976 | 151 | p->metric, | |
| 977 | target_nlist, | ||
| 978 | 151 | p->fan_out, | |
| 979 | 151 | p->centroid_format, | |
| 980 | rq_params, | ||
| 981 | global_mean, | ||
| 982 | node_store, | ||
| 983 | first_posting, | ||
| 984 | first_centroid, | ||
| 985 | prism_write_leaf_head, | ||
| 986 | &headctx, | ||
| 987 | collector); | ||
| 988 | 151 | prism_pbuild_blobstore_end(node_store); | |
| 989 | 151 | prism_head_write_ctx_cleanup(&headctx); | |
| 990 | 151 | pfree(bs->samples); | |
| 991 | 151 | bs->samples = NULL; | |
| 992 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 151 times.
|
151 | if (root == InvalidBlockNumber) |
| 993 | ✗ | ereport(ERROR, | |
| 994 | (errcode(ERRCODE_INTERNAL_ERROR), | ||
| 995 | errmsg("streaming centroid write failed"))); | ||
| 996 | |||
| 997 | /* | ||
| 998 | * Page-backed routing: build a PrismIndexBase from the just-written index | ||
| 999 | * so the scan routes each row exactly as the query/insert do | ||
| 1000 | * (prism_query_route over the centroid pages). Assignment uses no in-RAM | ||
| 1001 | * tree. The base is stack-local but outlives the scan (all within this | ||
| 1002 | * function); qs holds it by pointer until prism_query_state_cleanup below. | ||
| 1003 | */ | ||
| 1004 | 151 | PrismIndexBase idx_base; | |
| 1005 | /* Route the build for accuracy, not query speed: the build-time routing | ||
| 1006 | * constants rather than the query-tuned GUCs (see PRISM_BUILD_CENTROID_* | ||
| 1007 | * in posting_build.h). */ | ||
| 1008 | 302 | prism_build_router_base_init( | |
| 1009 | &idx_base, | ||
| 1010 | rq_params, | ||
| 1011 | storage, | ||
| 1012 | dim, | ||
| 1013 | 151 | (uint8_t)plan.nlevels, | |
| 1014 | root, | ||
| 1015 | 151 | p->metric, | |
| 1016 | 151 | p->centroid_format, | |
| 1017 | prism_fastscan_bits, | ||
| 1018 | PRISM_BUILD_CENTROID_ERROR_SCALE, | ||
| 1019 | PRISM_BUILD_CENTROID_BEAM_SCALE, | ||
| 1020 | bs->params.fan_out, | ||
| 1021 | nlist, | ||
| 1022 | rabitq_seed, | ||
| 1023 | global_mean, | ||
| 1024 | 151 | palloc(vec_nbytes)); | |
| 1025 | /* Build-only accuracy hook: score the internal tree levels against the | ||
| 1026 | * exact centroids collected during the streaming write (the query and | ||
| 1027 | * insert paths never set this; see PrismExactInternalCentroids). */ | ||
| 1028 | 151 | PrismExactInternalCentroids exact_view; | |
| 1029 |
2/2✓ Branch 0 taken 32 times.
✓ Branch 1 taken 119 times.
|
151 | if (collector != NULL) |
| 1030 | { | ||
| 1031 | 32 | prism_exact_centroid_view(collector, &exact_view); | |
| 1032 | 32 | idx_base.exact_internal = &exact_view; | |
| 1033 | } | ||
| 1034 | 151 | prism_query_state_init(&bs->qs, &idx_base, 1, PRISM_SECONDARY_TOPK); | |
| 1035 | |||
| 1036 | /* | ||
| 1037 | * When the sample was budget-limited AND the leaves are sample-thin | ||
| 1038 | * (below prism.leaf_refine_threshold samples per leaf -- a leaf's encode | ||
| 1039 | * reference is a sample mean whose error shrinks with its count, so | ||
| 1040 | * well-fed leaves gain nothing from the extra scan), refine each leaf's | ||
| 1041 | * encode reference on the full table (page-backed, bounded) before the | ||
| 1042 | * encode scan. | ||
| 1043 | */ | ||
| 1044 |
4/4✓ Branch 0 taken 4 times.
✓ Branch 1 taken 147 times.
✓ Branch 2 taken 3 times.
✓ Branch 3 taken 1 times.
|
151 | if (subsampled && bs->nsamples >= bs->max_samples && |
| 1045 |
1/2✓ Branch 0 taken 3 times.
✗ Branch 1 not taken.
|
3 | prism_leaf_refine_threshold > 0 && |
| 1046 | 3 | (uint32_t)bs->nsamples / Max(nlist, 1) < | |
| 1047 |
1/2✓ Branch 0 taken 3 times.
✗ Branch 1 not taken.
|
3 | (uint32_t)prism_leaf_refine_threshold) |
| 1048 | { | ||
| 1049 | 3 | instr_time t_ref_start; | |
| 1050 | 3 | INSTR_TIME_SET_CURRENT(t_ref_start); | |
| 1051 | 3 | prism_build_report_phase(bs->prog, PRISM_BUILD_PHASE_REFINE); | |
| 1052 | |||
| 1053 | 3 | PrismHeadWriteCtx rhead; | |
| 1054 | 6 | prism_head_write_ctx_init( | |
| 1055 | 3 | &rhead, storage, rq_params, dim, p->fastscan, first_posting); | |
| 1056 | 3 | serial_refine_heads(bs, &rhead, &bs->qs, first_posting, nlist); | |
| 1057 | 3 | prism_head_write_ctx_cleanup(&rhead); | |
| 1058 | |||
| 1059 | 3 | instr_time t_ref_end; | |
| 1060 | 3 | INSTR_TIME_SET_CURRENT(t_ref_end); | |
| 1061 | 3 | INSTR_TIME_SUBTRACT(t_ref_end, t_ref_start); | |
| 1062 |
1/2✓ Branch 1 taken 3 times.
✗ Branch 2 not taken.
|
3 | elog(LOG, |
| 1063 | "prism: page-backed leaf refinement %.1fms -- structure from a " | ||
| 1064 | "%d-sample subsample, %u leaf encode references refined on the " | ||
| 1065 | "full " | ||
| 1066 | "table", | ||
| 1067 | INSTR_TIME_GET_MILLISEC(t_ref_end), | ||
| 1068 | bs->max_samples, | ||
| 1069 | nlist); | ||
| 1070 | } | ||
| 1071 | |||
| 1072 | /* | ||
| 1073 | * Cluster-keyed sorter: the scan streams every posting entry here (keyed | ||
| 1074 | * by cluster); prism_posting_build_lists then reads them back grouped by | ||
| 1075 | * cluster and builds each list with a single resident page builder. Memory | ||
| 1076 | * is bounded by maintenance_work_mem (the sort spills if exceeded), | ||
| 1077 | * so no O(nlist) array of resident builders exists. | ||
| 1078 | * | ||
| 1079 | * region == NULL: a plain, non-parallel cluster-keyed sort. Same seam the | ||
| 1080 | * parallel leader uses. | ||
| 1081 | */ | ||
| 1082 | 151 | bs->rq_params = rq_params; | |
| 1083 | 151 | bs->sorter = prism_pbuild_sort_begin( | |
| 1084 | NULL, | ||
| 1085 | NULL, | ||
| 1086 | 0, | ||
| 1087 | 0, | ||
| 1088 | true, | ||
| 1089 | 151 | (uint32_t)prism_posting_entry_size(dim), | |
| 1090 | maintenance_work_mem); | ||
| 1091 | |||
| 1092 | /* Shared route+encode+emit context (same helper the parallel workers use). | ||
| 1093 | */ | ||
| 1094 | 151 | prism_build_route_ctx_init( | |
| 1095 | &bs->route, | ||
| 1096 | &bs->qs, | ||
| 1097 | bs->sorter, | ||
| 1098 | rq_params, | ||
| 1099 | storage, | ||
| 1100 | first_posting, | ||
| 1101 | dim, | ||
| 1102 | 151 | p->soar_lambda, | |
| 1103 | 151 | p->boundary_epsilon); | |
| 1104 | |||
| 1105 | /* Reports the scan phase and fires the "prism-build-load" test hook (see | ||
| 1106 | * the seam): lets an isolation test observe the in-progress serial build. | ||
| 1107 | */ | ||
| 1108 | 151 | prism_build_report_phase(bs->prog, PRISM_BUILD_PHASE_SCAN); | |
| 1109 | |||
| 1110 | 151 | instr_time t_serial_start; | |
| 1111 | 151 | INSTR_TIME_SET_CURRENT(t_serial_start); | |
| 1112 | |||
| 1113 | 151 | double heap_tuples = table_index_build_scan( | |
| 1114 | bs->heap, | ||
| 1115 | bs->index, | ||
| 1116 | bs->index_info, | ||
| 1117 | true, | ||
| 1118 | true, | ||
| 1119 | build_callback, | ||
| 1120 | (void *)bs, | ||
| 1121 | NULL); | ||
| 1122 | |||
| 1123 | 151 | instr_time t_serial_scan; | |
| 1124 | 151 | INSTR_TIME_SET_CURRENT(t_serial_scan); | |
| 1125 | 151 | INSTR_TIME_SUBTRACT(t_serial_scan, t_serial_start); | |
| 1126 | |||
| 1127 | /* | ||
| 1128 | * Sort entries by cluster, then build each cluster's posting list with a | ||
| 1129 | * single resident page builder, in cluster order. Empty clusters still get | ||
| 1130 | * an (empty) head page, matching the previous per-cluster behavior. | ||
| 1131 | */ | ||
| 1132 | 151 | prism_build_report_phase(bs->prog, PRISM_BUILD_PHASE_POSTING); | |
| 1133 | 151 | prism_posting_build_lists( | |
| 1134 | bs->sorter, | ||
| 1135 | storage, | ||
| 1136 | nlist, | ||
| 1137 | dim, | ||
| 1138 | 151 | p->fastscan, | |
| 1139 | rq_params, | ||
| 1140 | first_posting); | ||
| 1141 | 151 | bs->sorter = NULL; /* ended by prism_posting_build_lists */ | |
| 1142 | |||
| 1143 | 151 | bs->indtuples = bs->route.indtuples; | |
| 1144 | 151 | bs->soar_dupes = bs->route.soar_dupes; | |
| 1145 | 151 | prism_build_route_ctx_cleanup(&bs->route); | |
| 1146 | 151 | prism_query_state_cleanup(&bs->qs); | |
| 1147 |
2/2✓ Branch 0 taken 32 times.
✓ Branch 1 taken 119 times.
|
151 | if (collector != NULL) |
| 1148 | 32 | prism_exact_centroid_collector_cleanup(collector); | |
| 1149 | |||
| 1150 |
1/2✓ Branch 1 taken 151 times.
✗ Branch 2 not taken.
|
151 | elog(LOG, |
| 1151 | "prism: serial build scan %.1fms, " | ||
| 1152 | "%.0f tuples, %.0f soar_dupes, %u clusters", | ||
| 1153 | INSTR_TIME_GET_MILLISEC(t_serial_scan), | ||
| 1154 | bs->indtuples, | ||
| 1155 | bs->soar_dupes, | ||
| 1156 | nlist); | ||
| 1157 | |||
| 1158 | /* Metadata page (needs the final tuple count). first_centroid is the root | ||
| 1159 | * block, which the streaming write placed last. This path finalizes the | ||
| 1160 | * centroid pages + metadata itself, so the caller's tail does WAL only. */ | ||
| 1161 | 151 | write_meta_page( | |
| 1162 | storage, | ||
| 1163 | dim, | ||
| 1164 | 151 | (uint8_t)plan.nlevels, | |
| 1165 | 151 | (uint8_t)p->fan_out, | |
| 1166 | root, | ||
| 1167 | first_posting, | ||
| 1168 | plan.centroid_pages, | ||
| 1169 | nlist, | ||
| 1170 | 151 | p->centroid_format, | |
| 1171 | 151 | p->metric, | |
| 1172 | rabitq_seed, | ||
| 1173 | 151 | p->fastscan, | |
| 1174 | global_mean); | ||
| 1175 | 151 | *out_nlist = nlist; | |
| 1176 | 151 | *out_tree_nlevels = (uint8_t)plan.nlevels; | |
| 1177 | 151 | *out_global_mean = global_mean; | |
| 1178 | 151 | *out_heap_tuples = heap_tuples; | |
| 1179 | 151 | *out_indtuples = bs->indtuples; | |
| 1180 | 151 | *out_soar_dupes = bs->soar_dupes; | |
| 1181 | 151 | } | |
| 1182 | |||
| 1183 | /* ---------------------------------------------------------------- | ||
| 1184 | * Main build entry point | ||
| 1185 | * ---------------------------------------------------------------- */ | ||
| 1186 | |||
| 1187 | IndexBuildResult * | ||
| 1188 | 198 | prism_build(Relation heap, Relation index, struct IndexInfo *index_info) | |
| 1189 | { | ||
| 1190 | /* | ||
| 1191 | * Refuse expression indexes. | ||
| 1192 | * | ||
| 1193 | * Rerank fetches the heap tuple and reads the indexed value with | ||
| 1194 | * slot_getattr(slot, indkey.values[0]) -- and for an expression index | ||
| 1195 | * PostgreSQL stores 0 there, keeping the expression tree in indexprs | ||
| 1196 | * instead. attnum 0 reads tts_values[-1] and hands the garbage to | ||
| 1197 | * pg_detoast_datum, which segfaults the backend (measured: the crash is | ||
| 1198 | * in pg_rerank_readstream, not theoretical). | ||
| 1199 | * | ||
| 1200 | * Supporting them means doing at rerank time what the build already does | ||
| 1201 | * -- evaluating the index expression over the fetched tuple -- which | ||
| 1202 | * needs an ExprState and ExprContext live across the whole rerank loop | ||
| 1203 | * and an evaluation per candidate, in the hottest path there is. Worth | ||
| 1204 | * doing; not worth smuggling in as a crash fix. | ||
| 1205 | */ | ||
| 1206 |
2/2✓ Branch 0 taken 1 times.
✓ Branch 1 taken 197 times.
|
198 | if (index->rd_index->indkey.values[0] == InvalidAttrNumber) |
| 1207 |
1/2✓ Branch 1 taken 1 times.
✗ Branch 2 not taken.
|
1 | ereport(ERROR, |
| 1208 | (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), | ||
| 1209 | errmsg("prism indexes do not support index expressions"), | ||
| 1210 | errhint("Index the column directly, or materialize the " | ||
| 1211 | "expression into a column and index that."))); | ||
| 1212 | |||
| 1213 | 197 | MemoryContext caller_ctx = CurrentMemoryContext; | |
| 1214 | 197 | MemoryContext build_ctx = AllocSetContextCreate( | |
| 1215 | CurrentMemoryContext, "prism build", ALLOCSET_DEFAULT_SIZES); | ||
| 1216 | 197 | MemoryContextSwitchTo(build_ctx); | |
| 1217 | |||
| 1218 | /* 1. Initialize build state and resolve parameters */ | ||
| 1219 | 197 | PrismBuildState bs = {0}; | |
| 1220 | 197 | bs.heap = heap; | |
| 1221 | 197 | bs.index = index; | |
| 1222 | 197 | bs.index_info = index_info; | |
| 1223 | 197 | bs.build_ctx = build_ctx; | |
| 1224 | 197 | bs.tmp_ctx = AllocSetContextCreate( | |
| 1225 | build_ctx, "prism build tuple", ALLOCSET_DEFAULT_SIZES); | ||
| 1226 | |||
| 1227 | 197 | double est_rows; | |
| 1228 | 197 | resolve_build_params(heap, index, &bs.params, &est_rows); | |
| 1229 | |||
| 1230 | 195 | const PrismBuildParams *p = &bs.params; | |
| 1231 | 195 | Dimension dim = p->dim; | |
| 1232 | 195 | uint64_t rabitq_seed = VS_RABITQ_BUILD_SEED; | |
| 1233 | |||
| 1234 | /* Resolved once, directly -- no metadata page exists yet for the | ||
| 1235 | * per-backend cache to read. The per-tuple callbacks follow the pointer. | ||
| 1236 | */ | ||
| 1237 | 195 | bs.input = vec32_access(prism_index_type_info(index), dim, build_ctx); | |
| 1238 | |||
| 1239 | /* | ||
| 1240 | * Build introspection: one reporting context the serial and parallel paths | ||
| 1241 | * share, so pg_stat_progress_create_index advances through the same phases | ||
| 1242 | * either way and (when prism.log_build_stats is on) per-phase stats land | ||
| 1243 | * in the server log. | ||
| 1244 | */ | ||
| 1245 | 195 | PrismBuildStats stats = {0}; | |
| 1246 | 195 | PrismBuildProgress prog; | |
| 1247 | 195 | prism_build_progress_begin( | |
| 1248 | &prog, | ||
| 1249 | 195 | index_info->ii_ParallelWorkers > 0, | |
| 1250 | prism_log_build_stats, | ||
| 1251 | build_ctx, | ||
| 1252 | &stats, | ||
| 1253 | est_rows); | ||
| 1254 | 195 | bs.prog = &prog; | |
| 1255 | |||
| 1256 | 195 | VsPgStorage storage; | |
| 1257 | 195 | vs_pg_storage_init(&storage, index, NULL, p->metric); | |
| 1258 | 195 | storage.build_mode = true; | |
| 1259 | |||
| 1260 | 195 | double heap_tuples = 0; | |
| 1261 | 195 | double indtuples = 0; | |
| 1262 | 195 | double soar_dupes = 0; | |
| 1263 | 195 | float *global_mean = NULL; | |
| 1264 | /* Posting-area start block, surfaced by the parallel build so the finalize | ||
| 1265 | * can record it in the metadata page (vacuum skips the centroid region). | ||
| 1266 | */ | ||
| 1267 | 195 | BlockNumber meta_first_posting = InvalidBlockNumber; | |
| 1268 | |||
| 1269 | /* The streamed tree's shape, from whichever build path ran (0 leaves = | ||
| 1270 | * empty heap). The serial path writes its own metadata page; the | ||
| 1271 | * parallel path leaves it for the finalize below (it needs the final | ||
| 1272 | * tuple count). */ | ||
| 1273 | 195 | uint32_t built_nlist = 0; | |
| 1274 | 195 | uint8_t built_nlevels = 0; | |
| 1275 | |||
| 1276 | /* Try parallel build first (sampling + k-means + posting) */ | ||
| 1277 | 195 | bool did_parallel = false; | |
| 1278 |
2/2✓ Branch 0 taken 45 times.
✓ Branch 1 taken 150 times.
|
195 | if (index_info->ii_ParallelWorkers > 0) |
| 1279 | { | ||
| 1280 | /* | ||
| 1281 | * Estimate the leaf count up front. | ||
| 1282 | * | ||
| 1283 | * A parallel build allocates its shared-memory (DSM) regions before | ||
| 1284 | * the workers run, and DSM segments cannot be resized once created. | ||
| 1285 | * Several of those regions are sized per leaf (centroids, posting | ||
| 1286 | * heads, per-cluster assignment state), so the leader has to commit to | ||
| 1287 | * a leaf count at allocation time — but the real count is only known | ||
| 1288 | * after k-means clusters the sample, which happens inside the | ||
| 1289 | * workers. hkmeans targets `nlist` leaves at this fan_out but can | ||
| 1290 | * produce up to fan_out^nlevels of them, so we size every per-leaf | ||
| 1291 | * region for that worst-case upper bound here, then narrow to the | ||
| 1292 | * actual built leaf count below. | ||
| 1293 | * | ||
| 1294 | * nlist and fan_out are already resolved (resolve_build_params). | ||
| 1295 | */ | ||
| 1296 | 45 | uint32_t requested_nlist = p->nlist; | |
| 1297 | 45 | uint32_t max_nlist = prism_max_nlist(p->nlist, p->fan_out); | |
| 1298 | |||
| 1299 | 45 | bs.params.nlist = max_nlist; | |
| 1300 | |||
| 1301 | 45 | PrismBuildConfig cfg = { | |
| 1302 | 45 | .dim = bs.params.dim, | |
| 1303 | 45 | .metric = bs.params.metric, | |
| 1304 | 45 | .centroid_format = bs.params.centroid_format, | |
| 1305 | .nlist = bs.params.nlist, | ||
| 1306 | 45 | .fan_out = bs.params.fan_out, | |
| 1307 | 45 | .soar_lambda = bs.params.soar_lambda, | |
| 1308 | 45 | .boundary_epsilon = bs.params.boundary_epsilon, | |
| 1309 | 45 | .fastscan = bs.params.fastscan, | |
| 1310 | 45 | .concurrent = index_info->ii_Concurrent, | |
| 1311 | }; | ||
| 1312 | |||
| 1313 | /* do_parallel_build reports every phase through the seam (sampling, | ||
| 1314 | * k-means, subtrees, refine, setup, scan, posting) and fires | ||
| 1315 | * the per-phase test hooks, so no phase is set here. */ | ||
| 1316 | 45 | did_parallel = do_parallel_build( | |
| 1317 | heap, | ||
| 1318 | index, | ||
| 1319 | index_info, | ||
| 1320 | &cfg, | ||
| 1321 | &storage.base, | ||
| 1322 | &prog, | ||
| 1323 | &built_nlist, | ||
| 1324 | &built_nlevels, | ||
| 1325 | &heap_tuples, | ||
| 1326 | &indtuples, | ||
| 1327 | &soar_dupes, | ||
| 1328 | /* do_parallel_build routes page-backed: it wrote the centroid | ||
| 1329 | * + head pages into the index before the scan and returns the | ||
| 1330 | * global mean it used (so the metadata write below matches). | ||
| 1331 | * It does NOT write the metadata page (needs the final tuple | ||
| 1332 | * count), so the finalize below still runs. */ | ||
| 1333 | &global_mean, | ||
| 1334 | /* Surfaced for the metadata page: the posting-area start block | ||
| 1335 | * lets vacuum skip the centroid region without scanning it. */ | ||
| 1336 | &meta_first_posting); | ||
| 1337 | |||
| 1338 |
3/4✓ Branch 0 taken 43 times.
✓ Branch 1 taken 1 times.
✗ Branch 2 not taken.
✓ Branch 3 taken 43 times.
|
44 | if (did_parallel && built_nlist == 0) |
| 1339 | { | ||
| 1340 | /* | ||
| 1341 | * The parallel path found no indexable rows (a non-empty heap | ||
| 1342 | * whose rows are all dead or NULL can reach it). It wrote no | ||
| 1343 | * pages for the zero-leaf layout, so the serial fallback below | ||
| 1344 | * -- whose sampler substitutes a synthetic sample and emits a | ||
| 1345 | * valid single-cluster index -- starts from a clean relation. | ||
| 1346 | * Enforced at runtime: falling back onto a relation the | ||
| 1347 | * parallel path already wrote to would corrupt the index. | ||
| 1348 | */ | ||
| 1349 | ✗ | if (RelationGetNumberOfBlocks(index) > 1) | |
| 1350 | ✗ | ereport(ERROR, | |
| 1351 | (errcode(ERRCODE_INTERNAL_ERROR), | ||
| 1352 | errmsg("parallel index build reported no leaves " | ||
| 1353 | "but wrote %u pages", | ||
| 1354 | RelationGetNumberOfBlocks(index)))); | ||
| 1355 | |||
| 1356 | /* The zero-leaf return leaves out_global_mean NULL today; | ||
| 1357 | * free defensively so a future change to that contract | ||
| 1358 | * cannot leak the parallel mean when the serial build | ||
| 1359 | * replaces it. */ | ||
| 1360 | ✗ | if (global_mean != NULL) | |
| 1361 | { | ||
| 1362 | ✗ | pfree(global_mean); | |
| 1363 | ✗ | global_mean = NULL; | |
| 1364 | } | ||
| 1365 | did_parallel = false; | ||
| 1366 | } | ||
| 1367 | |||
| 1368 | 43 | if (did_parallel && built_nlist > 0) | |
| 1369 | 43 | bs.params.nlist = built_nlist; | |
| 1370 | else if (!did_parallel) | ||
| 1371 | { | ||
| 1372 | /* The worst-case bound exists only to size the parallel DSM | ||
| 1373 | * regions. The serial fallback clusters from scratch, so it must | ||
| 1374 | * target the requested partition count, not the sizing bound. */ | ||
| 1375 | 1 | bs.params.nlist = requested_nlist; | |
| 1376 | } | ||
| 1377 | } | ||
| 1378 | |||
| 1379 |
2/2✓ Branch 0 taken 1 times.
✓ Branch 1 taken 43 times.
|
44 | if (!did_parallel) |
| 1380 | { | ||
| 1381 | /* Serial fallback: sample, cluster, build. Always emits at least | ||
| 1382 | * one leaf, even when the heap has no indexable tuples. */ | ||
| 1383 | 151 | do_serial_build( | |
| 1384 | &bs, | ||
| 1385 | &storage.base, | ||
| 1386 | rabitq_seed, | ||
| 1387 | &built_nlist, | ||
| 1388 | &built_nlevels, | ||
| 1389 | &global_mean, | ||
| 1390 | &heap_tuples, | ||
| 1391 | &indtuples, | ||
| 1392 | &soar_dupes); | ||
| 1393 | } | ||
| 1394 | |||
| 1395 | /* Both paths produce at least one leaf now: a heap with no indexable | ||
| 1396 | * rows builds a single-cluster index around a synthetic centroid (see | ||
| 1397 | * sample_for_build), so inserts route and scans return empty. | ||
| 1398 | * Enforced at runtime: finalizing a zero-leaf layout would produce a | ||
| 1399 | * broken index. */ | ||
| 1400 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 194 times.
|
194 | if (built_nlist == 0) |
| 1401 | ✗ | ereport(ERROR, | |
| 1402 | (errcode(ERRCODE_INTERNAL_ERROR), | ||
| 1403 | errmsg("index build produced no clusters"))); | ||
| 1404 | |||
| 1405 |
2/2✓ Branch 0 taken 128 times.
✓ Branch 1 taken 66 times.
|
194 | if (soar_dupes > 0) |
| 1406 |
1/2✓ Branch 1 taken 128 times.
✗ Branch 2 not taken.
|
128 | elog(LOG, |
| 1407 | "prism: replicated %.0f vectors " | ||
| 1408 | "(%.1f%% of %.0f, lambda=%.4g)", | ||
| 1409 | soar_dupes, | ||
| 1410 | 100.0 * soar_dupes / indtuples, | ||
| 1411 | indtuples, | ||
| 1412 | bs.params.soar_lambda); | ||
| 1413 | |||
| 1414 | /* Finalize: write centroid pages + metadata (unless the build path already | ||
| 1415 | * did — the serial page-backed path writes them itself), then WAL-log the | ||
| 1416 | * whole index. */ | ||
| 1417 | { | ||
| 1418 |
2/2✓ Branch 0 taken 43 times.
✓ Branch 1 taken 151 times.
|
194 | if (did_parallel) |
| 1419 | { | ||
| 1420 | /* The parallel build wrote its centroid + head pages before the | ||
| 1421 | * scan but left the metadata page for here (it needs the final | ||
| 1422 | * tuple count); the serial path wrote its own. first_posting | ||
| 1423 | * comes from the build and lets vacuum skip the centroid region | ||
| 1424 | * without scanning it. */ | ||
| 1425 | 43 | BlockNumber fc = 1; | |
| 1426 | |||
| 1427 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 43 times.
|
43 | Assert(global_mean != NULL); |
| 1428 | |||
| 1429 | 43 | write_meta_page( | |
| 1430 | &storage.base, | ||
| 1431 | dim, | ||
| 1432 | built_nlevels, | ||
| 1433 | 43 | (uint8_t)p->fan_out, | |
| 1434 | fc, | ||
| 1435 | meta_first_posting, | ||
| 1436 | meta_first_posting - PRISM_FIRST_CENTROID_BLKNO, | ||
| 1437 | built_nlist, | ||
| 1438 | 43 | p->centroid_format, | |
| 1439 | 43 | p->metric, | |
| 1440 | rabitq_seed, | ||
| 1441 | 43 | p->fastscan, | |
| 1442 | global_mean); | ||
| 1443 | } | ||
| 1444 | 194 | prism_build_report_phase(&prog, PRISM_BUILD_PHASE_WAL); | |
| 1445 | /* | ||
| 1446 | * WAL-log the built pages only when the relation is WAL-logged, the | ||
| 1447 | * same guard every core index AM (GiST/GIN/SP-GiST) applies here. | ||
| 1448 | * Temp relations and same-transaction builds under wal_level=minimal | ||
| 1449 | * skip this: emitting WAL for them is wasted volume, and doing so | ||
| 1450 | * against a session-local temp relfilenode is something no core AM | ||
| 1451 | * does. (Unlogged tables never reach this path -- their build is | ||
| 1452 | * rejected up front in prism_buildempty.) | ||
| 1453 | * | ||
| 1454 | * page_std = false, and it has to be. The flag promises the standard | ||
| 1455 | * page layout, which lets the full-page image omit everything between | ||
| 1456 | * pd_lower and pd_upper as free space. prism keeps its entries there | ||
| 1457 | * -- pd_lower stays at SizeOfPageHeaderData and pd_upper at pd_special | ||
| 1458 | * on a posting page -- so a standard image would carry the header and | ||
| 1459 | * the opaque and nothing else. The primary would be fine, since its | ||
| 1460 | * pages are already on disk; a streaming standby, which has only the | ||
| 1461 | * WAL, would restore every page with its entries zeroed while the | ||
| 1462 | * opaque still reported them present. See t/001_replication.pl. | ||
| 1463 | * | ||
| 1464 | * The insert path solves the same problem the other way, by covering | ||
| 1465 | * the hole (pd_lower = pd_upper) before GenericXLog; that is not an | ||
| 1466 | * option here because centroid pages use pd_lower for their own | ||
| 1467 | * capacity accounting. | ||
| 1468 | */ | ||
| 1469 |
5/8✓ Branch 0 taken 192 times.
✓ Branch 1 taken 2 times.
✓ Branch 2 taken 1 times.
✓ Branch 3 taken 191 times.
✗ Branch 4 not taken.
✓ Branch 5 taken 1 times.
✗ Branch 6 not taken.
✗ Branch 7 not taken.
|
194 | if (RelationNeedsWAL(index)) |
| 1470 | 191 | log_newpage_range( | |
| 1471 | index, | ||
| 1472 | MAIN_FORKNUM, | ||
| 1473 | 0, | ||
| 1474 | RelationGetNumberOfBlocks(index), | ||
| 1475 | false); | ||
| 1476 | } | ||
| 1477 | |||
| 1478 | /* Flush the final phase timing + emit the build summary (heap_ctx is read | ||
| 1479 | * here, so this must run before build_ctx is deleted below). */ | ||
| 1480 | 194 | prism_build_progress_end(&prog); | |
| 1481 | |||
| 1482 | /* Cleanup */ | ||
| 1483 | |||
| 1484 | 194 | MemoryContextSwitchTo(caller_ctx); | |
| 1485 | 194 | MemoryContextDelete(build_ctx); | |
| 1486 | |||
| 1487 | 194 | IndexBuildResult *result = palloc0(sizeof(IndexBuildResult)); | |
| 1488 | 194 | result->heap_tuples = heap_tuples; | |
| 1489 | 194 | result->index_tuples = indtuples; | |
| 1490 | 194 | return result; | |
| 1491 | } | ||
| 1492 | |||
| 1493 | char * | ||
| 1494 | 7 | prism_buildphasename(int64 phasenum) | |
| 1495 | { | ||
| 1496 | /* Single source of truth for the phase names (index/build_progress.c), | ||
| 1497 | * shared with the build logs so the two never drift. */ | ||
| 1498 | 7 | return unconstify(char *, prism_build_phase_name((int)phasenum)); | |
| 1499 | } | ||
| 1500 |