Skip to content
Great engineers know 5 min read · Hudi

Apache Hudi: keys, indexes and table services

Apache Hudi was built at Uber to apply a firehose of upserts to a data lake within minutes. Its design follows from that: every record has a key, an index finds which file holds it, and a timeline records every action on the table.

Read first: Open table formats · Copy-on-write vs merge-on-read (skip if you know them)

Keys, ordering and file groups

ConceptWhat it is
Record keyUniquely identifies a record (within a partition, or globally). Upserts and deletes are by key.
Ordering (precombine) fieldDecides which version wins when several records share a key, for example the source update timestamp. A late, older record does not overwrite a newer one.
Partition pathOptional folder layout, like Hive partitions.
File groupAll versions of a set of records live in one file group: a base ParquetParquet: A columnar file format: values of each column are stored together, so queries read only the columns they need. Learn more → file plus, for merge-on-read, log files of later changes. Each record key maps to exactly one file group.

This is the main difference from Delta and Iceberg: Hudi tracks records by key, so it knows where a record lives before it writes. That makes it strong at high-volume upserts and incremental reads.

Copy-on-write and merge-on-read

Copy-on-write (COW)

An update rewrites the base file of the affected file group. Reads are plain Parquet scans. Good for read-heavy tables with moderate update rates.

Merge-on-read (MOR)

Updates are appended to row-based log files next to the base file, which is fast to write. Readers merge base and logs, until compaction folds the logs into a new base file. Good for write-heavy and streaming upserts.
Query typeSeesUse
Snapshotsnapshot: A complete, consistent version of a table at one moment. Readers always see one snapshot, never a half-finished write. Learn more →The latest state (base files merged with logs on MOR)Default for most queries
Read-optimizedBase files only, so MOR changes not yet compacted are missingFast reads where slight staleness is acceptable
IncrementalRecords changed since a given commit timeDownstream pipelines that only process changes

Writing from Spark

PySpark · An upsert into a merge-on-read table
hudi_options = {
    "hoodie.table.name": "customers",
    "hoodie.datasource.write.table.type": "MERGE_ON_READ",
    "hoodie.datasource.write.recordkey.field": "customer_id",
    "hoodie.datasource.write.precombine.field": "updated_at",   # the ordering field
    "hoodie.datasource.write.partitionpath.field": "country",
    "hoodie.datasource.write.operation": "upsert",
}
changes.write.format("hudi").options(**hudi_options).mode("append").save("s3://lake/customers")

Other write operations: insert (no key lookup, may create duplicates), bulk_insert (fast initial loads), delete, and insert_overwrite for partitions. The upsert path is where the index matters.

The index

To upsert, Hudi must find which file group already holds each incoming key. Scanning every file would be slow, so Hudi keeps an index:

IndexHow it finds a keyFits
SimpleJoins incoming keys with keys read from the affected partitionsThe default on Spark; fine for moderate tables and update rates
BloomBloom filters (and optionally key ranges) built from record keys, to skip files that cannot hold a keyKeys with some ordering, such as time-prefixed ids, so ranges prune
BucketHashes the key to a fixed bucket, so no lookup is neededVery high write throughput; bucket count fixed up front
Record-levelAn exact key-to-file-group map in the metadata table (global since 0.14, also partitioned since 1.1)Large tables with random updates

The timeline and table services

Every action is an instant on the timeline in the .hoodie folder, with a time, a type (commit, deltacommit, compaction, clean, rollback, replacecommit for clustering) and a state (requested, inflight, completed). Readers only see completed instants, which gives atomic commits; a failed write is rolled back.

ServiceWhat it does
Compaction (MOR)Merges log files into new base files, so snapshot reads stay fast
ClusteringRewrites data to improve layout: small files into large ones, sorting by query columns
CleaningDeletes file versions older than the retention, which bounds time travel like VACUUM
ArchivalMoves old instants out of the active timeline

Services can run inline (in the writer, adding latency), asynchronously alongside a streaming writer, or as separate jobs. Skipping compaction on a MOR table is the most common cause of slow reads.

Note: Hudi 1.0 (December 2024) added secondary indexes, partial updates, a scalable LSM-based timeline and non-blocking concurrency control, first for concurrent Flink streaming writers. Hudi also reads and writes metadata for other formats through Apache XTable, and remains common in CDC-heavy platforms, notably on AWS.

Common mistakes

  • No ordering field, or the wrong one — A late, older record can overwrite a newer one.
  • Using insert instead of upsert for CDC — insert skips the index lookup and can create duplicate keys.
  • Never compacting a MOR table — Every read merges an ever-growing pile of log files.
  • Choosing a bloom index for random UUID keys — Key ranges do not prune; most files are checked. Consider record-level or bucket indexes.

Interview prep

THE QUESTION

"When would you choose Hudi over Delta or Iceberg?"

When the workload is dominated by high-volume, key-based upserts and deletes, often streaming CDC, and by incremental reads of what changed. Hudi tracks records by key with an index that maps keys to file groups, supports merge-on-read with log files for fast writes, and offers incremental queries natively. Mention the trade-offs: more configuration (keys, ordering field, index type, table services), compaction to manage on MOR tables, and narrower engine support than Iceberg in some ecosystems.

Avoid saying: "Hudi is the streaming one" with nothing about keys, indexes and merge-on-read. Those are why it is good at upserts.

What the interviewer asks next. Answer out loud first, then open the strong answer.

Follow-up"Read-optimized vs snapshot queries on a MOR table?"
Snapshot merges base files with log files and returns the latest data. Read-optimized reads only base files: faster, but it misses changes not yet compacted. Fine for some analytics, wrong for anything that needs the latest state.
Scenario"Upserts into a large Hudi table got slow as the table grew. What do you look at?"
The index first: with random keys a bloom index cannot prune, and simple index joins grow with the table; a record-level or bucket index can help. Then file sizing and clustering, compaction lag on MOR tables, and whether the partition path matches how updates arrive.
Trap"Hudi incremental queries need Change Data Feed enabled, like Delta."
Incremental queries are built in: they read records changed after a given commit time from the timeline. Hudi also has a separate CDC mode for before and after images, but plain incremental pulls do not need it.

What you learned

  • Hudi's record keys, ordering field and file groups
  • Copy-on-write vs merge-on-read tables, and the three query types
  • What the index is for, and the main index types
  • The timeline and the table services that keep a Hudi table healthy

Key takeaways

  • Hudi tracks records by key; an index maps keys to file groups.
  • COW rewrites base files on update; MOR appends log files and compacts later.
  • Snapshot, read-optimized and incremental are the three query types.
  • The timeline records every action; compaction, clustering and cleaning keep the table healthy.

Check yourself

3 questions

What does the ordering (precombine) field decide?

Show the answer

Which record wins when several share a key. It prevents an older record from overwriting a newer one.

What does a read-optimized query on a MOR table miss?

Show the answer

Changes still in log files that have not been compacted. Read-optimized reads only base files.

Why does Hudi need an index?

Show the answer

To find which file group already holds each incoming key during upserts. Upserts must locate existing records without scanning every file.

Keep going

Up next · lesson 23 of 30 · 4 min read
Deletion vectors
Mark rows as deleted without rewriting whole files, and the cost readers pay for it.

Related lessons

Previous: Copy-on-write vs merge-on-read

Primary sources: Apache Hudi concepts · Hudi table types · Hudi indexing