skip to content

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

level: seniorimportance: nice to knowfreq 28%

answer

  1. the checkpoint is a number on the timeline
  2. that number is stamped when work begins
  3. a slow writer finishes out of order
  4. the reader has already moved past its slot
  5. a later version records when it finished, too

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.

solid answer

~50 s

An incremental query reads records whose `_hoodie_commit_time` falls between a begin and an optional end instant, and the consumer checkpoints the highest instant it has consumed. In the 0.x timeline an instant is stamped when the action *starts*. Writer A begins at instant 100 and takes ten minutes; writer B begins at 101 and finishes immediately. The consumer reads up to 101 and checkpoints 101. When A finally completes, its commit sits at instant 100 — behind the checkpoint — so the next incremental read starting after 101 skips it entirely. This is why Hudi 1.x records a completion time on each instant and lets incremental reads position on when a commit became *visible* rather than when it began. Mitigations before that: a single writer, or checkpointing at the oldest inflight instant rather than the newest completed one.

code

python · 7 lines
python
df = (spark.read.format("hudi")
    .option("hoodie.datasource.query.type", "incremental")
    .option("hoodie.datasource.read.begin.instanttime", last_checkpoint)
    .load(table_path))

# checkpoint below the oldest still-inflight instant, not the newest completed one
next_checkpoint = min(oldest_inflight_instant, df_max_commit_time)

go deeper

for a junior

Know that a Hudi incremental query returns only records written after a given point on the timeline, and that each row carries the commit time that produced it.

for a middle

Explain that the consumer checkpoints an instant and reads forward from it, and that an instant's time marks when the write started rather than when it became visible.

for a senior

Reason about the concurrency gap end to end: why a slow writer's commit lands behind the checkpoint, what completion-time positioning fixes, and which mitigation you would pick on the version you actually run.

for a principal

Own the guarantee for the platform: decide whether tables allow concurrent writers at all, set timeline and cleaner retention against your slowest consumer's downtime, and define the documented recovery when a consumer falls outside the window.

## What an incremental query is A Hudi incremental query returns only the records written between two points on the timeline, rather than the current state of the table. You set the query type to incremental and give a begin instant, optionally an end instant; Hudi consults the timeline to find the commits in that range and reads the affected file slices, returning the records those commits wrote. Every row carries a `_hoodie_commit_time` metadata column identifying the instant that produced it, which is what makes the filter possible. This is Hudi's headline capability for chained pipelines: a downstream job pulls only what changed, so a bronze-to-silver step costs the size of the delta instead of the size of the table. ## The concurrency gap in the 0.x timeline The consumer keeps a checkpoint — the last instant it consumed — and uses it as the next read's begin instant. The subtlety is what the instant time means. In the long-standing 0.x timeline, an instant's time is allocated when the action *starts*, and it becomes visible only when the completed instant file appears. Those are two different moments, and with concurrent writers they can be far apart. Concretely: writer A starts a large batch and is allocated instant 100. Writer B starts a small batch at instant 101 and completes in seconds. The consumer runs, sees 101 completed, reads it, and checkpoints 101. Ten minutes later A completes — as instant 100. The next incremental read asks for instants greater than 101 and never looks back, so A's rows are silently lost from the incremental stream. They are in the table, and a snapshot query returns them, but the downstream table built from the increments is missing data. Silent divergence between the incremental consumer and the source is the worst shape this bug takes. ## What Hudi 1.x changes Hudi 1.x reworked the timeline so completed instants record a completion time in addition to the action's start time, and incremental reads can be positioned on completion time. Ordering by when a commit became visible is exactly the property a streaming consumer needs: a commit can never become visible before the point you have already consumed, so nothing can appear "behind" your checkpoint. This, together with the non-blocking concurrency work in the same line, is the structural fix. State which line you are on when you answer, because the correct mitigation differs. ## Mitigations without upgrading - **Single writer per table.** If only one process writes, start order and completion order coincide and the gap cannot open. This is the most common production answer, and the reason many teams route all ingest for a table through one job. - **Checkpoint conservatively.** Instead of checkpointing the newest completed instant, checkpoint below the oldest instant that is still inflight. You may re-read some commits, so the downstream write must be idempotent — which, with a record key and precombine field, an upsert already is. - **Serialise the writers.** Where a second writer exists for a different purpose (a backfill, a compaction-triggering job), keeping it from overlapping the ingest window sidesteps the issue. ## The other way incremental reads break Separately from concurrency, an incremental query fails when its begin instant has fallen off the retained timeline. The active timeline is bounded by archival retention, and the cleaner removes old file slices the query would need. A consumer that has been down for longer than that window cannot resolve its start point, and the honest recovery is a snapshot read to rebuild plus a fresh checkpoint — not quietly moving the checkpoint forward, which loses whatever was written in between. Sizing retention against your slowest consumer's tolerable downtime is the operational lesson. ## What good answers include A strong answer distinguishes *start time* from *visibility time* and explains why a checkpoint over the former is unsafe under concurrency; names the 1.x completion-time behaviour as the structural fix; offers single-writer or conservative-checkpoint mitigations for older versions; and separates this failure from the retention-window failure, which looks similar from the outside — rows missing downstream — but has a completely different cause and cure. It is also worth saying explicitly that a snapshot query is unaffected: the data is in the table throughout, and only the incremental stream skipped it.

  • How would you make an incremental Hudi consumer safe without upgrading the timeline?
    Either guarantee a single writer per table, so start and completion order coincide, or checkpoint below the oldest inflight instant rather than the newest completed one. The second re-reads some commits, so the downstream write must be idempotent — an upsert keyed on the record key with a precombine field already is, which makes the overlap harmless.
  • What happens when an incremental query's begin instant has been archived or its files cleaned?
    The read cannot be served: the instant is no longer on the active timeline and the file slices it referenced may be gone. The correct recovery is a snapshot read to rebuild the downstream table and a fresh checkpoint. Simply advancing the checkpoint to the current instant loses everything written during the outage.
  • Would a snapshot query also miss the slow writer's rows?
    No. Once the writer's instant completes, its file slices are part of the table and every snapshot query returns them. The gap is entirely in the incremental stream's positioning: the rows exist, but the consumer's checkpoint has already passed the instant that carries them, so it never asks for that range again.

saying these in an interview costs you the question

  • Assumes commits become visible in instant-time order under concurrency
  • Thinks the missing rows are absent from the table, not just from the stream
  • Fixes a lagging consumer by advancing the checkpoint to now
  • Confuses this with the cleaner having removed the needed files
  • Says incremental reads are unaffected by how many writers there are

context