Object storage is cheap and durable but has no notion of updating one row. If a stream of order events says order 42 is now shipped, a plain Parquet lake has to rewrite whole files or keep every version and make readers sort it out. Apache Hudi adds what is missing: a primary key per record, an index from key to file, a log of committed changes called the timeline, and background services that keep files a sensible size. That turns a directory of Parquet into a table that accepts upserts and deletes at streaming rates and can itself be read as a stream.
This article looks at Hudi from the streaming side. The Spark write and read paths are covered in Spark + Hudi architecture; here the focus is a continuously running writer, usually Flink, feeding a merge-on-read table, the services that run alongside it, and downstream jobs that consume the table incrementally. It ends with a worked sizing example and the failure modes that show up in production.
The storage model: file groups, slices, base and log files
A Hudi table is a directory tree on storage plus a metadata folder, .hoodie. Within each partition, records are spread across file groups, each identified by a file ID. A record key, once placed in a file group, stays there until clustering moves it; that stickiness is what makes an update a local operation.
A file group is a sequence of file slices. Each slice has an optional base file in a columnar format such as Parquet and, for merge-on-read tables, zero or more log files holding row-oriented blocks of later inserts, updates and deletes. A new slice begins whenever compaction or a copy-on-write update writes a new base file.
| Table type | An update does | Write cost | Read cost |
|---|---|---|---|
| Copy-on-write (COW) | Rewrites the base file containing the key | High: a whole file per touched group | Low: read Parquet only |
| Merge-on-read (MOR) | Appends a block to a log file | Low: append only | Higher until compaction merges logs into a new base file |
For streams that update existing keys, MOR is usually the right type. A COW table receiving updates spread across many file groups rewrites hundreds of megabytes per commit to change a few kilobytes, and commit latency grows with the spread of the keys rather than the volume of data.
The timeline is the source of truth
Every change to a Hudi table is an instant on the timeline. An instant has an action and a state. Actions include commit and deltacommit (COW and MOR writes), compaction, clean, rollback, replacecommit (used by clustering and insert-overwrite), savepoint and restore. States move from requested to inflight to completed.
Readers consider only completed instants. Files written by an inflight or failed instant may sit on storage, but no query returns their rows, and a later rollback deletes them. This is how Hudi gets atomic commits on storage that has no transactions: the data files are written first, and a small metadata file completing the instant is the switch that makes them visible.
For streaming consumers the timeline is also a change feed. An incremental query asks for records changed between two instants, and a streaming reader polls the timeline for instants it has not yet seen.
The Flink write path and exactly-once
The diagram shows a Flink job writing a MOR table. Records are keyed by record key, a bucket-assign step maps each key to a file group, and write tasks append records to the log files of their file groups. A coordinator on the JobManager owns the timeline: it starts an instant, and when a Flink checkpoint completes it collects the write metadata from every task and completes the instant.
Tying commits to checkpoints is what gives end-to-end exactly-once behaviour with a replayable source such as Kafka. If the job fails before the checkpoint completes, the source offsets roll back to the previous checkpoint, the inflight instant is rolled back, and the same records are written again into a new instant. The exactly-once guarantees and their limits are discussed generally in exactly-once processing in streaming.
The consequence for latency is direct: data becomes visible no sooner than the checkpoint that commits it. A one-minute checkpoint interval means readers see changes within roughly a minute plus commit time. Shortening checkpoints improves freshness but creates more instants, more small log blocks and more compaction work.
A minimal Flink SQL table definition for an upserting MOR sink looks like this.
CREATE TABLE orders_hudi (
order_id STRING,
customer_id STRING,
status STRING,
amount DECIMAL(12, 2),
updated_at TIMESTAMP(3),
dt STRING,
PRIMARY KEY (order_id) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
'connector' = 'hudi',
'path' = 's3a://lake/orders_hudi',
'table.type' = 'MERGE_ON_READ',
'precombine.field' = 'updated_at', -- newest version wins on key collisions
'index.type' = 'BUCKET',
'hoodie.bucket.index.num.buckets' = '64',
'compaction.async.enabled' = 'true',
'compaction.delta_commits' = '5'
);
INSERT INTO orders_hudi SELECT * FROM orders_kafka;The precombine field decides which version wins when two records with the same key arrive in one batch, so it must be a field that increases with each change, such as an update timestamp or a source log position. Recent Hudi documentation also shows this setting under the name ordering.fields; check the name your release expects.
Choosing an index for a stream
Before writing a record, Hudi has to know whether its key already exists and in which file group. For streaming writers that lookup sits on the hot path, so the index choice decides throughput.
| Index | How it finds the file group | Good for | Watch out for |
|---|---|---|---|
FLINK_STATE (Flink default) | Key-to-file-group map kept in Flink keyed state | Updates spread across all history | State grows with distinct keys; a new job on an existing table needs index.bootstrap.enabled to load it |
BUCKET | Hash of the key modulo a fixed bucket count per partition | Predictable cost, no lookup state | Bucket count (default 4) is fixed per partition; too few makes huge file groups, too many makes tiny ones |
| Record-level index (metadata table) | Key-to-location map stored in Hudi's metadata table | Large tables with random updates across writers | Extra write work on every commit; check engine support in your release |
The bucket index trades flexibility for simplicity. Because the file group is a pure function of the key, nothing needs to be looked up, and multiple writers agree on placement. Size the bucket count so that each bucket's base file lands near your target file size, typically a few hundred megabytes, at the partition's expected volume.
Table services in a streaming job
Three services keep a streaming table healthy.
- Compaction merges a file group's log files into a new base file. In Flink it runs asynchronously by default for MOR tables, and the default trigger schedules a compaction after
compaction.delta_commitsdelta commits, which defaults to 5. Without it, snapshot reads merge more and more log blocks and slow down steadily. - Clustering rewrites file groups to change layout: merging small files or sorting by query columns. It produces a
replacecommitand is usually run as a separate job rather than inside the ingest job. - Cleaning deletes file slices that are no longer needed. In Flink,
clean.retain_commitsdefaults to 30, and the documentation notes this also bounds how far back an incremental reader can pull.
That last point is easy to miss. If a downstream incremental job falls behind by more than the retained commits, the files for the instants it needs have been cleaned and it can no longer catch up incrementally; it has to restart from a snapshot. Set retention from the longest outage a consumer must survive, not from storage cost alone.
Running compaction inside the ingest job is simplest but competes for task slots with writing. Large tables often schedule compaction in the ingest job and execute it in a separate job, so a slow compaction never back-pressures ingestion. Back pressure itself is discussed in streaming back pressure.
HoodieStreamer: ingestion without writing a job
If you do not want to own a Flink job, Hudi ships a Spark-based ingestion utility. It was formerly called DeltaStreamer and is now org.apache.hudi.utilities.streamer.HoodieStreamer. In continuous mode it loops fetch, transform and write, with an optional minimum interval between syncs.
spark-submit \
--class org.apache.hudi.utilities.streamer.HoodieStreamer \
hudi-utilities-bundle.jar \
--table-type MERGE_ON_READ \
--op UPSERT \
--source-class org.apache.hudi.utilities.sources.JsonKafkaSource \
--source-ordering-field updated_at \
--target-base-path s3a://lake/orders_hudi \
--target-table orders_hudi \
--props /etc/hudi/orders.properties \
--continuous \
--min-sync-interval-seconds 60The properties file holds the record key, partition path, Kafka settings and schema provider. HoodieStreamer suits teams that already run Spark and want micro-batch latency of a minute or more; Flink suits lower latency and pipelines that already transform data in Flink.
Reading a Hudi table as a stream
Hudi tables are sources as well as sinks. A Flink streaming read polls the timeline for new completed instants and emits the records they changed.
CREATE TABLE orders_changes WITH (
'connector' = 'hudi',
'path' = 's3a://lake/orders_hudi',
'table.type' = 'MERGE_ON_READ',
'read.streaming.enabled' = 'true',
'read.start-commit' = '20261001000000', -- instant time to start from
'read.streaming.check-interval' = '30' -- seconds; default 60
) LIKE orders_hudi (EXCLUDING OPTIONS);By default a streaming read emits the latest state of each changed record per instant. Setting changelog.enabled keeps intermediate changes, and cdc.enabled produces a change-data-capture stream with before and after images, at the cost of extra data written. Pick based on whether downstream logic needs every transition or only the latest value.
Restarting a streaming reader from a Flink savepoint resumes from the instant it last emitted; see Flink savepoints for upgrades and rescaling.
Worked example: sizing an order table
Suppose an orders topic carries 20,000 events per second at peak, about 60% updates to orders from the last week and 40% new orders, at roughly 500 bytes each. That is about 10 MB per second, or 600 MB per minute before compression.
- Checkpoint interval. Analysts accept two minutes of staleness, so checkpoint every 60 seconds. Each checkpoint produces one delta commit of about 1.2 million records.
- Index. Updates cluster in recent days, keys are stable and there may later be a second writer for backfills, so use the bucket index with daily partitions.
- Bucket count. Peak is not average. If the daily average is 5,000 events per second, a day brings about 430 million events, of which about 170 million are new orders. The table stores current state, so a day partition (partitioned by order date) holds about 170 million rows at around 100 bytes compressed each, or about 17 GB. At a 512 MB target that is about 34 buckets; choose 64 to leave growth headroom.
- Compaction. With
compaction.delta_commitsat 5, each touched file group is compacted roughly every five minutes. Read-optimized queries are therefore up to about five minutes staler than snapshot queries. - Cleaning. The downstream incremental job must survive a four-hour outage: 240 commits at one per minute. Set
clean.retain_commitsto around 300, and budget the extra storage for retained slices.
Validate each number with a day of real traffic: watch commit duration, log blocks per file group before compaction and compaction duration. If commit duration approaches the checkpoint interval, the job will fall behind.
Failure modes
- Compaction falls behind. Log files pile up, snapshot reads slow down, and memory pressure rises on merge. Alert on pending compaction instants and run compaction in a separate job if needed.
- Checkpoint timeouts. Writes are flushed at checkpoints, so slow storage or skewed buckets make checkpoints time out, nothing commits and the job restarts in a loop. Look for one hot bucket before adding parallelism.
- Lost index state. Starting a new Flink job with the state index and no bootstrap creates duplicate keys, because existing keys are unknown and get inserted into new file groups.
- Wrong precombine field. Using a field that does not increase, such as ingestion time on a replayed topic, lets stale versions overwrite newer ones.
- Cleaned-away increments. An incremental reader paused longer than the retention window fails or silently skips. Monitor consumer lag in instants, not just time.
- Concurrent writers without locking. Two jobs writing the same table need a concurrency control mode and a lock provider configured; without them, commits can conflict or corrupt the timeline.
Trade-offs against other table formats
Hudi's distinguishing features for streaming are record-level indexing, merge-on-read with built-in async compaction, and incremental queries driven by the timeline. Apache Iceberg and Delta Lake also support row-level changes and change feeds, with different mechanisms and engine coverage; Iceberg's architecture is a useful comparison. Apache Paimon takes an LSM-tree approach aimed at Flink-first streaming. Choose by the engines that must read and write the table, how update-heavy the stream is, and which services your team is prepared to operate, and benchmark on your own data before committing.
What to do next
- Classify your stream by update ratio and key locality; choose MOR for update-heavy streams and COW only for append-mostly ones.
- Pick the precombine field from a value that always increases per key, and test it by replaying the same topic twice.
- Choose the index: bucket for predictable placement, Flink state for keys spread across history, and size buckets from current-state volume per partition.
- Set the checkpoint interval from your freshness target, then confirm commit duration stays well below it under peak load.
- Set compaction cadence and cleaning retention from read latency targets and the longest consumer outage you must survive.
- Add alerts for pending compactions, checkpoint failures, commit duration and incremental consumer lag measured in instants.