A Hudi Merge-on-Read table's snapshot queries slow down every hour — what would you check?
answer
- the write got cheaper, the read got the bill
- nothing folded the appended data back in
- count what has landed since the last one finished
- a scheduled instant that never completed is the clue
- inline versus a separate job that may be dead
basics
~20 sAlmost 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.
solid answer
~40 sIn a `MERGE_ON_READ` table, ingestion appends log blocks and records `deltacommit` instants; only **compaction** folds those logs into a new base file. If compaction is not scheduled, is failing, or is slower than ingestion, log files accumulate per file slice and every snapshot query pays a larger read-time merge — which looks exactly like "queries get slower every hour". Diagnose from `.hoodie`: count `deltacommit` instants since the newest completed compaction `commit`, and look for `compaction.requested` or `.inflight` instants that never finished. Then check whether compaction is inline or run by a separate async job, whether that job is alive, whether its parallelism and resources match ingest volume, and whether concurrent-writer configuration is blocking it. Immediate relief: point heavy readers at the read-optimized view while you drain the backlog.
code
text · 7 lines.hoodie/20240115060000000.commit <- last completed compaction, 06:00
.hoodie/20240115060500111.deltacommit
.hoodie/20240115061000222.deltacommit
... <- 140 more deltacommits
.hoodie/20240115173000999.deltacommit
.hoodie/20240115120000000.compaction.requested <- scheduled at 12:00, never completed
.hoodie/20240115121500000.rollbackgo deeper
Know that Merge-on-Read defers merge work to read time and that compaction is what clears it. Recognize that steadily worsening query time on such a table points at compaction, not at the query.
Explain the timeline evidence — delta commits since the last completed compaction, a requested instant that never finished — and the difference between inline and async compaction as failure modes.
Diagnose end to end: inspect log depth per file group, find stuck or rolled-back compaction instants, judge whether it is scheduling, resources or skew, and apply the read-optimized view as a temporary lever while communicating its staleness.
Own compaction as a funded, monitored service with explicit indicators: age of the newest completed compaction and outstanding log per file group. Decide when a table's read profile means it should never have been Merge-on-Read.
## The mechanism behind the symptom A Merge-on-Read table stores each file group as a **file slice**: a Parquet base file plus log files appended since it was written. Ingestion is cheap precisely because it never rewrites the base file — it appends log blocks and lands a `deltacommit` instant on the `.hoodie` timeline. That deferred work does not disappear. A **snapshot** query must open the base file and replay every log block on top, merging by record key, before it returns a row. Read cost therefore grows monotonically with the amount of unmerged log data on each slice. Compaction is the only thing that resets it: it reads a slice's base file plus its logs and writes a new base file with no logs. It is scheduled as a `compaction` instant (requested, then inflight) and lands as a completed `commit`. If compaction stops or falls behind ingestion, queries degrade steadily and predictably — hour by hour, exactly as described. ## Diagnose from the timeline first List `.hoodie` and answer three questions before touching any config: 1. **When did compaction last complete?** Find the newest compaction instant that reached completion. If it is hours or days old on a table ingesting every few minutes, you have your answer. 2. **How many `deltacommit` instants have landed since?** That number, times the average per-commit log volume, is the backlog a snapshot query is now merging. 3. **Is there a stuck instant?** A `compaction.requested` with no matching completion, or an instant sitting inflight, means compaction was scheduled and then failed, was killed, or is blocked. Look for `rollback` instants around it — repeated rollbacks mean compaction jobs keep dying, often on memory. A directory listing confirms it: pick a hot partition, group files by file ID, and count log files per base file. A slice with a dozen logs is a slice whose reads are a dozen merges deep. ## Then check how compaction is supposed to run Compaction runs one of two ways, and the failure modes differ. **Inline.** The writer performs compaction itself as part of a write, triggered after a configured number of delta commits (`hoodie.compact.inline` enables it; `hoodie.compact.inline.max.delta.commits` sets the trigger). It is simple and self-healing, but it makes some commits dramatically slower than others, and if ingestion has a tight latency budget the operator often disables it and then forgets to run anything else. That is the single most common root cause. **Asynchronous.** Compaction runs as a separate process — a service inside a long-running streaming writer, or a scheduled batch job pointed at the table. Ingestion latency stays flat, which is why streaming deployments prefer it, but now compaction is an independent job that can be down, under-resourced, or silently failing without anyone noticing, because ingestion keeps succeeding. Check that the job is actually running, that its resources scale with ingest volume, and that its failures alert someone. Either way, when several processes write to the same table you need Hudi's concurrency control configured correctly. A mismatch there can cause compaction to keep aborting on conflicts, which shows up as those repeated rollback instants. ## Other contributors worth ruling out - **Compaction is running but cannot keep up.** Completed compactions exist, yet the backlog still grows. Then it is a throughput problem: raise compaction parallelism and resources, compact more often, or reduce the write volume per commit. - **Too many small file groups.** If ingestion opens new groups constantly instead of packing into existing base files, compaction has far more slices to process. Revisit the file-sizing settings. - **Skewed file groups.** A few hot keys concentrate updates into a handful of groups, so those slices grow enormous logs while the rest are fine. Look at log volume per file ID, not just the total. - **The cleaner.** Cleaning is a separate action from compaction; a stalled cleaner inflates storage and listing cost but is not what makes the merge deeper. Do not confuse the two. ## Fixing it Short term, restore query latency by routing heavy readers to the read-optimized view, which reads base files only — accepting that its results are as of the last compaction and that recently deleted rows will still appear. Say that out loud to whoever consumes the data; silently serving stale results is worse than a slow query. Medium term, drain the backlog with a dedicated compaction run sized for it, then re-establish a steady schedule. Compact often enough that log depth per slice stays bounded, and make that bound your published freshness guarantee for the read-optimized view. Long term, monitor it as a service-level indicator, not as a folklore. The two metrics that matter are the age of the newest completed compaction and the number of delta commits (or bytes of log) outstanding per file group. Alert on both. And ask the design question: if this table's read pattern cannot tolerate any read-time merge, it may simply belong in a Copy-on-Write table, paying the cost on the write side instead.
- How do you tell a compaction backlog apart from ordinary data growth as the cause of the slowdown?Compare log depth per file slice over time, not table size. If base files are stable and log files per file ID keep climbing while the newest completed compaction ages, it is a backlog. If both base data and query time grow proportionally and compaction is current, you are simply scanning more data and the fix is pruning or clustering, not compaction.
- What is the risk of just switching heavy readers to the read-optimized view and leaving it there?You are serving data as of the last compaction, and if compaction stays broken that staleness is unbounded. Worse, recently deleted rows remain visible, because the delete lives only in an unread log block. It is a valid emergency lever and a bad steady state unless compaction is healthy and its interval is published as a freshness SLA.
- When would you conclude the table should not be Merge-on-Read at all?When the read pattern cannot tolerate any read-time merge and the write pattern does not actually need sub-compaction-interval freshness. If the table is written in batches and read constantly by latency-sensitive consumers, Copy-on-Write pays the merge cost once at write time and removes the whole class of incident — at the price of write amplification.
- How does clustering differ from compaction here, and would it help?Clustering rewrites base files into new file groups sorted by chosen columns to improve pruning and file sizing, and it lands as a `replacecommit`. It does not fold log files into base files, so it will not clear a merge backlog. Compaction is the fix; clustering addresses a different symptom, namely poor data skipping or badly sized files.
saying these in an interview costs you the question
- Blames the cleaner or retention instead of compaction lag
- Says Merge-on-Read compacts automatically so nothing can fall behind
- Proposes clustering to clear accumulated log files
- Treats the read-optimized view as a permanent fix with no staleness cost
- Reads only total table size instead of log depth per file slice