For a new Hudi-based lakehouse, how would you decide Copy-on-Write versus Merge-on-Read per table?
answer
- it is a design decision, not a tuning knob
- start from a conservative default
- three conditions must hold together, not one
- ask who will be paged for compaction
- check what your query engines can actually read
basics
~20 sDecide 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.
solid answer
~40 sTreat it as a per-table decision with a platform-wide default of `COPY_ON_WRITE`, because most analytics tables are written in batches and read constantly, and CoW gives plain-Parquet read latency with no background service to run. Move a table to `MERGE_ON_READ` when three things hold: commits are frequent and updates are scattered enough that rewriting base files dominates write time; consumers genuinely need freshness tighter than the batch cadence; and your team will operate compaction as a monitored service. Also check engine support — some connectors expose only the read-optimized view of a Merge-on-Read table, which silently changes what your consumers see. Because switching types means rewriting the table, decide before the first load, document the rationale, and publish the freshness contract each type implies.
code
properties · 8 lines# batch-loaded, read-heavy fact table: platform default
hoodie.datasource.write.table.type=COPY_ON_WRITE
hoodie.datasource.write.operation=upsert
# streaming CDC table: exception, with compaction owned and monitored
hoodie.datasource.write.table.type=MERGE_ON_READ
hoodie.compact.inline=false
hoodie.compact.schedule.inline=truego deeper
Know that the two types trade write cost against read cost and that the choice is made when the table is created. You are not expected to set platform policy yet.
Be able to argue a specific table's case: how often it is written, how scattered the updates are, how fresh consumers need it, and what compaction would have to run.
Demonstrate that you would measure write amplification rather than guess, verify engine support for the snapshot view, and refuse Merge-on-Read where nobody will operate compaction.
Own the default and the exception process across many tables, including the cost model on both sides, the irreversibility of the choice, and the freshness contracts each type lets you publish to consumers.
## Why this is a design decision, not a tuning knob An Apache Hudi table is created as either `COPY_ON_WRITE` or `MERGE_ON_READ`, and the type is recorded in the table's properties. It is not a setting you flip on a loaded table — changing it means rewriting the data into a new table and repointing the catalog. So it belongs in the design review before the first load, with the same seriousness as choosing a partitioning scheme. On a platform with hundreds of tables, the leverage is in the *default* and in a small set of explicit exceptions, not in agonizing over each table. ## Set the default to Copy-on-Write Most analytics tables are written on a batch cadence and read far more often than they are written. For those, Copy-on-Write is the correct default and should require no justification: readers get self-contained Parquet with no read-time merge, every engine supports it identically, and there is no background compaction service whose failure degrades queries. Fewer moving parts is a real platform property, not a cop-out. ## The three conditions for Merge-on-Read Require an exception to be argued on three axes together, because any one alone is a bad reason. **Write pattern.** Merge-on-Read pays off when rewriting base files dominates write time — frequent commits, and updates scattered across many file groups rather than concentrated in a few recent partitions. Measure it: bytes rewritten per commit divided by bytes actually changed is the write-amplification ratio, and a large ratio is the quantitative case for MoR. Conversely, if updates land almost entirely in the newest partition, Copy-on-Write rewrites little and MoR buys much less than it appears to. **Freshness requirement.** MoR only helps if someone actually needs data fresher than the batch cadence and will read the snapshot view to get it. If every consumer reads on a schedule and tolerates hourly data, MoR's cheap writes solve a problem no one has, and you have bought a compaction service for nothing. **Operational capacity.** Merge-on-Read is only correct if compaction is run and monitored. Inline compaction is self-healing but makes some commits much slower; asynchronous compaction keeps ingest latency flat but becomes an independent job that can die while ingestion keeps succeeding. Either way you need indicators — age of the newest completed compaction, outstanding log volume per file group — and someone who is paged. A team that will not fund that should not run MoR tables. ## Check the read fleet before standardizing A Merge-on-Read table presents two views: **snapshot**, which merges base files with logs and is fully fresh, and **read-optimized**, which reads base files only and is as of the last compaction — including still showing rows deleted since. Engine support for the snapshot view varies by connector, and if a major consumer can only read the read-optimized view, that consumer's freshness is bounded by your compaction interval whether you intended it or not. Inventory which engines your consumers use before making MoR a platform pattern, and publish the freshness contract per view. ## Cost, not just latency Both types cost money, in different places. Copy-on-Write spends it on write compute and rewritten bytes, and on the storage of superseded file slices until the cleaner removes them. Merge-on-Read spends it on compaction compute and on every snapshot query's merge, plus higher memory for readers. On object storage, request counts matter too: MoR's extra log files raise listing and open costs per query. Model both against your actual ingest and query volumes rather than assuming MoR is universally cheaper — it moves cost, it does not delete it. ## Reversibility and blast radius Because the type is baked in, treat a wrong choice as a migration project: create the new table, backfill, dual-write or cut over, repoint the catalog, retire the old one. That cost is the reason to be conservative. It is also the reason to keep the exception list small and documented — a platform where any team may choose MoR ad hoc ends up with compaction backlogs on tables nobody owns. ## A workable policy Default every table to Copy-on-Write. Allow Merge-on-Read for tables that pass all three conditions, and attach three obligations to the exception: a named owner for compaction, a documented compaction cadence that becomes the read-optimized freshness SLA, and monitoring on compaction age and log backlog. Review the exception list periodically — ingestion patterns drift, and a table that needed MoR when it was fed by a streaming CDC pipeline may not need it after the pipeline moved to hourly batches. Finally, remember that this decision interacts with file sizing and clustering: whichever type you choose, badly sized files will hurt, and neither type is a substitute for a sane layout.
- What single measurement most strongly argues for Merge-on-Read on a given table?Write amplification: bytes rewritten per commit divided by bytes actually changed. A table where a few megabytes of updates rewrite tens of gigabytes because the changes are scattered across many file groups is the clear case. If updates concentrate in the newest partitions, Copy-on-Write rewrites little and the argument collapses.
- If a team asks to switch an existing large Copy-on-Write table to Merge-on-Read, what does that project involve?A rewrite, not a setting change. You create a new Merge-on-Read table, backfill history, run both in parallel long enough to compare results, cut consumers over, repoint the catalog, then retire the old table. Because of that cost, the honest first question is whether better file sizing or a different commit cadence solves the problem instead.
- How do you keep Merge-on-Read from spreading across the platform by default?Make it an exception with obligations attached: a named compaction owner, a documented compaction cadence that becomes the published read-optimized freshness SLA, and alerts on compaction age and log backlog. Review the exception list periodically, since ingestion patterns drift and the original justification may have expired.
saying these in an interview costs you the question
- Recommends Merge-on-Read platform-wide because writes are faster
- Treats the table type as switchable later with a config change
- Ignores whether the query engines support the snapshot view
- Assumes Merge-on-Read is cheaper overall rather than cost-shifted
- Chooses on write latency alone without a freshness requirement