skip to content

Apache Hudi

You will learn Apache Hudi, the table format built first for fast upserts and incremental pulls rather than for read-only analytics: Copy-on-Write versus Merge-on-Read storage, a timeline of instants, and record-level indexing. Interviewers bring up Hudi when the scenario is CDC ingestion, where its write path is the strongest of the three formats.

on this pageshow

explore

questions

12

In an Apache Hudi write, what do the record key and precombine field control?

level: juniorimportance: must knowfreq 55%

answer

  1. two configs decide identity and recency
  2. one answers 'which row is this?'
  3. the other answers 'which version wins?'
  4. the index is keyed by the first one
  5. a duplicate key keeps the larger of the second

basics

~20 s

The record key is a row's identity, so Hudi can update or delete that row instead of only appending. The precombine field breaks ties between duplicates of a key, keeping the record with the higher value.

solid answer

~40 s

`hoodie.datasource.write.recordkey.field` names the column, or comma-separated columns, that identify a row. Hudi stores it in the hidden `_hoodie_record_key` column and uses it on every upsert to ask the index which file group already holds that key, so the write becomes an update of that file group rather than a blind append. `hoodie.datasource.write.partitionpath.field` scopes that identity to a partition unless the index is global. `hoodie.datasource.write.precombine.field` orders two versions of the same key: when a key arrives twice in one batch, Hudi keeps the record with the larger precombine value, normally an event timestamp or a CDC sequence number. Under event-time ordering the same comparison also stops an out-of-order older record from overwriting the stored one; under commit-time ordering the newest commit simply wins.

go deeper

for a junior

Be able to name the two write options and say what each is for: one identifies the row so it can be updated, the other picks the winner when the same row appears twice. Know that a null record key fails the write.

for a middle

Explain the tagging step: Hudi looks the key up in the index, routes tagged records to the file group that already holds them, and treats the rest as inserts. Be able to contrast event-time and commit-time ordering of the precombine value.

for a senior

Show judgment about choosing the fields. Argue for a source-produced, monotonic precombine value so replays and backfills are idempotent, and explain why the key definition is effectively immutable once data exists.

for a principal

Own this as a platform contract. Standardise how teams derive record keys and precombine values from CDC sources, since a bad choice is silently unrecoverable and forces a full table rewrite across every downstream consumer.

## The problem these two configs solve A directory of Parquet files has no notion of row identity — you can append files to it, but you cannot say "row 42 changed". Hudi's upsert path exists to give a lake table that ability, and it rests on two write configurations you set the first time you write the table: which columns identify a row, and which column decides which version of that row is newer. ## The record key `hoodie.datasource.write.recordkey.field` names one column, or several comma-separated columns, that uniquely identify a row: a source primary key, a CDC key, a UUID. Hudi materialises the value into a hidden metadata column, `_hoodie_record_key`, which it adds to every base file alongside `_hoodie_commit_time`, `_hoodie_commit_seqno`, `_hoodie_partition_path` and `_hoodie_file_name`. A composite key is stored as a list of `field:value` pairs so it round-trips exactly. The record key is what the index is keyed by. When an `upsert` batch arrives, Hudi *tags* each incoming record: it asks the index which file group already contains that key. Tagged records are updates and are routed to that file group — merged into a new version of the base file on a Copy-on-Write table, appended to a log file on a Merge-on-Read table. Untagged records are inserts and are placed into a file group with room. With no record key there is nothing to tag against, and every write can only append. Two practical rules follow. The key value may not be null — Hudi fails the write rather than guessing. And the key definition is effectively immutable for the life of the table: change which columns make up the key and every stored row now has a different identity, so no future update can ever match it. That is a full rewrite, not a config change. ## The partition path `hoodie.datasource.write.partitionpath.field` decides the directory a record lands in. Uniqueness is enforced per (partition path, record key) unless you configure a global index, which is why a record whose partition value changes can end up existing twice. ## The precombine field `hoodie.datasource.write.precombine.field` names the column that orders two versions of the same key. It is consulted in two places: 1. **Within the incoming batch.** A CDC batch commonly carries three changes to the same row. Hudi collapses them before writing, keeping the record with the largest precombine value, so one commit never writes the same key twice. 2. **Against the record already stored.** Here the behaviour depends on merge semantics. With *event-time ordering* Hudi compares the incoming precombine value against the stored one and keeps the larger, so replaying an old file or a lagging stream partition cannot resurrect a stale version. With *commit-time ordering* the newest commit wins outright regardless of the value. Hudi 1.x exposes this choice as an explicit record merge mode; earlier versions expressed it through the record payload class you configured. Choose a monotonically non-decreasing value produced by the source: the database transaction timestamp, a binlog sequence number, a Kafka offset made comparable across partitions. Ingestion wall-clock time is a poor choice, because a backfill that re-reads history will stamp old rows with a *newer* precombine value and overwrite good data. ## Deletes The same key drives deletion. A hard delete is a write with `hoodie.datasource.write.operation=delete`, which removes the row entirely; a soft delete sends the record with the `_hoodie_is_deleted` marker column set so the row stays present with its non-key fields nulled. Both need the record key to find the file group holding the row. ## Common failure modes - **A key that is not actually unique.** Hudi will happily collapse two genuinely different business rows into one because they share a key, and the loss is silent. - **No precombine field, or a constant one.** Duplicates in a batch then resolve arbitrarily, and re-running the same batch can flip the winner, which makes the pipeline non-deterministic. - **Assuming precombine also orders rows on disk.** It does not — it is a comparison used during merge, not a sort order. Clustering and the bulk-insert sort modes control physical ordering. - **Assuming the newest write always wins.** Under event-time ordering it does not: an incoming record with a smaller precombine value is discarded even though its commit is newer, which is exactly the protection you asked for and exactly the surprise when a backfill "does nothing". In interviews these two configs are the entry point to everything else on the Hudi timeline-and-index topic: the key is what the index maps, and the precombine value is what makes replaying an upstream stream safe.

  • Can a Hudi record key span more than one column?
    Yes — `hoodie.datasource.write.recordkey.field` accepts a comma-separated list, and Hudi uses a complex key generator that stores the composite as a list of field:value pairs in `_hoodie_record_key`. Every component must be non-null. Keep composite keys narrow: the key is materialised into every base file and compared on every upsert, so wide composites cost storage and tagging time.
  • What happens if you change which columns make up the record key after the table already has data?
    Existing rows keep their old `_hoodie_record_key` values, so new records under the new definition never match them. Updates land as inserts, and the table quietly accumulates duplicates of the same business entity. There is no in-place migration: you rewrite the table under the new key definition, typically with a bulk_insert into a fresh table.
  • Why is an event timestamp a better precombine field than ingestion time?
    Precombine decides which version of a key survives. Ingestion time reflects when your pipeline saw the row, so a backfill or replay stamps old data as newest and overwrites current values. A source-produced event timestamp, transaction time or log sequence number stays correct no matter when or how often you reprocess, which makes the pipeline idempotent.

saying these in an interview costs you the question

  • Thinks the record key only deduplicates within a single batch
  • Says any column can be the record key even if it is not unique
  • Assumes you can change the record key definition later without a rewrite
  • Claims the precombine field controls row ordering inside the data file
  • Believes the newest commit always wins regardless of precombine value

context

open as a page

In Apache Hudi, how does an update reach storage differently under Copy-on-Write and Merge-on-Read?

level: middleimportance: must knowfreq 70%

basics

~20 s

Copy-on-Write merges the update into the base Parquet file and rewrites it as a new file version. Merge-on-Read appends the update as a log block in the same file group, leaving the merge to read time or to a later compaction.

open as a page

In Apache Hudi, what is an instant on the .hoodie timeline?

level: middleimportance: must knowfreq 50%

basics

~20 s

An instant is one entry on the timeline in a Hudi table's .hoodie directory: an action such as commit, deltacommit, compaction or clean, stamped with a monotonic time and a state of requested, inflight or completed.

open as a page

A Hudi Merge-on-Read table's snapshot queries slow down every hour — what would you check?

level: seniorimportance: must knowfreq 55%

basics

~20 s

Almost always compaction is not keeping up, so log files pile up on each file slice and every snapshot query merges more of them. Check the timeline for delta commits since the last completed compaction, then fix scheduling, parallelism or the writer's compaction service.

open as a page

Which Apache Hudi table type is the default, and what does it write on an update?

level: juniorimportance: should knowfreq 52%

basics

~20 s

Copy-on-Write is the default Hudi table type. Updating a record rewrites the whole Parquet base file that holds it as a new file version, so readers just scan the newest base file per file group with no merging.

open as a page

In a Hudi Merge-on-Read table, how do snapshot and read-optimized queries differ?

level: middleimportance: should knowfreq 50%

basics

~20 s

A snapshot query merges each file slice's base file with its log files at read time, so it returns the latest committed state. A read-optimized query reads base files only, so it is faster but shows data as of the last compaction.

open as a page

In Apache Hudi, how does the bulk_insert operation differ from upsert?

level: middleimportance: should knowfreq 42%

basics

~20 s

upsert looks every incoming record up in the index and merges it into the file group already holding that key. bulk_insert skips the index lookup and small-file sizing entirely and writes new files fast, so it cannot deduplicate against existing rows.

open as a page

Why does an Apache Hudi upsert create a duplicate when a record's partition value changes?

level: seniorimportance: should knowfreq 32%

basics

~20 s

A non-global Hudi index enforces key uniqueness only within a partition. When the partition value changes, the lookup in the new partition finds nothing, so the record is inserted there while the original row stays in its old partition.

open as a page

An Apache Hudi upsert job's index tagging slows as the table grows — what would you change?

level: seniorimportance: should knowfreq 38%

basics

~20 s

Tagging cost tracks how many candidate files each key must be probed against. Switch from a bloom index to a record-level index for direct key-to-file lookups, or a bucket index for hash-derived placement with no lookup at all, and keep keys range-friendly.

open as a page

For a new Hudi-based lakehouse, how would you decide Copy-on-Write versus Merge-on-Read per table?

level: principalimportance: should knowfreq 38%

basics

~20 s

Decide per table from the write pattern, the freshness target and the read profile: Copy-on-Write for batch-loaded, heavily read tables, Merge-on-Read for frequent scattered updates where you can operate compaction. The choice is effectively irreversible without a rewrite.

open as a page

In Apache Hudi, what is a file group and how does a file slice relate to it?

level: middleimportance: nice to knowfreq 32%

basics

~20 s

A Hudi file group is a set of records inside a partition, identified by a stable file ID and owning those record keys over time. A file slice is one version of that group: a base Parquet file plus, in Merge-on-Read, its log files.

open as a page

Why can an Apache Hudi incremental query miss rows committed by a slow concurrent writer?

level: seniorimportance: nice to knowfreq 28%

basics

~20 s

An incremental read positions on instant time, the moment a write started. A writer that started earlier but finished later lands a commit behind the reader's checkpoint, which has already moved past it, so those rows are never returned.

open as a page