Incremental processing in Hudi: reading only what changed, and the timestamp that betrays you
Hudi's name contains "incrementals" and it is the capability the format was built around. Here is what an incremental read actually returns, why it is not a change log, how CDC mode differs, and the instant-versus-completion-time change in Hudi 1.x that makes the obvious checkpoint loop re-read every batch.
- The timeline is the whole mechanism
- What an incremental read returns
- The timestamp that betrays you
- CDC mode: the actual change log
- The
hudi_table_changesfunction - How far back can you read?
- Copy-on-write and merge-on-read
- Common misconceptions
- Frequently asked questions
- Conclusion
- References
- Trademarks
TL;DR
- An incremental read returns the current state of records that changed in a range. It is not a change log: a record updated three times appears once, and a deleted record does not appear at all.
- CDC mode is the change log. The same three commits produced 5 CDC rows — three inserts, an update with
beforeandafter, and a delete — against 2 rows from the incremental read.- CDC must be switched on when the table is created. Asking an ordinary table for CDC fails with
It isn't a CDC hudi table; there is no retrofit.- The trap. Hudi 1.x orders its timeline by completion time, and
begin.instanttimeis now a completion time — but_hoodie_commit_timestill holds the instant time. Feeding the latter back in as your next checkpoint re-reads the whole previous batch: 4 rows instead of 3 in the run below.- The fix is a supported one.
CALL show_commitsexposesstate_transition_time, which is the completion time, and that is what a checkpoint should store.hoodie.datasource.read.start.commitis not the config. Hudi 1.1.1 refuses it and namesbegin.instanttimein the error.
Hudi is an acronym: Hadoop Upserts Deletes and Incrementals. The first three get most of the attention, and the fourth is the one the design is actually organised around — a table that knows what changed since a point in time lets a downstream job process only that, instead of rescanning everything.
The mechanism is simple and the semantics are not, which is where this post spends its time.
The timeline is the whole mechanism
Every write to a Hudi table produces an instant on the timeline, under
.hoodie/. An incremental read is then a filter over that timeline: give me the
records touched by instants in this range.
In Hudi 1.x the timeline files carry two timestamps, and the filename says so:
.hoodie/timeline/
20260928041728131.commit.requested
20260928041728131.inflight
20260928041728131_20260928041743271.commit <- instant _ completion
20260928041743990_20260928041748181.commit
20260928041748709_20260928041752941.commit
The first number is when the write started; the second is when it completed. That pair is a 1.x change, and it exists so the timeline can be ordered by completion — which is what makes non-blocking concurrency possible, because a writer that started earlier but finished later must not appear earlier in the ordering.
Hold on to that distinction. It is the subject of the section that matters most here.
What an incremental read returns
Three commits against a table keyed on id:
# commit 1
[(1,"a",1), (2,"b",1), (3,"c",1)] # insert
# commit 2
[(1,"a2",2), (4,"d",2)] # update id=1, insert id=4
# commit 3
[(2,"b3",3)] # update id=2
Reading incrementally from the beginning:
(spark.read.format("hudi")
.option("hoodie.datasource.query.type", "incremental")
.option("hoodie.datasource.read.begin.instanttime", <instant>)
.load(PATH))
returns the latest state of every record touched in the range — not each change:
| id | val | _hoodie_commit_time |
|---|---|---|
| 1 | a2 |
commit 2 |
| 2 | b3 |
commit 3 |
| 4 | d |
commit 2 |
Record 1 was written twice and appears once, with its final value. There is no row describing what it used to be. And critically — a record deleted in the range does not appear at all, because there is no latest state to return.
That is the right shape for the common job, which is “re-apply the changed rows downstream”. It is the wrong shape for anything needing an audit trail, a before image, or deletes.
The timestamp that betrays you
Here is the pattern every incremental pipeline is written with: process a batch, record the last commit time you saw, use it as the start of the next run.
last = df.agg(F.max("_hoodie_commit_time")).first()[0] # checkpoint
# next run
.option("hoodie.datasource.read.begin.instanttime", last)
On Hudi 1.x that re-reads the entire previous batch, and nothing warns you.
begin.instanttime is compared against the completion time. But
_hoodie_commit_time is the instant time — the first number in the filename,
not the second. A commit that started at ...728131 and completed at ...743271
is not excluded by a filter of “completion time after ...728131”, because its
completion is later than that.
Measured on the table above:
begin = commit-1 INSTANT (from _hoodie_commit_time)
[(1,'a2'), (2,'b3'), (3,'c'), (4,'d')] <- 4 rows, commit 1 included
begin = commit-1 COMPLETION (from the timeline)
[(1,'a2'), (2,'b3'), (4,'d')] <- 3 rows, commit 1 excluded
Record 3 was only ever written by commit 1. Its presence in the first result is the whole bug: the batch you already processed comes back, every run, forever.
Whether that is harmless or catastrophic depends entirely on what the downstream job does. An idempotent upsert into another Hudi table absorbs it and merely wastes work. An append into a warehouse table, or anything that increments a counter, silently double-counts.
Getting the right value
You do not have to parse filenames. show_commits exposes both times, and the
second column is the one to checkpoint:
CALL show_commits(table => 'hudidemo.incr', limit => 10);
| commit_time | state_transition_time | action |
|---|---|---|
| 20260928041748709 | 20260928041752941 | commit |
| 20260928041743990 | 20260928041748181 | commit |
| 20260928041728131 | 20260928041743271 | commit |
state_transition_time is the completion time — it matches the second half
of the timeline filename exactly. Store that as your checkpoint, not
_hoodie_commit_time.
And while you are here: hoodie.datasource.read.start.commit is not the
option, whatever a tutorial written against another version says. Hudi 1.1.1
rejects it and tells you which to use:
org.apache.hudi.exception.HoodieException: Specify the start completion time
to pull from using option hoodie.datasource.read.begin.instanttime
CDC mode: the actual change log
When you need before images and deletes, incremental is the wrong query type. CDC mode is a different thing with a different schema, and it has to be enabled at table creation:
"hoodie.table.cdc.enabled": "true",
"hoodie.table.cdc.supplemental.logging.mode": "data_before_after"
Asking an ordinary table for CDC does not degrade gracefully:
It isn't a CDC hudi table on s3a://warehouse/hudi_incr_demo
There is no retrofit. A table that might ever need CDC has to be created with it on, which is a decision worth making deliberately rather than discovering later.
With it enabled, three commits — insert two rows, update one and insert another, delete one — produce a proper change log:
SELECT op, ts_ms, before, after
FROM hudi_table_changes('hudidemo.cdc', 'cdc', 'earliest');
op=i before=None after=(1,'a')
op=i before=None after=(2,'b')
op=i before=null after=(3,'c')
op=u before=(1,'a') after=(1,'a2')
op=d before=(2,'b') after=null
Five rows, one per change, with op naming it and both images present where they
exist. The same table read incrementally returns two rows — the current state
of records 1 and 3 — and says nothing at all about record 2 having been deleted.
Choosing between them
Incremental (latest_state) |
CDC | |
|---|---|---|
| Returns | current value of changed records | one row per change |
| Repeated updates | collapsed to the final value | one row each |
| Deletes | invisible | op=d with the before image |
| Before image | no | yes, with data_before_after |
| Enable when | always available | table creation only |
| Storage cost | none | supplemental logging on every write |
| Good for | re-applying changes downstream, idempotently | audit, replication, anything order-sensitive |
The rule of thumb: if the downstream operation is an idempotent upsert keyed the same way, incremental is cheaper and sufficient. If you need to know that something was deleted, or what it was before, you needed CDC enabled before you started.
There is also a supplemental.logging.mode cheaper than data_before_after —
storing only keys, or only the after image — which trades what CDC can tell you
for what it costs to write. Choose it at creation too.
The hudi_table_changes function
Both modes have a SQL surface that avoids the DataFrame options entirely:
-- current state of what changed
SELECT * FROM hudi_table_changes('db.tbl', 'latest_state', 'earliest');
-- the change log
SELECT * FROM hudi_table_changes('db.tbl', 'cdc', 'earliest');
-- bounded
SELECT * FROM hudi_table_changes('db.tbl', 'latest_state', <begin>, <end>);
'earliest' is a useful literal for the first run of a pipeline, when there is no
checkpoint yet and you want everything the timeline still holds. Which brings up
the limit on all of this.
How far back can you read?
Only as far as the timeline goes, and the timeline is deliberately bounded by two services that exist to keep it that way:
- Cleaning removes old file versions. Once the versions an instant referred to are gone, that instant cannot be served incrementally.
- Archival moves old instants out of the active timeline entirely.
So an incremental checkpoint is perishable. A pipeline that stops for longer than your retention window cannot resume from where it left off, and the honest recovery is a full re-read rather than a silently short one. This is the same retention arithmetic covered in table maintenance, seen from the reader’s side: retention policy is not only about storage, it is the bound on how long a downstream consumer may be offline.
Worth setting deliberately, and worth alerting on — a consumer whose checkpoint has fallen outside the window is a real incident, not a slow job.
Copy-on-write and merge-on-read
Both support incremental reads and the semantics above are identical. What differs is where the changed records are read from: on copy-on-write the reader picks up new base files, on merge-on-read it must also read log files for instants not yet compacted.
The practical consequence is that on a merge-on-read table, an incremental read competes with compaction for the same state, and a read spanning a compaction boundary does more work. That is the same read-side cost described in Hudi Merge-on-Read, and it is not a correctness difference.
Common misconceptions
“An incremental read gives me a change log.” It gives you the latest state of changed records. No before images, no per-change rows, and no deletes.
“I can turn CDC on when I need it.” It is a table-creation property. An existing table rejects CDC queries outright.
“_hoodie_commit_time is what I checkpoint.” On Hudi 1.x it is the instant
time and begin.instanttime wants the completion time. Checkpoint
state_transition_time from show_commits.
“start.commit is the modern option name.” Hudi 1.1.1 rejects it and names
begin.instanttime in the exception.
“Incremental reads are free.” They read the files the relevant instants touched. A range spanning many commits, or a merge-on-read table with uncompacted logs, is not obviously cheaper than a filtered snapshot scan.
“I can always resume from my checkpoint.” Only inside the retention window. Cleaning and archival will eventually make an old checkpoint unusable.
Frequently asked questions
How do I do the very first run with no checkpoint?
Use 'earliest' with hudi_table_changes, or the earliest instant from
show_commits. Then store the completion time you reached.
Does an incremental read see a record that was inserted and then deleted inside the range? No. There is no latest state to return, so it is absent — and indistinguishable from a record that never existed. CDC mode shows both events.
Can I read incrementally from a specific partition only? Yes, normal predicates apply on top of the incremental relation, and partition pruning still works.
Is the ordering of incremental results guaranteed?
No. You get a DataFrame; order it yourself if order matters. With CDC, ts_ms
and the commit ordering are what you sort on.
What about Iceberg’s equivalent?
Iceberg has incremental reads between snapshots with similar “changed rows”
semantics, and a changelog procedure closer to CDC. The
interoperability post
covers where the formats line up.
Conclusion
The incremental read is the capability Hudi is named for, and it does what it says as long as you know what “what changed” means here: the current value of touched records, with deletes absent and history collapsed. When that is not enough, CDC is a different query type with a different schema and a decision you must have made at table creation.
The part worth carrying away is the timestamp. Hudi 1.x moved the timeline to
completion-time ordering for good reasons, and the consequence is that the column
you naturally reach for — _hoodie_commit_time — is no longer the value
begin.instanttime expects. The failure is silent, it re-reads the previous
batch on every run, and it is one show_commits call away from being correct.
References
- Hudi incremental processing for the query types and their options
- Hudi CDC for the supplemental logging modes and the change-log schema
- Hudi timeline for instants, actions and states
- SQL procedures for
show_commitsand the rest of the inspection surface - The Hudi index for how a record is located in the first place
- Apache Hudi architecture for what a commit writes underneath all of this
Trademarks
Apache Hudi, Apache Spark, Apache Iceberg, Apache Hadoop, Apache and the Apache feather logo are either registered trademarks or trademarks of The Apache Software Foundation in the United States and other countries.
Found this useful?
These posts and tools are free. If one saved you an afternoon, you can buy me a coffee.