chrevanthreddy commented on code in PR #19309: URL: https://github.com/apache/hudi/pull/19309#discussion_r3677438520
########## rfc/rfc-109/rfc-109.md: ########## @@ -0,0 +1,594 @@ +<!-- + Licensed to the Apache Software Foundation (ASF) under one or more + contributor license agreements. See the NOTICE file distributed with + this work for additional information regarding copyright ownership. + The ASF licenses this file to You under the Apache License, Version 2.0 + (the "License"); you may not use this file except in compliance with + the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. +--> + +# RFC-109: Native Vector Search Support in Apache Hudi + +## Proposers + +@chrevanthreddy + +## Approvers + +- TBD + +## Status + +Umbrella issue: [apache/hudi#19094](https://github.com/apache/hudi/issues/19094) + +Related: [apache/hudi#18676](https://github.com/apache/hudi/issues/18676) + +State: UNDER REVIEW + +--- + +## Table of Contents + +- [Abstract](#abstract) +- [1. Goals and Non-Goals](#1-goals-and-non-goals) +- [2. Architecture](#2-architecture) +- [3. IVF + RaBitQ Index Algorithm](#3-ivf--rabitq-index-algorithm) +- [4. Metadata Table Storage Model: the Posting Block](#4-metadata-table-storage-model-the-posting-block) +- [5. Bootstrap and Write Path](#5-bootstrap-and-write-path) +- [6. Read Path](#6-read-path) +- [7. Maintenance, Rebalancing, and Cleaner](#7-maintenance-rebalancing-and-cleaner) +- [8. Spark API Surface](#8-spark-api-surface) +- [9. Correctness and Consistency](#9-correctness-and-consistency) +- [10. Test Plan](#10-test-plan) +- [11. Rollout and MVP Scope](#11-rollout-and-mvp-scope) +- [12. References](#12-references) + +--- + +## Abstract + +This RFC proposes native approximate nearest-neighbor (ANN) vector search in Apache Hudi. +Tables increasingly carry embedding columns (`ARRAY<FLOAT>` produced by ML models) next to +their business data, and users want to ask *"find the K rows most similar to this query +vector"* — for semantic search, recommendations, RAG, and deduplication — without copying +data into a separate vector database. + +Today the only option on a Hudi table is a brute-force scan: read every vector, compute +every distance. That is correct but scales linearly with table size (tens of seconds at a +billion rows). This RFC adds an index so that vector queries read only a small, targeted +fraction of the index and the table, return results with high recall, and stay +transactionally consistent with the table under upserts and deletes — all with **no new +storage system**. The index lives in the Hudi Metadata Table (MDT), like Hudi's existing +record-level and secondary indexes, and is maintained by the same table services. + +The design combines three well-understood pieces — IVF clustering, RaBitQ quantization, and +exact re-ranking — with one storage innovation that makes them practical on an immutable, +columnar, object-store-resident lakehouse: + +> **The posting block.** Instead of one MDT record per indexed vector, the index packs +> ~1–4K vectors into a single MDT record laid out column-wise (structure-of-arrays), keyed +> so that one IVF cluster forms one contiguous, prefix-scannable key range. This reduces MDT +> record count by roughly three orders of magnitude, turns "scan a cluster" into a single +> contiguous range read, and lets a query touch only the columns a given scan pass needs. + +The base table remains the source of truth for exact vector values. The MDT stores only +routing, pruning, and approximate-scoring metadata; final ranking always reads exact vectors +from the base table. + +Prototype measurements on a **1-billion-row, 128-dimensional table** show exact-reranked +recall@10 = 0.985 at nprobe=128 with query latency in low single-digit seconds on a modest +Spark cluster, versus ~15s for brute force. + +--- + +## 1. Goals and Non-Goals + +### 1.1 Goals + +1. Keep authoritative vector values in the base table (`ARRAY<FLOAT>` / `VECTOR(D)` column). +2. Store the vector index in the MDT, maintained by Hudi metadata-table commits, compaction, + and cleaning — no hidden or generated columns in base-table files. +3. Make candidate discovery cheap and targeted: probe a few clusters, scan contiguous key + ranges, score on compressed codes with a provable pruning bound. +4. Make results trustworthy: approximate math selects candidates; **exact** distance on + base-table vectors ranks them. +5. Stay transactionally consistent: snapshot-pinned reads, correct behavior under inserts, + updates, deletes, and clustering. +6. Be engine-neutral in design; Spark is the first implementation. +7. Maintain the index incrementally (no global rebuild for normal churn) and support + versioned, zero-downtime rebuilds. + +### 1.2 Non-Goals (initial landing) + +- ANN families beyond IVF + RaBitQ (e.g. HNSW, DiskANN). +- Filtered search (arbitrary predicate + kNN) as a first-class planned operation. +- Time-travel-consistent index reads for historical snapshots. +- Engines beyond Spark, and GPU-accelerated encoding. +- Workload-specific auto-tuning of `nprobe` / refine factor. + +--- + +## 2. Architecture + + + +The design splits responsibilities the way Hudi already does between the data table and the +metadata table: + +```text +DATA TABLE (parquet/orc) METADATA TABLE (vector_index partition) + authoritative vectors + payload ←── the index: centroids, quantizer, posting blocks, + read only for final re-ranking cluster manifests, generation manifest + read for candidate generation +``` + +Each vector index is one MDT partition. Creating an index: + +```sql +CREATE INDEX embedding_idx +ON products +USING VECTOR (embedding) +OPTIONS ( + 'vector.dimension' = '768', + 'vector.metric' = 'cosine', + 'vector.quantizer' = 'IVF_RABITQ', + 'vector.num_clusters'= '4096' +); +``` + +creates: + +```text +.hoodie/metadata/vector_index_embedding_idx/ +``` + +and does not change the base-table schema: + +```text +products/category=electronics/<file-group>.parquet +├── _hoodie_record_key = p001 +├── _hoodie_partition_path = category=electronics +├── user columns = ... +└── embedding ARRAY<FLOAT> = [0.12, -0.08, ...] +``` + +The query path uses the MDT first to discover candidates, then reads base-table vectors for +exact re-ranking: + +```text +query vector + → compare to centroids, pick nprobe clusters (in-memory, ms) + → MDT prefix-scan those clusters' posting blocks (targeted range reads) + → two-pass RaBitQ scoring, keep refineFactor·K best (bit math + error bounds) + → validate candidate freshness via Record Level Index (batched point lookups) + → fetch ONLY those rows from the base table by position (page-level reads) + → exact distance on real vectors → final top-K +``` + +--- + +## 3. IVF + RaBitQ Index Algorithm + +Every practical ANN index answers two questions: **where to look** (avoid scanning +everything) and **how to compare cheaply** (avoid full-precision math on what is scanned). + +### 3.1 Where to look: IVF routing + +Inverted File (IVF) indexing clusters vectors with KMeans into `numClusters` groups (e.g. +~4K–64K). Each vector belongs to its nearest centroid. A query compares against the centroids +only (thousands, not billions), selects the `nprobe` nearest clusters, and scans only those +clusters' entries. `nprobe` is the recall dial. + +IVF is the right fit for a lakehouse-resident index because a cluster's entries can be stored +**contiguously**, which maps directly onto sorted key ranges in Hudi's MDT (§4). Graph +indexes (HNSW) give excellent in-memory recall but require random traversal of the whole +graph, which fights columnar, immutable, object-store storage. + +### 3.2 How to compare cheaply: RaBitQ quantization + +Inside a probed cluster there are still thousands of full vectors. Quantization stores a +small *code* per vector plus a few correction scalars, so most comparisons run on compressed +codes and only the best few hundred candidates are re-checked against real vectors. + +RaBitQ is chosen over scalar (SQ), product (PQ), and plain binary quantization for four +reasons: + +1. **Unbiased estimator with a provable per-vector error bound.** Each code carries scalars + that turn a cheap bit-level dot product into an *unbiased* estimate of the true distance, + plus a bound on how wrong it can be. The bound enables **safe pruning**: skip a vector + only when even its best plausible distance cannot make top-K. +2. **No codebooks.** RaBitQ needs only a random rotation (a seed) and the centroids — both + tiny, both versioned in the index metadata. Nothing to retrain when data drifts. +3. **Tunable precision.** B = 1 bit/dim is a fast coarse filter; B = 4 bits/dim gives + near-SQ quality at ~8× less space. Both are used in a two-pass scan (§6.2). +4. **Metric-flexible.** One stored code serves L2, cosine, and dot-product; the metric is + applied at query time. + +The four ideas, precisely: + +- **Residual.** Store each vector as its difference from its centroid, `r = v − c`. + Residuals are small and centered, which is what lets few bits go far. +- **Rotation.** Apply one fixed random orthonormal rotation `R` (derived from a per-generation + seed) to everything first: `x = R·r`. This spreads information evenly across dimensions, + which is what makes the error bound hold for any data distribution. Only the seed is stored. +- **Code.** Quantize `x` to B bits per dimension, stored as **bit planes** — plane 0 holds + bit 0 of every dimension, plane 1 bit 1, etc. The top plane alone is a 1-bit sign sketch. + Scoring a plane against the (transformed) query is `AND`/`XOR` + `popcount` — a few CPU + instructions per 64 dimensions. +- **Factors.** A handful of small scalars per vector (residual norm, two rescale factors, an + error term, a centroid correction) that convert plane math into an unbiased distance + estimate plus its confidence interval. + +At query time the query vector is transformed the same way (`R·(q − c)` per probed cluster), +and the per-plane popcounts combine with the stored factors into the estimate and its bound. +The estimate builds the shortlist; exact base-table distances produce the final ranking (§6). + +### 3.3 Why this fits Hudi + +- Centroids are small enough to load at planning time: `K × D` floats (~12 MB for K=4096, + D=768). +- Codes are compact: a 1B × 128-dim float table's raw vectors are ~512 GB; the RaBitQ index + including keys and locators is ~136 GB, and the scanned portion per query is tens of MB. +- Posting keys are prefix-scannable by generation, cluster, and shard (§4). +- Quantizer state is stable: a seed and centroids, no per-generation learned codebook. +- Exact re-ranking preserves correctness for returned candidates. + +--- + +## 4. Metadata Table Storage Model: the Posting Block + +This section describes the core storage contribution of this RFC. + +### 4.1 The posting block + +A naive index would write one MDT record per indexed vector. At a billion rows that is a +billion MDT records per generation — prohibitive to write, compact, scan, and clean. + +Instead, the index sorts entries by `(cluster, fileGroup, rowPosition)` and packs +~1–4K vectors' worth of codes and metadata into a single MDT record — a **posting block** — +laid out column-wise (structure-of-arrays) so each scan pass touches only the columns it +needs: + +```text +POSTING BLOCK (~512 KB target, one MDT record) +┌───────────────────────────────────────────────┐ +│ S1 sign planes ← pass 1 touches this │ +│ S2 extra bit planes ← pass 2, survivors │ +│ S3 factor arrays ← both passes │ +│ S4 row locators ← finalists only │ +│ S5 dictionaries (file groups, partitions) │ +│ S6 record keys ← finalists only │ +└───────────────────────────────────────────────┘ +key: 0x10 | generation | clusterId | shardId | blockId +``` + +Consequences: + +- **~1000× fewer MDT records.** One block record replaces ~1–4K per-vector records, cutting + write amplification, compaction cost, and cleaner load by three orders of magnitude. +- **Contiguous cluster scans.** The binary key scheme makes "scan cluster 12345" a single + contiguous HFile range read rather than thousands of point lookups. +- **Pay only for promise.** Column-wise layout means pass 1 reads only sign planes + factors; + only survivors touch extra planes; only finalists touch locators and keys (§6.2). + +### 4.2 Row families + +The `vector_index_<name>` partition holds several record families under one binary-sorted key +scheme, so one prefix scan of a cluster returns its blocks and any fresh deltas together: + +| Key family | Cardinality | Purpose | +|---|---:|---| +| `__manifest__` | 1 | Active generation pointer used by readers. | +| `__centroids__` | 1 | Serialized `K × D` centroid matrix. | +| `__quantizer__` | 1 | RaBitQ type, bit width `B`, rotation seed, normalization flag. | +| `M\|<generation>` | generations | Immutable generation-level metadata. | +| `C\|<generation>\|<cluster>` | K per generation | Cluster manifest: shard count, vector count, candidate file groups, delta/tombstone counters. | +| `P\|<generation>\|<cluster>\|<shard>\|<blockId>` | blocks per generation | **Posting block** (packed codes, factors, locators, keys). | +| `P\|...\|<DELTA>` | deltas | Small per-record delta records appended to a cluster range between compactions. | + +### 4.3 Posting shards + +Large clusters are split into posting shards so one hot cluster does not become one oversized +prefix range. The cluster manifest stores `shardCount`; writers compute +`shardId = hash(record_key) % shardCount`. Within a shard, entries are packed into blocks by Review Comment: Addressed in §4.3 and §4.5 in commit `b902152fa6`. `shardCount` is not changed in place while old and new mappings coexist. Because `hash(record_key) % shardCount` remaps essentially every key, maintenance must rewrite the whole cluster under a new `routingVersion` and publish the cluster manifest atomically with that rewrite. Readers therefore never combine two routing versions for one cluster. For a large remap, the revision requires a new complete generation rather than an in-place rewrite. Generation remains the only independent version axis; unchanged rows may be copied under the new generation id, while the affected cluster is re-encoded with the new shard mapping before atomic activation. ########## rfc/rfc-109/rfc-109.md: ########## @@ -0,0 +1,594 @@ +<!-- + Licensed to the Apache Software Foundation (ASF) under one or more + contributor license agreements. See the NOTICE file distributed with + this work for additional information regarding copyright ownership. + The ASF licenses this file to You under the Apache License, Version 2.0 + (the "License"); you may not use this file except in compliance with + the License. You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. +--> + +# RFC-109: Native Vector Search Support in Apache Hudi + +## Proposers + +@chrevanthreddy + +## Approvers + +- TBD + +## Status + +Umbrella issue: [apache/hudi#19094](https://github.com/apache/hudi/issues/19094) + +Related: [apache/hudi#18676](https://github.com/apache/hudi/issues/18676) + +State: UNDER REVIEW + +--- + +## Table of Contents + +- [Abstract](#abstract) +- [1. Goals and Non-Goals](#1-goals-and-non-goals) +- [2. Architecture](#2-architecture) +- [3. IVF + RaBitQ Index Algorithm](#3-ivf--rabitq-index-algorithm) +- [4. Metadata Table Storage Model: the Posting Block](#4-metadata-table-storage-model-the-posting-block) +- [5. Bootstrap and Write Path](#5-bootstrap-and-write-path) +- [6. Read Path](#6-read-path) +- [7. Maintenance, Rebalancing, and Cleaner](#7-maintenance-rebalancing-and-cleaner) +- [8. Spark API Surface](#8-spark-api-surface) +- [9. Correctness and Consistency](#9-correctness-and-consistency) +- [10. Test Plan](#10-test-plan) +- [11. Rollout and MVP Scope](#11-rollout-and-mvp-scope) +- [12. References](#12-references) + +--- + +## Abstract + +This RFC proposes native approximate nearest-neighbor (ANN) vector search in Apache Hudi. +Tables increasingly carry embedding columns (`ARRAY<FLOAT>` produced by ML models) next to +their business data, and users want to ask *"find the K rows most similar to this query +vector"* — for semantic search, recommendations, RAG, and deduplication — without copying +data into a separate vector database. + +Today the only option on a Hudi table is a brute-force scan: read every vector, compute +every distance. That is correct but scales linearly with table size (tens of seconds at a +billion rows). This RFC adds an index so that vector queries read only a small, targeted +fraction of the index and the table, return results with high recall, and stay +transactionally consistent with the table under upserts and deletes — all with **no new +storage system**. The index lives in the Hudi Metadata Table (MDT), like Hudi's existing +record-level and secondary indexes, and is maintained by the same table services. + +The design combines three well-understood pieces — IVF clustering, RaBitQ quantization, and +exact re-ranking — with one storage innovation that makes them practical on an immutable, +columnar, object-store-resident lakehouse: + +> **The posting block.** Instead of one MDT record per indexed vector, the index packs +> ~1–4K vectors into a single MDT record laid out column-wise (structure-of-arrays), keyed +> so that one IVF cluster forms one contiguous, prefix-scannable key range. This reduces MDT +> record count by roughly three orders of magnitude, turns "scan a cluster" into a single +> contiguous range read, and lets a query touch only the columns a given scan pass needs. + +The base table remains the source of truth for exact vector values. The MDT stores only +routing, pruning, and approximate-scoring metadata; final ranking always reads exact vectors +from the base table. + +Prototype measurements on a **1-billion-row, 128-dimensional table** show exact-reranked +recall@10 = 0.985 at nprobe=128 with query latency in low single-digit seconds on a modest +Spark cluster, versus ~15s for brute force. + +--- + +## 1. Goals and Non-Goals + +### 1.1 Goals + +1. Keep authoritative vector values in the base table (`ARRAY<FLOAT>` / `VECTOR(D)` column). +2. Store the vector index in the MDT, maintained by Hudi metadata-table commits, compaction, + and cleaning — no hidden or generated columns in base-table files. +3. Make candidate discovery cheap and targeted: probe a few clusters, scan contiguous key + ranges, score on compressed codes with a provable pruning bound. +4. Make results trustworthy: approximate math selects candidates; **exact** distance on + base-table vectors ranks them. +5. Stay transactionally consistent: snapshot-pinned reads, correct behavior under inserts, + updates, deletes, and clustering. +6. Be engine-neutral in design; Spark is the first implementation. +7. Maintain the index incrementally (no global rebuild for normal churn) and support + versioned, zero-downtime rebuilds. + +### 1.2 Non-Goals (initial landing) + +- ANN families beyond IVF + RaBitQ (e.g. HNSW, DiskANN). +- Filtered search (arbitrary predicate + kNN) as a first-class planned operation. +- Time-travel-consistent index reads for historical snapshots. +- Engines beyond Spark, and GPU-accelerated encoding. +- Workload-specific auto-tuning of `nprobe` / refine factor. + +--- + +## 2. Architecture + + + +The design splits responsibilities the way Hudi already does between the data table and the +metadata table: + +```text +DATA TABLE (parquet/orc) METADATA TABLE (vector_index partition) + authoritative vectors + payload ←── the index: centroids, quantizer, posting blocks, + read only for final re-ranking cluster manifests, generation manifest + read for candidate generation +``` + +Each vector index is one MDT partition. Creating an index: + +```sql +CREATE INDEX embedding_idx +ON products +USING VECTOR (embedding) +OPTIONS ( + 'vector.dimension' = '768', + 'vector.metric' = 'cosine', + 'vector.quantizer' = 'IVF_RABITQ', + 'vector.num_clusters'= '4096' +); +``` + +creates: + +```text +.hoodie/metadata/vector_index_embedding_idx/ +``` + +and does not change the base-table schema: + +```text +products/category=electronics/<file-group>.parquet +├── _hoodie_record_key = p001 +├── _hoodie_partition_path = category=electronics +├── user columns = ... +└── embedding ARRAY<FLOAT> = [0.12, -0.08, ...] +``` + +The query path uses the MDT first to discover candidates, then reads base-table vectors for +exact re-ranking: + +```text +query vector + → compare to centroids, pick nprobe clusters (in-memory, ms) + → MDT prefix-scan those clusters' posting blocks (targeted range reads) + → two-pass RaBitQ scoring, keep refineFactor·K best (bit math + error bounds) + → validate candidate freshness via Record Level Index (batched point lookups) + → fetch ONLY those rows from the base table by position (page-level reads) + → exact distance on real vectors → final top-K +``` + +--- + +## 3. IVF + RaBitQ Index Algorithm + +Every practical ANN index answers two questions: **where to look** (avoid scanning +everything) and **how to compare cheaply** (avoid full-precision math on what is scanned). + +### 3.1 Where to look: IVF routing + +Inverted File (IVF) indexing clusters vectors with KMeans into `numClusters` groups (e.g. +~4K–64K). Each vector belongs to its nearest centroid. A query compares against the centroids +only (thousands, not billions), selects the `nprobe` nearest clusters, and scans only those +clusters' entries. `nprobe` is the recall dial. + +IVF is the right fit for a lakehouse-resident index because a cluster's entries can be stored +**contiguously**, which maps directly onto sorted key ranges in Hudi's MDT (§4). Graph +indexes (HNSW) give excellent in-memory recall but require random traversal of the whole +graph, which fights columnar, immutable, object-store storage. + +### 3.2 How to compare cheaply: RaBitQ quantization + +Inside a probed cluster there are still thousands of full vectors. Quantization stores a +small *code* per vector plus a few correction scalars, so most comparisons run on compressed +codes and only the best few hundred candidates are re-checked against real vectors. + +RaBitQ is chosen over scalar (SQ), product (PQ), and plain binary quantization for four +reasons: + +1. **Unbiased estimator with a provable per-vector error bound.** Each code carries scalars + that turn a cheap bit-level dot product into an *unbiased* estimate of the true distance, + plus a bound on how wrong it can be. The bound enables **safe pruning**: skip a vector + only when even its best plausible distance cannot make top-K. +2. **No codebooks.** RaBitQ needs only a random rotation (a seed) and the centroids — both + tiny, both versioned in the index metadata. Nothing to retrain when data drifts. +3. **Tunable precision.** B = 1 bit/dim is a fast coarse filter; B = 4 bits/dim gives + near-SQ quality at ~8× less space. Both are used in a two-pass scan (§6.2). +4. **Metric-flexible.** One stored code serves L2, cosine, and dot-product; the metric is + applied at query time. + +The four ideas, precisely: + +- **Residual.** Store each vector as its difference from its centroid, `r = v − c`. + Residuals are small and centered, which is what lets few bits go far. +- **Rotation.** Apply one fixed random orthonormal rotation `R` (derived from a per-generation + seed) to everything first: `x = R·r`. This spreads information evenly across dimensions, + which is what makes the error bound hold for any data distribution. Only the seed is stored. +- **Code.** Quantize `x` to B bits per dimension, stored as **bit planes** — plane 0 holds + bit 0 of every dimension, plane 1 bit 1, etc. The top plane alone is a 1-bit sign sketch. + Scoring a plane against the (transformed) query is `AND`/`XOR` + `popcount` — a few CPU + instructions per 64 dimensions. +- **Factors.** A handful of small scalars per vector (residual norm, two rescale factors, an + error term, a centroid correction) that convert plane math into an unbiased distance + estimate plus its confidence interval. + +At query time the query vector is transformed the same way (`R·(q − c)` per probed cluster), +and the per-plane popcounts combine with the stored factors into the estimate and its bound. +The estimate builds the shortlist; exact base-table distances produce the final ranking (§6). + +### 3.3 Why this fits Hudi + +- Centroids are small enough to load at planning time: `K × D` floats (~12 MB for K=4096, + D=768). +- Codes are compact: a 1B × 128-dim float table's raw vectors are ~512 GB; the RaBitQ index + including keys and locators is ~136 GB, and the scanned portion per query is tens of MB. +- Posting keys are prefix-scannable by generation, cluster, and shard (§4). +- Quantizer state is stable: a seed and centroids, no per-generation learned codebook. +- Exact re-ranking preserves correctness for returned candidates. + +--- + +## 4. Metadata Table Storage Model: the Posting Block + +This section describes the core storage contribution of this RFC. + +### 4.1 The posting block + +A naive index would write one MDT record per indexed vector. At a billion rows that is a +billion MDT records per generation — prohibitive to write, compact, scan, and clean. + +Instead, the index sorts entries by `(cluster, fileGroup, rowPosition)` and packs +~1–4K vectors' worth of codes and metadata into a single MDT record — a **posting block** — +laid out column-wise (structure-of-arrays) so each scan pass touches only the columns it +needs: + +```text +POSTING BLOCK (~512 KB target, one MDT record) +┌───────────────────────────────────────────────┐ +│ S1 sign planes ← pass 1 touches this │ +│ S2 extra bit planes ← pass 2, survivors │ +│ S3 factor arrays ← both passes │ +│ S4 row locators ← finalists only │ +│ S5 dictionaries (file groups, partitions) │ +│ S6 record keys ← finalists only │ +└───────────────────────────────────────────────┘ +key: 0x10 | generation | clusterId | shardId | blockId +``` + +Consequences: + +- **~1000× fewer MDT records.** One block record replaces ~1–4K per-vector records, cutting + write amplification, compaction cost, and cleaner load by three orders of magnitude. +- **Contiguous cluster scans.** The binary key scheme makes "scan cluster 12345" a single + contiguous HFile range read rather than thousands of point lookups. +- **Pay only for promise.** Column-wise layout means pass 1 reads only sign planes + factors; + only survivors touch extra planes; only finalists touch locators and keys (§6.2). + +### 4.2 Row families + +The `vector_index_<name>` partition holds several record families under one binary-sorted key +scheme, so one prefix scan of a cluster returns its blocks and any fresh deltas together: + +| Key family | Cardinality | Purpose | +|---|---:|---| +| `__manifest__` | 1 | Active generation pointer used by readers. | +| `__centroids__` | 1 | Serialized `K × D` centroid matrix. | +| `__quantizer__` | 1 | RaBitQ type, bit width `B`, rotation seed, normalization flag. | +| `M\|<generation>` | generations | Immutable generation-level metadata. | +| `C\|<generation>\|<cluster>` | K per generation | Cluster manifest: shard count, vector count, candidate file groups, delta/tombstone counters. | +| `P\|<generation>\|<cluster>\|<shard>\|<blockId>` | blocks per generation | **Posting block** (packed codes, factors, locators, keys). | +| `P\|...\|<DELTA>` | deltas | Small per-record delta records appended to a cluster range between compactions. | + +### 4.3 Posting shards + +Large clusters are split into posting shards so one hot cluster does not become one oversized +prefix range. The cluster manifest stores `shardCount`; writers compute +`shardId = hash(record_key) % shardCount`. Within a shard, entries are packed into blocks by +`blockId`. Sharding changes physical layout only, never vector semantics. + +### 4.4 Delta records + +Between compactions, per-commit vector writes append small **delta records** (`blockId` +marked `DELTA`) at the end of the same cluster key range. Because they share the prefix, a +single cluster prefix scan sees packed blocks and fresh deltas in one pass (§6, §7). + +### 4.5 Generation model + +A generation is a consistent set of centroid, quantizer, cluster, and posting-block metadata: + +```text +__manifest__ -> active generation id +M|<gen> -> immutable generation metadata +__centroids__ -> centroid matrix for the active generation +__quantizer__ -> RaBitQ seed / width / normalization +C|<gen>|... -> cluster manifests +P|<gen>|... -> posting blocks + deltas +``` + +A query uses one generation consistently. Publishing a generation is atomic from a reader's +perspective: write the generation's rows, commit them to MDT, then flip `__manifest__`. Old +generations are retained until no retained table snapshot needs them. This enables +blue/green rebuilds with clean rollback. + +--- + +## 5. Bootstrap and Write Path + + + +### 5.1 Spark bootstrap + +Bootstrap builds a complete generation from a table snapshot in two distributed phases: Review Comment: Addressed in §4.5, §5.1, and §10 in commit `b902152fa6`. Generation construction is idempotent at generation scope. The builder allocates a generation id once and writes deterministic `M|`, `T|`, `C|`, and `P|` keys, so a retry overwrites/reuses the same logical records rather than exposing a second partial generation. A `BUILDING` generation is invisible to readers. Before activation, bootstrap validates centroid chunk counts/checksums, vector schema, cluster/block counts, and memory budgets. Only one completed MDT commit changes the generation to `ACTIVE` and flips `__manifest__`; a failed or invalid build never becomes query-visible. Abandoned `BUILDING` generations can be garbage-collected safely. The revised test plan calls out partial-write retry, validation-before-activation, and abandoned-generation cleanup explicitly. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
