All posts

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.

11 min read Hudi

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 before and after, 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.instanttime is now a completion time — but _hoodie_commit_time still 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_commits exposes state_transition_time, which is the completion time, and that is what a checkpoint should store.
  • hoodie.datasource.read.start.commit is not the config. Hudi 1.1.1 refuses it and names begin.instanttime in 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

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.

Buy me a coffee