vinothchandar commented on code in PR #19309:
URL: https://github.com/apache/hudi/pull/19309#discussion_r3920105004


##########
rfc/rfc-109/rfc-109.md:
##########
@@ -0,0 +1,1005 @@
+<!--
+  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)
+- [Limitations](#limitations)
+- [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, Compatibility, and 
Freshness](#9-correctness-compatibility-and-freshness)
+- [10. Test Plan](#10-test-plan)
+- [11. Rollout and MVP Scope](#11-rollout-and-mvp-scope)
+- [12. References](#12-references)
+- [Appendix A. Controlled BIGANN Research 
Evidence](#appendix-a-controlled-bigann-research-evidence)
+
+---
+
+## Abstract
+
+This RFC proposes native approximate nearest-neighbor (ANN) vector search in 
Apache Hudi.
+Tables increasingly carry embedding columns represented by Hudi's 
fixed-dimension
+`VECTOR(D[, elementType])` logical type, 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. A **posting** is the index entry 
for one logical
+record: its generation-scoped cluster membership, compressed scoring material, 
record key,
+and source locator. A **generation** is an atomically published, versioned 
namespace for one
+complete index state: schema contract, centroids, quantizer, routing, 
postings, and freshness
+frontier. Its model identity is immutable after activation; postings and 
frontier advance only
+through snapshot-atomic MDT commits within that namespace.
+
+> **The posting block.** Instead of one MDT record per posting, 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.
+
+Artifacted measurements on a **1-billion-row, 128-dimensional table** show 
approximate-only
+recall@10 of 0.822, 0.857, and 0.871 at `nprobe=16`, `32`, and `64`; exact 
reranking reaches
+0.960 at `nprobe=32` with `refineFactor=50` (Appendix A). These modes serve 
different goals:
+approximate-only provides the latency floor, while exact reranking is the 
recommended
+quality mode.
+
+---
+
+## 1. Goals and Non-Goals
+
+### 1.1 Goals
+
+1. Keep authoritative vector values in a top-level base-table `VECTOR(D[, 
elementType])` 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: prune to requested table 
partitions, probe a few
+   clusters, scan contiguous key ranges, and score 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.
+- Native non-Spark generation construction, GPU encoding, and 
workload-specific auto-tuning.
+
+### 1.3 Alternatives considered
+
+- **One index record per vector** is simple but creates billions of MDT 
records and excessive
+  write, compaction, and scan amplification; posting blocks preserve MDT 
ownership while
+  amortizing that overhead.
+- **Dedicated index files in the table** permit specialized layouts but 
introduce a second
+  commit, cleaning, and snapshot protocol. MDT records reuse Hudi's existing 
transaction and
+  table-service machinery.
+- **External sidecar indexes or vector databases** may offer richer serving 
features but lose
+  atomic Hudi snapshot semantics and require another storage system. They 
remain valid when
+  independent serving infrastructure is desired.
+- **Native ANN libraries** improve local kernels but do not define durable 
object-store layout,
+  multi-writer maintenance, or snapshot visibility. They may be used behind 
the interfaces in
+  this RFC without changing the persisted contract.
+
+---
+
+## Limitations
+
+The first implementation deliberately fails closed or falls back when it 
cannot prove the
+following contracts:
+
+- **Vector schema evolution.** Changing the indexed vector's dimension or 
element type makes
+  the active generation incompatible. Existing-index maintenance and indexed 
queries fail
+  closed after such a change; rebuilding or recreating the index establishes 
and atomically
+  activates a compatible generation. This RFC does not claim transparent 
online migration
+  across incompatible vector schemas.
+- **Filtered ANN.** The initial implementation is intended to support 
partition pruning: posting
+  dictionaries preserve partition paths so the query path can exclude postings 
outside the
+  requested table partitions before candidate-heap selection and exact 
reranking. This produces
+  partition-local top-K rather than filtering a global top-K after search. 
Arbitrary
+  predicate-plus-kNN planning over non-partition columns is not part of the 
initial implementation.
+- **Index quality.** IVF quality depends on the training sample and corpus 
geometry. A fixed
+  seed makes fitting repeatable for fixed inputs and partitioning, but does 
not make centroids
+  invariant to changed file/RDD partition layouts. Recall must therefore be 
measured for each
+  built generation.
+- **Local incremental rebalancing.** LIRE split/merge is design-specified but 
not part of the
+  first implementation. Until it lands, sustained skew or centroid drift is 
corrected by a
+  full generation rebuild and atomic activation.
+- **Maintenance scale evidence.** Bootstrap and query paths are measured at 
one billion rows.
+  Incremental COW/MoR mutation, compaction/clustering relocation, and 
multi-writer races require
+  the test evidence in §10 before their implementations are declared complete.
+- **MoR exact-fetch cost.** Log-resident finalists are resolved by key through 
the pinned
+  merged file slice. This preserves correctness but can cost more than 
base-file positional
+  fetch and requires explicit scale measurement.
+
+---
+
+## 2. Architecture
+
+![RFC-109 architecture overview](diagrams/01-architecture-overview.svg)
+
+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
+```
+
+### 2.1 Design principles
+
+- **Correct before optimized.** Delta postings and RLI arbitration make every 
completed update
+  immediately discoverable and prevent stale versions from being served before 
any table
+  service runs. Exact reranking resolves MoR log-resident finalists through 
the pinned merged
+  file slice; compaction may reduce that additional read cost but is never a 
prerequisite for
+  candidate visibility, exact materialization, or stale-version suppression.
+- **Side-effect-free queries.** The query path never writes the MDT or data 
table. Workload
+  signals such as log-resident finalist count, stale-candidate crowding, and 
fetch-key mismatch
+  travel only through metrics. Maintenance may consume an externally 
aggregated, decayed
+  snapshot of those signals and must remain correct when it is absent.
+- **One pinned view.** MDT, RLI, file-slice, base-table, routing, and cache 
reads for a query use
+  one pinned snapshot and one ACTIVE generation.
+- **Derived state is replaceable.** The data table remains authoritative. 
Index artifacts may
+  be compacted, locally rebalanced, rebuilt, or retired without changing 
logical table data.
+
+Each vector index is one MDT partition. RFC-109 adds no table properties, 
common writer-
+dispatch changes, timeline-semantics changes, or behavior visible to 
non-vector readers or
+writers. Every consistency mechanism is either a record in the vector-index 
MDT partition or
+logic in the RFC-owned vector indexer and query planner. Creating an index:
+
+```sql
+CREATE INDEX embedding_idx
+ON products
+USING VECTOR (embedding)
+OPTIONS (
+  'vector.metric'      = 'cosine',
+  'vector.quantizer'   = 'IVF_RABITQ',
+  'vector.num_clusters'= '4096'
+);
+```
+
+creates:
+
+```text
+.hoodie/metadata/vector_index_embedding_idx/
+```
+
+and does not add generated columns to the base-table schema. RFC-109 consumes 
rather than
+redefines the RFC-99 vector contract: the source must be a top-level Hudi
+`VECTOR(D[, elementType])`. The table schema is authoritative for `D` and 
element type;
+index definitions and generation manifests repeat them only for integrity 
validation.
+
+The current RFC-99 storage backing is fixed-width bytes: Avro `FIXED` and 
Parquet
+`FIXED_LEN_BYTE_ARRAY(D × elementWidth)`. Engine adapters expose idiomatic 
values—Spark uses
+an annotated `ArrayType(FloatType|DoubleType|ByteType)`—and convert at the 
storage boundary.
+This fixed width also avoids Parquet LIST repetition-level traversal during 
positional exact
+fetch. A plain `ARRAY<FLOAT>` is not implicitly indexable and requires 
explicit migration or
+backfill to `VECTOR(D)`; index creation neither reinterprets nor rewrites it.
+
+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 and persisted-format version. 
|
+| `M\|<generation>` | generations | Generation state, vector schema, bootstrap 
baseline, verified-contiguous frontier, chunk counts/checksums, and quantizer 
metadata. |
+| `T\|<generation>\|<chunk>` | chunks per generation | Size-bounded chunks of 
the serialized `K × D` centroid matrix. |
+| `C\|<generation>\|<cluster>` | K per generation | Cluster manifest: routing 
version, shard count, vector count, candidate file groups, and 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 between 
compactions. |
+| `F\|<generation>\|<dataInstant>` | data writes after baseline | Atomic proof 
that the generation incorporated one source data-write instant. |
+
+### 4.3 Posting shards
+
+A **posting shard** is only an intra-cluster MDT key-range split; it is not an 
RLI shard, a
+data-table bucket or file group, or a Spark partition. Large clusters are 
split into these
+posting shards so one hot cluster does not become one oversized prefix range. 
The cluster
+manifest stores `shardCount` and `routingVersion`; writers compute
+`shardId = hash(record_key) % shardCount`. Changing `shardCount` changes every 
key's mapping,

Review Comment:
   Do we define how we split clusters uniformly across MDT vector index shards? 
I prefer using the term `posting segment` to not overload MDT shards. Range 
partitioning a key space is pretty tricky and new vectors need to be assigned 
to a posting segment/RLI shard while preserving balance? do we cover those 
here? I think RFC is not covering the mapping function between posting 
segment/shard and MDT vector search index shard



-- 
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]

Reply via email to