In Apache Hudi, how does an update reach storage differently under Copy-on-Write and Merge-on-Read?
answer
- both start by finding the same file group
- one pays in bytes rewritten, one in merge work
- the timeline action name gives it away
- look for hidden log files in the partition
- deltacommit versus commit
basics
~20 sCopy-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.
solid answer
~40 sBoth types index the record key to the same **file group**, then diverge. Copy-on-Write reads that group's current Parquet base file, merges the incoming records in, writes a whole new base file, and records a `commit` instant on the `.hoodie` timeline — readers see finished Parquet and do no merging. Merge-on-Read instead appends the records as blocks to a **log file** attached to that file group's current base file and records a `deltacommit`. Nothing is rewritten, so the write is cheap and fast, but a snapshot reader must now merge base file plus logs on the fly, and a background **compaction** must eventually fold the logs into a new base file. So the same logical update costs bytes-rewritten on the write side under CoW, and merge work on the read side under MoR.
code
text · 8 lines-- COPY_ON_WRITE partition after two commits
country=US/a1b2c3d4-0_0-24-25_20240115103000123.parquet
country=US/a1b2c3d4-0_0-31-44_20240115110000789.parquet
-- MERGE_ON_READ partition after one load plus two streaming deltacommits
country=US/a1b2c3d4-0_0-24-25_20240115103000123.parquet
country=US/.a1b2c3d4-0_20240115103000123.log.1_0-31-44
country=US/.a1b2c3d4-0_20240115103000123.log.2_0-38-55go deeper
Know that both types locate the record the same way, then one rewrites a Parquet file while the other appends to a log file. Recognize log files in a listing as the Merge-on-Read signature.
Explain the full write path for each type, including the precombine merge, the atomic commit, and the distinct timeline action names. State clearly where the cost lands in each case.
Be ready to reason about the operational consequences: unbounded log growth without compaction, memory pressure in the read-time merge, and how delete semantics differ before compaction. Diagnose from the timeline and a directory listing.
Own the pairing of ingestion pattern to table type across many tables, including that the wrong pairing shows up as either hour-long upserts or steadily degrading queries. Budget compaction as a first-class, funded workload.
## The same first step Whatever the table type, an Apache Hudi upsert starts identically: the incoming records are tagged. Hudi consults its index with each record key to find which **file group** — a stable, file-ID-identified set of records inside a partition — already holds that key. Records that match an existing file group are updates; the rest are inserts. Only after tagging does the table type decide what physically gets written. ## Copy-on-Write: rewrite the base file In a `COPY_ON_WRITE` table each file group version is a single Parquet **base file**. For every file group that received at least one update, the writer streams the existing base file, merges the incoming records into that stream (the precombine field breaks ties between two records with the same key), and writes a brand-new base file for the group. File groups no record landed in are not touched at all. On success Hudi appends a completed `commit` instant to the timeline in `.hoodie`; before that instant exists, the new files are invisible, which is what makes the write atomic. The read side is then trivial: pick the newest base file per file group and read Parquet. The cost sits entirely on the writer, and it is proportional to the *bytes* in the touched files rather than to the number of changed rows. ## Merge-on-Read: append a log block In a `MERGE_ON_READ` table a file group version — a **file slice** — is a base file *plus zero or more log files*. An update is not merged into the base file. Instead the writer appends the records as blocks into a log file associated with that file group and that base instant; log blocks are Avro-encoded by default, and a delete arrives as a delete block rather than a rewritten row. The commit is recorded as a `deltacommit` instant, a distinct action name from CoW's `commit`, so you can tell from the timeline alone which path a write took. No large Parquet file is read or rewritten, so ingestion latency drops sharply and write amplification largely disappears. Two costs replace it: - **Read cost.** A snapshot query over the table must now open the base file *and* replay its log blocks, merging by record key in memory, before it can return rows. The more log files accumulate on a slice, the slower and more memory-hungry that merge becomes. - **Compaction.** A background job must periodically read a file slice's base file plus its logs and write a new base file, producing a fresh file slice with no logs. Compaction is scheduled as a `compaction` instant and lands as a `commit` when it completes. Until it runs, the read cost keeps growing. ## Reading the difference off the directory listing A Copy-on-Write partition contains only `.parquet` base files. A Merge-on-Read partition contains base files plus hidden log files whose names carry the same file ID as the base file they belong to. If you see log files, the table is Merge-on-Read and there is unmerged data on disk; if the timeline shows `deltacommit` instants, the same conclusion follows. ## Delete handling Under Copy-on-Write a delete is just another merge input: the row is dropped while the base file is rewritten, so the file that comes out no longer contains it. Under Merge-on-Read a delete is a delete block appended to a log file — the row is still physically present in the base file, and only the merge at read time (or the next compaction) removes it from results. This is why a Merge-on-Read table can look like it "still returns deleted rows" to a reader using the wrong query type. ## Do not cross the wires Two naming traps are worth calling out. First, Hudi's Merge-on-Read is a *table type* covering every write to the table — it is not the same thing as a per-operation write mode chosen for a single DML statement in other table formats. Second, Hudi's log files are its own append log format carrying record-level blocks; they are not equivalent to a positional delete file or a bitmap of deleted row positions. Hudi merges by *record key*, which is why a record key and a precombine field are mandatory concepts in Hudi and optional or absent elsewhere. ## Choosing between them Copy-on-Write suits batch-loaded tables read far more than written: you pay once at load and every query is fast. Merge-on-Read suits streaming and CDC ingestion where commits arrive every few minutes and updates are scattered: you pay a little on every read and you must operate compaction as a real, monitored service. The wrong pairing is what produces the two classic incidents — hour-long batch upserts on a CoW table with scattered keys, and steadily degrading queries on a MoR table whose compaction stopped.
- Looking only at the Hudi timeline, how can you tell which table type wrote a commit?Ingestion into a Copy-on-Write table lands as a `commit` instant; ingestion into a Merge-on-Read table lands as a `deltacommit`. A Merge-on-Read table also shows `compaction` instants that complete as `commit` instants. So a timeline containing `deltacommit` actions is unambiguously Merge-on-Read.
- In a Merge-on-Read table, what happens to a deleted record before compaction runs?The row stays physically present in the base file, and a delete block is appended to the file group's log. A snapshot query replays the log and correctly omits the row; a read-optimized query, which reads base files only, still returns it. Compaction later writes a new base file with the row genuinely gone.
- Does Merge-on-Read remove write amplification entirely?No, it defers and amortizes it. Ingestion writes are cheap because nothing is rewritten, but compaction later reads the base file plus all its log blocks and rewrites a new base file. The total bytes rewritten can be lower than Copy-on-Write because many updates are folded in one pass, but the work still happens — on a schedule you control.
saying these in an interview costs you the question
- Thinks Merge-on-Read never rewrites base files at all
- Says Copy-on-Write appends deltas that readers merge later
- Calls Hudi's log files the same thing as positional delete files
- Believes Merge-on-Read is a per-statement write mode, not a table type
- Assumes both types record the same instant action on the timeline