apache iceberg v3 is the third revision of the Iceberg table format specification, and it is the release that turns Iceberg from "a better way to lay out Parquet in object storage" into a format that can track every row's identity across time, delete rows without rewriting gigabytes of data, and carry semi-structured and geospatial values natively. The format-version integer stamped in a table's metadata is not cosmetic: it is a contract between every engine that reads and writes the table, and bumping it from 2 to 3 unlocks three headline capabilities — deletion vectors that overhaul merge-on-read, row lineage that gives each row a durable identity and a last-modified sequence number, and a set of new binary types (variant, geometry, geography, nanosecond timestamps) that let a lakehouse store JSON and spatial data without bolting on a side table.
This guide is the walkthrough you wished existed the first time an interviewer asked "what actually changed between Iceberg v2 and v3?", or "why did the community replace positional delete files with deletion vectors?", or "how would you build incremental CDC on top of Iceberg row lineage?" It opens the format in layers: the four-level metadata tree (table format metadata, manifest lists, manifests, and data files) that every version shares, the snapshots and sequence numbers that give Iceberg its atomic commits and time travel, the partition evolution that lets you re-partition without a rewrite, and then the three v3 features in depth. Each section pairs a teaching block with a Solution-Tail interview answer — code, a step-by-step trace, an output table, then a concept-by-concept breakdown of why it works. Examples use Spark SQL and PyIceberg, but the spec-level mental model carries to Flink, Trino, Dremio, and the managed Iceberg services in Snowflake, BigQuery, and Athena.
When you want hands-on reps immediately after reading, drill the database practice library →, work the storage-and-pipeline problems in the data-processing practice library →, and tune scan cost on the optimization practice library →.
On this page
- The Iceberg table format recap — metadata, manifests, snapshots
- Deletion vectors and merge-on-read
- Row lineage — row IDs and sequence numbers
- New types — variant, geometry, and binary
- Migration, compatibility, and engine support
- Cheat sheet — Apache Iceberg v3 recipes
- Frequently asked questions
- Practice on PipeCode
1. The Iceberg table format recap — metadata, manifests, snapshots
A table is a tree of immutable metadata files pointing at immutable data files — and the format version governs what that tree may contain
The one-sentence invariant: an Iceberg table is not a directory of data files but a versioned tree of immutable metadata — a single current metadata JSON that lists snapshots, each snapshot pointing at one manifest list, each manifest list pointing at manifest files, each manifest listing data files and delete files with their partition values and column statistics — and every write produces a brand-new metadata file swapped in with a single atomic pointer update, which is exactly why Iceberg gives you snapshot isolation, time travel, and safe concurrent writers without a locking table server. Understanding this tree is the prerequisite for understanding every v3 feature, because deletion vectors, row lineage, and the new types are all changes to what the tree is allowed to record, not changes to how the tree is walked.
The four layers, top to bottom.
- Catalog pointer. The catalog (Hive Metastore, AWS Glue, a REST catalog, Nessie, JDBC) stores exactly one thing per table: the path of the current metadata file. A commit is "atomically compare-and-swap the current-metadata pointer from old to new." That single atomic swap is the whole concurrency story.
-
Table metadata file. A JSON document (
v42.metadata.json) holding the table's schema(s), partition spec(s), sort order(s), properties, the list of snapshots, and the current snapshot id. It also holds theformat-versioninteger — 1, 2, or 3 — that gates which features are legal. - Manifest list. One Avro file per snapshot. Each row describes one manifest file plus summary stats (partition ranges, added/existing/deleted data-file counts) so a planner can prune whole manifests without opening them.
-
Manifest file. An Avro file listing individual data files and delete files, each entry carrying the file path, partition tuple, record count, column-level lower/upper bounds, null counts, and — critically — a
sequence_number. This is where scan planning does most of its pruning.
Snapshots, sequence numbers, and isolation.
- A snapshot is a complete logical table state. It references every data file (and delete file) that is live at that point in time. Time travel is "read the manifest list of snapshot N"; rollback is "set the current snapshot id back to N."
- Sequence numbers order commits. Every successful commit gets a monotonically increasing sequence number. Data files and delete files inherit the sequence number of the commit that added them. A delete file applies to a data file only when the data file's sequence number is less than or equal to the delete file's sequence number — this single rule is how Iceberg keeps merge-on-read correct without timestamps.
- Optimistic concurrency. Two writers can build new metadata trees in parallel; the loser of the compare-and-swap retries against the winner's new snapshot. No writer blocks a reader; a reader always sees a consistent snapshot.
Hidden partitioning and partition evolution.
-
Partitioning is a transform, not a directory convention. You declare
days(event_ts)orbucket(16, user_id); Iceberg records the partition value in each manifest entry and prunes on it. Queries never reference the partition column explicitly — noWHERE dt = '2026-09-05'folklore. -
Partition evolution changes the spec without rewriting data. Adding
hours(event_ts)to a table that was partitioned bydays(event_ts)creates a new partition spec id; old data files keep their old spec, new files use the new spec, and the planner prunes each file against the spec it was written with. - Format version is orthogonal but related. v1 was analytic tables only (no row-level deletes). v2 added merge-on-read delete files. v3 adds deletion vectors, row lineage, and the new types. Each bump is additive; a v3 table can still be read as the sum of its parts by an engine that understands v3.
What interviewers listen for.
- Do you describe a commit as "an atomic swap of the current-metadata pointer" rather than "writing files to a folder"? — senior signal.
- Do you name the sequence-number rule for when a delete applies to a data file? — required answer for any MoR discussion.
- Do you distinguish snapshot (a table state) from manifest (a list of files) from data file (Parquet/ORC/Avro)? — required vocabulary.
- Do you say partition evolution needs no rewrite because each file remembers its own spec id? — senior signal.
Worked example — reading the metadata tree of a real table
Detailed explanation. The fastest way to internalise the tree is to inspect one. Iceberg exposes metadata tables (.snapshots, .manifests, .files, .history) that you can query like any other table. Walk through inspecting an orders table and mapping each result back to a layer of the tree.
-
The catalog holds the pointer to
s3://warehouse/db/orders/metadata/v7.metadata.json. -
The metadata JSON lists three snapshots; the current one is snapshot
id = 30291. - The manifest list for that snapshot names two manifest files.
- The manifests together list five data files and one delete file.
Question. Given the orders table, produce the queries that reveal each layer and interpret the output.
Input.
| Layer | Metadata table | What it reveals |
|---|---|---|
| Snapshots | orders.snapshots |
snapshot id, parent, sequence number, manifest list |
| Manifests | orders.manifests |
manifest paths, added/deleted file counts |
| Files | orders.files |
data files, partition, record count, column bounds |
| History | orders.history |
which snapshot was current at which time |
Code.
-- Spark SQL against an Iceberg catalog named `prod`
-- 1. Snapshots: the list of table states
SELECT snapshot_id, parent_id, sequence_number, operation, manifest_list
FROM prod.db.orders.snapshots
ORDER BY committed_at;
-- 2. Manifests referenced by the current snapshot
SELECT path, added_data_files_count, existing_data_files_count, deleted_data_files_count
FROM prod.db.orders.manifests;
-- 3. Data + delete files, with partition and stats
SELECT content, file_path, record_count, partition, lower_bounds, upper_bounds
FROM prod.db.orders.files;
-- 4. History: current-snapshot timeline (time travel targets)
SELECT made_current_at, snapshot_id, is_current_ancestor
FROM prod.db.orders.history;
Step-by-step explanation.
-
orders.snapshotsreturns one row per commit. Thesequence_numbercolumn is the monotonic counter;manifest_listis the Avro file for that snapshot.operationisappend,overwrite,delete, orreplace— a one-word summary of what the commit did. -
orders.manifestslists the manifest files of the current snapshot. Theadded/existing/deletedcounts let you see at a glance whether the last commit added files (append) or rewrote them (compaction/overwrite). -
orders.filesis the leaf layer.content = 0is a data file,content = 1is a positional delete file,content = 2is an equality delete file. Thepartition,lower_bounds, andupper_boundscolumns are exactly what the planner uses to skip files. -
orders.historymaps wall-clock time to snapshot ids. This is the lookup you use forSELECT ... FOR SYSTEM_TIME AS OF '2026-09-01'— the engine resolves the timestamp to a snapshot id, then reads that snapshot's manifest list. - Nothing here mutates: every file named is immutable. A new commit writes new metadata JSON, a new manifest list, and (usually) new manifests, then swaps the catalog pointer. The old files remain for time travel until expiry.
Output.
| Layer | Example value | Role |
|---|---|---|
| Current metadata | v7.metadata.json |
root of the tree |
| Current snapshot | id=30291, seq=5 |
the live table state |
| Manifest list | snap-30291-...avro |
names 2 manifests |
| Manifest | manifest-a.avro |
lists 5 data files + 1 delete |
| Data file | data-0002.parquet |
120,000 rows, dt=2026-09-05
|
Rule of thumb. When debugging any Iceberg oddity — a slow scan, a "missing" row, a stuck compaction — start at snapshots, walk to manifests, then to files. The bug is almost always visible as a file the planner is (or isn't) pruning, and the metadata tables show you exactly which layer to blame.
Worked example — how a commit produces a new snapshot
Detailed explanation. The atomic-swap commit is the mechanism behind every Iceberg guarantee. Walk through what physically happens when a writer appends 10,000 rows to orders, so the "atomic pointer swap" stops being an abstraction.
-
Before. Current metadata is
v7.metadata.json, current snapshot30291, sequence number 5. -
The write. The engine writes one new
data-0006.parquet, then a new manifest listing it, then a new manifest list combining the old manifests plus the new one. -
The swap. The engine writes
v8.metadata.json(snapshot30292, sequence 6) and asks the catalog to compare-and-swap the pointer from v7 to v8.
Question. Trace the file writes and the commit for a 10,000-row append, and explain what happens if a second writer commits concurrently.
Input.
| Step | Artifact written | Immutable? |
|---|---|---|
| 1 |
data-0006.parquet (10k rows) |
yes |
| 2 |
manifest-new.avro (lists data-0006) |
yes |
| 3 |
snap-30292.avro (manifest list) |
yes |
| 4 |
v8.metadata.json (snapshot 30292) |
yes |
| 5 | catalog CAS: v7 → v8 | atomic |
Code.
# PyIceberg — append a batch and inspect the new snapshot
from pyiceberg.catalog import load_catalog
import pyarrow as pa
catalog = load_catalog("prod") # REST/Glue/Hive catalog config
table = catalog.load_table("db.orders")
before = table.current_snapshot().snapshot_id
print("before snapshot:", before, "seq:", table.current_snapshot().sequence_number)
batch = pa.table({
"id": pa.array(range(1, 10_001), pa.int64()),
"customer_id": pa.array([7] * 10_000, pa.int64()),
"total_cents": pa.array([1500] * 10_000, pa.int64()),
"status": pa.array(["pending"] * 10_000),
})
table.append(batch) # writes data file + manifest + snapshot, then CAS
table.refresh()
after = table.current_snapshot()
print("after snapshot:", after.snapshot_id, "seq:", after.sequence_number)
print("operation:", after.summary.operation) # 'append'
Step-by-step explanation.
- The engine writes the data file first. Nothing is visible to readers yet — the file exists in object storage but no metadata references it, so no snapshot includes it.
- It writes a new manifest listing
data-0006.parquetwith its partition tuple, record count, and column bounds. Then it writes a new manifest list that references the existing manifests plus the new manifest. - It writes
v8.metadata.jsonwith a new snapshot30292, parent30291, and sequence number 6 (one more than the parent's 5). - It issues the catalog compare-and-swap: "set the current pointer to v8 only if it is currently v7." If it succeeds, the append is now atomically visible; if it fails, another writer won the race.
- On a concurrent-writer conflict, the CAS fails. The engine refreshes to the winner's snapshot, re-validates that its append still makes sense (an append almost always does), rebuilds
v9.metadata.jsonon top of the winner's tree, and retries the CAS. Readers never see a half-written state because they only ever follow the committed pointer.
Output.
| Field | Before | After |
|---|---|---|
| Current metadata | v7.metadata.json |
v8.metadata.json |
| Snapshot id | 30291 | 30292 |
| Sequence number | 5 | 6 |
| Data files | 5 | 6 |
| Visibility | atomic on CAS success | atomic on CAS success |
Rule of thumb. Every Iceberg commit is copy-on-write at the metadata level even when it is merge-on-read at the data level. New metadata files are cheap (kilobytes); the atomic pointer swap in the catalog is the only serialization point, which is why Iceberg scales writers far better than a lock-per-table warehouse.
Worked example — partition evolution without a rewrite
Detailed explanation. Partition evolution is the feature that most differentiates Iceberg from Hive-style tables, and it is a favourite interview probe. A table partitioned by days(event_ts) starts receiving 100× the traffic; you want hours(event_ts) going forward but you cannot afford to rewrite two years of history. Iceberg lets you change the spec and keep both.
-
Old spec (id 0).
days(event_ts)— one partition per day. -
New spec (id 1).
hours(event_ts)— one partition per hour, applied to new writes only. - Reads transparently prune each file against the spec it was written with.
Question. Evolve events from daily to hourly partitioning and show that historical files are untouched while new files use the finer spec.
Input.
| Aspect | Before | After |
|---|---|---|
| Partition spec id 0 | days(event_ts) |
still used by old files |
| Partition spec id 1 | — |
hours(event_ts) for new files |
| Data rewrite | — | none required |
| Pruning | per-file, per-spec | per-file, per-spec |
Code.
-- 1. Evolve the partition spec (metadata-only operation)
ALTER TABLE prod.db.events
REPLACE PARTITION FIELD days(event_ts) WITH hours(event_ts);
-- 2. New writes now land in hourly partitions; old files keep daily
INSERT INTO prod.db.events
SELECT * FROM staging.events_incoming;
-- 3. Confirm both specs coexist and which files use which
SELECT spec_id, count(*) AS file_count
FROM prod.db.events.files
GROUP BY spec_id
ORDER BY spec_id;
Step-by-step explanation.
-
REPLACE PARTITION FIELDmutates only the table metadata — it appends a new partition spec (id 1) and marks it current. Not a single data file is read or rewritten; the operation commits in milliseconds regardless of table size. - New inserts are planned against spec id 1, so their manifest entries record an hour partition value. Old files still carry their day partition value under spec id 0.
- The
filesmetadata table exposesspec_idper file. Grouping by it shows the split: millions of legacy files under spec 0, new files under spec 1. - At read time, the planner evaluates a
WHERE event_ts BETWEEN ...predicate against each file using that file's own spec. A daily file is pruned at day granularity; an hourly file at hour granularity. Correctness holds because the predicate is on the source column, not on a derived partition column. - If you later want history at hourly granularity too, you run a
rewrite_data_filescompaction that re-partitions the old files under spec 1 — but that is an explicit, schedulable maintenance action, not a requirement for the schema change to take effect.
Output.
| spec_id | Transform | File count | Rewritten? |
|---|---|---|---|
| 0 | days(event_ts) |
730 (2 years) | no |
| 1 | hours(event_ts) |
growing | no |
Rule of thumb. Treat partition evolution as a free metadata change and treat re-partitioning history as a separate, optional compaction. Conflating the two — "we can't change partitioning without a full rewrite" — is the Hive-era assumption Iceberg exists to kill.
Data engineering interview question on the Iceberg metadata tree
A senior interviewer often opens with: "Explain how an Iceberg table stays consistent when three Spark jobs and a Flink job all write to it at once, and how a reader that started a five-minute query never sees a torn state. Then tell me what physically changes on disk and in the catalog for a single append, and where the format-version integer comes into it."
Solution Using atomic metadata swap, sequence-number ordering, and snapshot isolation
Iceberg commit + isolation model (whiteboard answer)
====================================================
WRITE PATH (per committing job)
1. Write data files -> immutable Parquet/ORC/Avro, not yet referenced
2. Write manifest(s) -> Avro; list new files + partition + stats + seq_no
3. Write manifest list -> Avro; the snapshot's file inventory
4. Write vN+1.metadata.json -> new snapshot, seq_no = parent.seq_no + 1
5. Catalog compare-and-swap -> pointer: vN -> vN+1 (the ONLY lock point)
success -> commit visible atomically
failure -> refresh to winner, re-validate, rebuild vN+2, retry (step 4-5)
READ PATH (per query)
1. Resolve current metadata pointer once, at planning time
2. Pin the snapshot id for the whole query
3. Plan against that snapshot's manifest list -> manifests -> files
4. Apply delete files where data_file.seq_no <= delete_file.seq_no
FORMAT VERSION
- Stored in vN.metadata.json as "format-version": 1 | 2 | 3
- Gates legal content: v2 enables delete files; v3 enables deletion
vectors, row lineage, and variant/geo/nanosecond types
Step-by-step trace.
| Actor | Sees / does | Consistency mechanism |
|---|---|---|
| Spark job A | CAS v7→v8 succeeds (seq 6) | wins the race |
| Spark job B | CAS v7→v8 fails | refresh, rebuild v9, retry |
| Flink job | CAS v9→v10 (seq 8) | serialized by the catalog |
| Reader R | pinned snapshot 30292 | never re-reads the pointer mid-query |
| Reader R | applies deletes seq ≤ 6 | sequence-number rule |
After the four writers finish, the catalog pointer has advanced through v8, v9, v10 in a strict serial order even though the jobs ran concurrently; each losing CAS simply retried on top of the winner. Reader R, which pinned snapshot 30292 at planning time, sees exactly the files live at sequence number 6 for its entire five-minute run, immune to every commit that lands after it started.
Output:
| Property | Guarantee | Source of the guarantee |
|---|---|---|
| Atomicity | commit is all-or-nothing | single catalog CAS |
| Isolation | reader sees one snapshot | snapshot pinned at plan time |
| Concurrency | writers never block readers | new metadata written out-of-band |
| Delete correctness | deletes apply to older data only |
data.seq ≤ delete.seq rule |
| Feature gating | v3 content only in v3 tables |
format-version integer |
Why this works — concept by concept:
- Atomic pointer swap — the catalog stores one mutable value, the current metadata path. Compare-and-swap on that single value is the entire serialization mechanism; everything else is immutable files written out-of-band, so there is no torn state to observe.
-
Sequence numbers — a monotonic per-commit counter inherited by every file. Ordering deletes against data by sequence number (delete applies when
data.seq ≤ delete.seq) makes merge-on-read deterministic without relying on wall-clock timestamps that could skew. - Snapshot isolation — a reader resolves the pointer once and pins the snapshot id. Because snapshots and their manifests are immutable, the reader's view cannot change mid-query no matter how many writers commit.
-
Format version gate — the
format-versioninteger tells every engine which spec features are legal in this table. It prevents a v2-only writer from silently corrupting a table that a v3 reader expects, and it is the switch the rest of this article turns on. - Cost — O(1) catalog CAS per commit plus O(new files) manifest writes; reads are O(matching files) after manifest-level pruning. The metadata tree adds kilobytes per commit, a negligible tax for atomic, isolated, time-travelable tables at object-storage scale.
SQL
Topic — database
Database internals and table-format problems
2. Deletion vectors and merge-on-read
A deletion vector is one compact bitmap per data file that marks deleted row positions — replacing v2's swarm of small positional delete files
The mental model in one line: a deletion vector is a per-data-file Roaring bitmap, stored as a blob inside a Puffin file, whose set bits are the row positions that have been logically deleted from that data file — a reader scans the data file, consults the single deletion vector attached to it, and skips exactly those positions, which is the v3 form of merge-on-read and which replaces v2's pattern of writing many small positional delete files that each named a handful of (file_path, position) pairs. The deletion vector is the single most impactful v3 change for tables with frequent row-level deletes and updates, because it collapses the read-time and compaction-time cost of merge-on-read from "open N tiny delete files" to "read one dense bitmap."
Copy-on-write vs merge-on-read — the two update strategies.
- Copy-on-write (CoW). An update or delete rewrites every data file that contains an affected row. Reads are as fast as a plain table (no merge step) but writes amplify badly — deleting one row in a 500 MB file rewrites all 500 MB. CoW is right for read-heavy tables with rare, batchy mutations.
- Merge-on-read (MoR). A delete writes a small marker instead of rewriting the data file; inserts write new data files; the reader merges data and delete markers at query time. Writes are cheap; reads pay a merge cost that grows with the number of delete markers. MoR is right for write-heavy tables with frequent point mutations.
-
The dial is per-operation. Table properties like
write.delete.mode,write.update.mode, andwrite.merge.modechoosecopy-on-writeormerge-on-readindependently. Deletion vectors are the v3 MoR delete representation.
How v2 did merge-on-read, and why it hurt.
-
Positional delete files named
(file_path, position)pairs — "indata-0002.parquet, positions 5 and 6 are deleted." Precise, but a streaming pipeline that deletes a few rows per minute produces thousands of tiny delete files per day. -
Equality delete files named predicates — "any row where
id = 42is deleted." Powerful for CDC (you delete by key without knowing the position), but expensive to apply because the reader must evaluate the predicate against every candidate row. - The small-file swarm. Both kinds accumulate. A reader of a hot data file might have to open dozens of delete files and union their effects. Compaction had to merge delete files constantly, and planning had to track which deletes applied to which data files by sequence number.
What a deletion vector changes.
- One vector per data file. v3 mandates at most one deletion vector per data file. New deletes to that file update the existing vector rather than writing another file, so the swarm never forms.
- Roaring bitmap encoding. The bitmap uses the Roaring compression scheme, which is dense for clustered deletes and sparse-efficient for scattered ones — a few bytes for a handful of deletes, still compact for millions.
-
Stored in Puffin. Deletion vectors live as
deletion-vector-v1blobs inside Puffin files (the same sidecar format Iceberg uses for statistics like Theta sketches). A manifest entry points a data file at its vector's Puffin blob by offset and length. - Positional deletes are superseded. v3 replaces positional delete files with deletion vectors. Equality deletes remain available for key-based CDC deletes, but the positional swarm is gone.
Common interview probes on deletion vectors.
- "Why did v3 replace positional delete files?" — required answer: to eliminate the small-delete-file swarm and give O(1) per-file delete lookup via one bitmap.
- "How many deletion vectors can a data file have?" — required answer: exactly one; new deletes update it in place.
- "Where is a deletion vector stored?" — a Puffin blob (
deletion-vector-v1), referenced from the manifest by offset/length. - "When does a deletion vector apply to a data file?" — same sequence-number rule as all deletes:
data_file.seq ≤ delete.seq.
Worked example — the bitmap mechanics of a deletion vector
Detailed explanation. A deletion vector is conceptually a set of integers — the deleted row positions within one data file. Roaring bitmaps store such a set compactly. Walk through what the bitmap holds for a data file whose rows 2, 5, and 6 (zero-based positions) have been deleted.
-
Data file.
data-0002.parquet, 8 rows, positions 0–7. - Deleted positions. 2, 5, 6.
- Bitmap. Bits 2, 5, 6 set; all others clear.
Question. Show the bitmap for the deleted set, and demonstrate that adding another delete updates the same vector rather than creating a new file.
Input.
| Position | Row key | State |
|---|---|---|
| 0 | id=100 | live |
| 2 | id=102 | deleted |
| 5 | id=105 | deleted |
| 6 | id=106 | deleted |
| 7 | id=107 | live |
Code.
# Illustrative Roaring-bitmap model of a deletion vector for one data file
from pyroaring import BitMap
# Deletion vector for data-0002.parquet: positions 2, 5, 6 are deleted
dv = BitMap([2, 5, 6])
print("cardinality:", len(dv)) # 3 deleted rows
print("is position 5 deleted?", 5 in dv) # True
print("is position 0 deleted?", 0 in dv) # False
# A later commit deletes id=107 (position 7): UPDATE THE SAME VECTOR
dv.add(7)
print("after new delete:", sorted(dv)) # [2, 5, 6, 7] — still one vector
# Serialized size stays tiny even as the file grows to millions of rows
print("serialized bytes:", len(dv.serialize()))
Step-by-step explanation.
- The deletion vector for
data-0002.parquetis the integer set{2, 5, 6}. A Roaring bitmap stores this in a few bytes; a positional delete file would have stored three(path, position)rows plus Avro framing overhead. - Membership is O(1): "is position 5 deleted?" is a single bitmap lookup. The reader iterates data-file rows 0..7 and skips any position present in the vector.
- When a later commit deletes position 7, v3 rewrites the same logical vector as
{2, 5, 6, 7}— one new Puffin blob replacing the old one for that data file. Crucially there is still exactly one vector per data file; no swarm accumulates. - Roaring's encoding means the serialized size grows sublinearly. A data file with a million rows and ten thousand scattered deletes still serializes to a small blob, and a file with a contiguous deleted range compresses to almost nothing.
- Because the vector is addressed by data-file path, a reader that opens
data-0002.parquetfetches exactly one vector, not "all delete files whose predicates might touch this file." That is the read-amplification win.
Output.
| Operation | Vector contents | Files on disk |
|---|---|---|
| Initial deletes (2,5,6) | {2,5,6} |
1 Puffin blob |
| Delete position 7 | {2,5,6,7} |
1 Puffin blob (rewritten) |
Read data-0002
|
skip {2,5,6,7} | 1 data + 1 vector |
Rule of thumb. Think of a deletion vector as "the delete file, but exactly one per data file and encoded as a bitmap." Every mental model you had for positional deletes carries over; you just replace "a growing pile of small files" with "one compact vector that gets rewritten in place."
Worked example — the merge-on-read scan path with a deletion vector
Detailed explanation. The point of deletion vectors is the read path. Walk through how a query engine plans and executes a SELECT against a table whose hot data file has a deletion vector, and contrast the work with the v2 positional-delete path.
-
Query.
SELECT count(*) FROM orders WHERE status = 'pending'. - Planning. The planner lists data files, and for each, the single deletion vector (if any) whose sequence number is ≥ the data file's.
- Execution. For each data file, scan rows, skip positions in the vector, apply the filter, count.
Question. Describe the scan of one data file with a deletion vector and compute the I/O compared with the v2 positional-delete approach.
Input.
| Item | v3 deletion vectors | v2 positional deletes |
|---|---|---|
| Delete markers for the file | 1 bitmap | up to N small files |
| Files opened per data file | 1 vector | N delete files |
| Position lookup | O(1) bitmap | merge of N sorted lists |
| Compaction burden | rewrite 1 vector | merge N files |
Code.
# Reader-side pseudocode for one data file under merge-on-read (v3)
def scan_data_file(data_file, deletion_vector, row_filter):
"""Yield live rows of a data file, skipping deleted positions."""
dv = deletion_vector # a Roaring bitmap, or None if the file has no deletes
for position, row in enumerate(read_parquet_rows(data_file)):
if dv is not None and position in dv:
continue # O(1) skip: this position is deleted
if row_filter(row): # WHERE status = 'pending'
yield row
# Planner picks, for each data file, the ONE deletion vector with seq >= data seq
def plan(snapshot):
for data_file in snapshot.data_files():
dv = snapshot.deletion_vector_for(data_file) # at most one
yield (data_file, dv)
Step-by-step explanation.
- The planner pairs each data file with its single deletion vector (or
None). There is no N-way match of "which delete files might touch this data file" — the vector is addressed directly by data-file path, so planning is O(files), not O(files × deletes). - Execution reads the data file's rows in position order. For each position, an O(1) bitmap membership check decides whether to skip it. There is no sort-merge of delete lists as v2 required.
- Only after the deletion skip does the engine apply the
WHEREfilter, so deleted rows never reach predicate evaluation, projection, or aggregation. The count reflects only live, matching rows. - Compaction cost drops correspondingly. To garbage-collect deletes,
rewrite_data_filesreads the data file, drops the vector's positions, and writes a clean data file with no vector — a single-vector read instead of an N-file merge. - The net effect on a streaming-delete workload is dramatic: v2 could accumulate thousands of positional delete files that every reader had to consider; v3 keeps one bitmap per data file, so read latency stays flat as delete volume grows.
Output.
| Metric | v3 deletion vectors | v2 positional deletes |
|---|---|---|
| Delete files opened per data file | 1 | up to thousands |
| Position check | O(1) bitmap | O(log) per delete list |
| Read latency vs delete volume | flat | grows |
| Compaction input | 1 vector | many small files |
Rule of thumb. If your Iceberg table takes frequent row-level deletes or MERGE updates, deletion vectors are the reason to move to v3 first. The read path stops degrading as deletes pile up, and compaction stops fighting a small-file swarm.
Worked example — MERGE with merge-on-read producing a deletion vector
Detailed explanation. A MERGE INTO upsert under merge-on-read produces two effects: new data files for inserted/updated rows, and deletion-vector updates that retire the old versions of updated rows. Walk through a CDC-style merge that applies 3 updates and 1 delete to orders.
-
Target.
orders(v3,write.merge.mode = merge-on-read). - Source. A change set: update orders 101, 102, 103; delete order 104.
- Result. New data file with new versions of 101–103; deletion vector marks old positions of 101–104.
Question. Run the merge and show which files and vectors the commit produces.
Input.
| Change | Order id | Effect on data | Effect on vector |
|---|---|---|---|
| update | 101 | new row appended | old position marked deleted |
| update | 102 | new row appended | old position marked deleted |
| update | 103 | new row appended | old position marked deleted |
| delete | 104 | none | old position marked deleted |
Code.
-- Ensure merge-on-read for this table (v3 uses deletion vectors for deletes)
ALTER TABLE prod.db.orders SET TBLPROPERTIES (
'format-version' = '3',
'write.merge.mode' = 'merge-on-read',
'write.update.mode' = 'merge-on-read',
'write.delete.mode' = 'merge-on-read'
);
-- CDC upsert: update 101-103, delete 104
MERGE INTO prod.db.orders AS t
USING staging.orders_changes AS s
ON t.id = s.id
WHEN MATCHED AND s.op = 'D' THEN DELETE
WHEN MATCHED AND s.op = 'U' THEN UPDATE SET
status = s.status, total_cents = s.total_cents
WHEN NOT MATCHED AND s.op = 'I' THEN INSERT (id, customer_id, total_cents, status)
VALUES (s.id, s.customer_id, s.total_cents, s.status);
Step-by-step explanation.
- The engine plans the merge against the current snapshot. Rows 101–104 live in an existing data file, say
data-0002.parquet, at known positions. Row 104 is a pure delete; 101–103 are updates (delete-old + insert-new). - For the three updates, the engine writes the new row versions into a fresh data file
data-0007.parquet— inserts are always new files under MoR. - For all four affected rows, the engine marks the old positions in
data-0002.parquet's deletion vector. If the file had no vector, one is created with{pos(101), pos(102), pos(103), pos(104)}; if it had one, those positions are added. - The commit produces a new snapshot referencing
data-0007.parquet(new versions) and the updated deletion vector fordata-0002.parquet(retiring old versions). Sequence numbers ensure the vector, committed at seq K, applies todata-0002which was committed at seq < K. - A subsequent reader sees the new versions of 101–103 (from
data-0007), does not see 104 at all (deleted, no new version), and skips the old versions of 101–104 indata-0002via the vector. The upsert is correct with no data-file rewrite of the 500 MB original.
Output.
| Order id | Old version (data-0002) | New version (data-0007) | Visible? |
|---|---|---|---|
| 101 | skipped by vector | present | new version |
| 102 | skipped by vector | present | new version |
| 103 | skipped by vector | present | new version |
| 104 | skipped by vector | none | not visible |
Rule of thumb. Under merge-on-read, an UPDATE is always "mark old position deleted + append new row." Deletion vectors make the "mark old position deleted" half nearly free, which is what makes high-frequency MERGE workloads practical on Iceberg v3.
Data engineering interview question on deletion vectors
A senior interviewer might ask: "Your streaming CDC pipeline into an Iceberg v2 table deletes and updates a few hundred rows per minute, and read latency has been climbing for weeks. You suspect the delete files. Explain what is happening, how Iceberg v3 deletion vectors fix it, what changes at write time and read time, and what maintenance job you still need."
Solution Using deletion vectors, sequence-number ordering, and periodic compaction
-- 1. Upgrade the table and switch row-level ops to merge-on-read
ALTER TABLE prod.db.orders SET TBLPROPERTIES (
'format-version' = '3',
'write.delete.mode' = 'merge-on-read',
'write.update.mode' = 'merge-on-read',
'write.merge.mode' = 'merge-on-read'
);
-- 2. Ongoing CDC merges now produce one deletion vector per touched data file
-- (no positional-delete-file swarm)
-- 3. Scheduled maintenance: compact data files AND rewrite/expire deletes
CALL prod.system.rewrite_data_files(
table => 'db.orders',
options => map('delete-file-threshold', '1') -- rewrite files carrying deletes
);
CALL prod.system.rewrite_position_delete_files(table => 'db.orders'); -- legacy cleanup
CALL prod.system.expire_snapshots(table => 'db.orders', older_than => TIMESTAMP '2026-08-29 00:00:00');
# 4. Read-latency check: count deletes per data file before/after
from pyiceberg.catalog import load_catalog
tbl = load_catalog("prod").load_table("db.orders")
files = tbl.inspect.files() # includes content=0 data, content=1 deletes/DVs
print(files.group_by("file_path").aggregate([("record_count", "sum")]))
Step-by-step trace.
| Step | v2 state (before) | v3 state (after) |
|---|---|---|
| Delete representation | many positional delete files | one deletion vector per data file |
| Files opened per hot data file | dozens–thousands | exactly one vector |
| Read latency trend | climbing with delete volume | flat |
| Write cost per delete | new small file | update one bitmap |
| Compaction input | merge many delete files | rewrite one vector |
| Sequence-number rule | unchanged | unchanged (data.seq ≤ delete.seq) |
After the upgrade, each CDC merge updates a single deletion vector per affected data file instead of dropping a new positional-delete file, so the swarm stops growing; read latency flattens because every data file consults at most one bitmap; and the scheduled rewrite_data_files job periodically materializes the deletes into clean data files so vectors never grow unbounded.
Output:
| Metric | Before (v2) | After (v3 + maintenance) |
|---|---|---|
| Delete files per hot data file | up to thousands | 1 vector |
| p95 read latency | rising weekly | stable |
| Write amplification per delete | 1 small file | in-place bitmap update |
| Compaction frequency needed | constant | periodic |
| Delete correctness | seq-number rule | seq-number rule |
Why this works — concept by concept:
- Deletion vector per data file — replacing a growing set of positional delete files with exactly one Roaring bitmap per data file turns delete lookup from an N-file merge into an O(1) bitmap check, so read latency no longer degrades as delete volume rises.
-
Merge-on-read modes — setting
write.{delete,update,merge}.mode = merge-on-readtells the engine to mark deletes in vectors and append new rows, keeping writes cheap on a high-frequency CDC workload instead of rewriting large data files copy-on-write. - Sequence-number ordering — a deletion vector committed at sequence K applies to any data file committed at sequence ≤ K, which is what lets updates (delete-old + insert-new) stay correct across many concurrent commits without timestamps.
-
Periodic compaction —
rewrite_data_filesmaterializes vector deletes into clean data files andexpire_snapshotsdrops the old ones, bounding both vector size and metadata growth; MoR trades a little scheduled maintenance for cheap writes. - Cost — write is O(changed rows) with an in-place bitmap update; read is O(scanned rows) with an O(1) skip per row; maintenance is O(files carrying deletes) on a schedule. Compared with v2's unbounded delete-file swarm, v3 keeps every axis bounded.
Data
Topic — data-processing
Merge-on-read and CDC-upsert problems
3. Row lineage — row IDs and sequence numbers
Every row gets a stable _row_id and a _last_updated_sequence_number, so you can follow a single row across snapshots
The mental model in one line: row lineage is the v3 feature that assigns every row a table-unique, immutable _row_id and a _last_updated_sequence_number recording the commit that last changed it — the _row_id is derived from a data file's first-row-id plus the row's position, and the pair together lets an engine follow one logical row as it is inserted, updated, and deleted across snapshots, which is the primitive that makes efficient incremental processing, change queries, and materialized-view maintenance possible without a separate CDC side channel. Row lineage is enabled by default on v3 tables and is the feature most likely to change how you build incremental pipelines on a lakehouse.
The two lineage fields.
-
_row_id. A table-scoped, monotonically assigned integer that identifies a logical row for the life of the table. It is assigned once, when the row is first written, and it travels with the row through every update. Two different rows never share a_row_id; a single row keeps its_row_ideven after it moves to a new data file during an update or compaction. -
_last_updated_sequence_number. The sequence number of the commit that most recently modified the row. A freshly inserted row's value equals the sequence number of its insert; an updated row's value advances to the update's sequence number. This is what lets "give me every row changed since sequence N" be answered by a metadata-aware scan.
How _row_id is assigned — first-row-id + position.
-
Files carry a
first-row-id. When a data file is committed, the table's running row-id counter reserves a contiguous block for it; the file records itsfirst-row-id. Row k in that file has_row_id = first_row_id + kunless it inherited an id from a previous version. -
Inheritance on rewrite. When compaction or an update rewrites a row into a new file, the row keeps its original
_row_id— the new file records the inherited ids explicitly rather than deriving them from position. This is what makes the id stable across physical rewrites. - The counter is durable. The next-row-id counter lives in table metadata and only ever advances, so ids are never reused even across snapshot expiry.
What lineage unlocks.
-
Incremental processing. A downstream job records "I have consumed up to sequence number N." On the next run it scans only rows with
_last_updated_sequence_number > N, using manifest-level sequence bounds to skip whole files. No full-table diff, no external watermark table. -
Change queries / CDC. Because a row keeps its
_row_idacross updates, you can emit "row 101 changed from X to Y" by joining the old and new versions on_row_id, rather than guessing identity from a business key. -
Materialized-view maintenance. An incremental view refresh consumes only changed rows since its last refresh sequence number and applies them by
_row_id, turning a full recompute into a delta apply.
Common interview probes on row lineage.
- "What are the two row-lineage fields?" — required answer:
_row_id(stable identity) and_last_updated_sequence_number(last-modified commit). - "How is
_row_idassigned?" —first-row-idof the data file plus row position, inherited on rewrite. - "Why is lineage better than a business key for CDC?" — it is engine-maintained, immune to key changes, and lets you skip files by sequence-number bounds.
- "Is lineage on by default in v3?" — yes; it is a table-level property enabled for v3 tables.
Worked example — assigning row IDs from first-row-id
Detailed explanation. Walk through how ids land on rows across two data files, so the first-row-id + position rule is concrete.
-
File A.
first-row-id = 0, 4 rows. Rows get_row_id0,1,2,3. -
File B.
first-row-id = 4, 3 rows. Rows get_row_id4,5,6. - Next-row-id counter. Advanced to 7 after both commits.
Question. Show the _row_id assigned to each row and the counter's value after each commit.
Input.
| File | first-row-id | rows | positions |
|---|---|---|---|
| data-A | 0 | 4 | 0,1,2,3 |
| data-B | 4 | 3 | 0,1,2 |
Code.
# Model of first-row-id assignment (what the writer does at commit time)
def assign_row_ids(files):
next_row_id = 0
assigned = {}
for name, row_count in files:
first_row_id = next_row_id # reserve a contiguous block
assigned[name] = [first_row_id + pos for pos in range(row_count)]
next_row_id += row_count # advance the durable counter
return assigned, next_row_id
files = [("data-A", 4), ("data-B", 3)]
ids, counter = assign_row_ids(files)
print(ids) # {'data-A': [0,1,2,3], 'data-B': [4,5,6]}
print(counter) # 7 -> stored in table metadata as next-row-id
Step-by-step explanation.
- The writer reserves a contiguous id block per data file.
data-Areserves[0,4); its rows get ids 0–3 by adding position tofirst-row-id = 0. -
data-Breserves the next block[4,7); its rows get 4–6. The blocks never overlap because the counter only advances. - The final counter value 7 is persisted as
next-row-idin table metadata. The next data file to be committed will start atfirst-row-id = 7. - Because the id derives from
first-row-id + position, a reader can compute ids on the fly for files that were never rewritten — the ids need not be materialized column-by-column until a rewrite forces inheritance. - This scheme is O(1) per file at write time (reserve a block) and needs no coordination beyond the single atomic commit that advances the counter, so it scales with concurrent writers exactly like every other Iceberg commit.
Output.
| Row | File | position | _row_id |
|---|---|---|---|
| id=100 | data-A | 0 | 0 |
| id=103 | data-A | 3 | 3 |
| id=104 | data-B | 0 | 4 |
| id=106 | data-B | 2 | 6 |
Rule of thumb. Read _row_id as "first-row-id of my file plus my position, unless I was rewritten, in which case my id was carried over explicitly." That one sentence covers the entire assignment model.
Worked example — following a row across snapshots
Detailed explanation. Lineage's payoff is tracking one row through change. Follow _row_id = 101 as it is inserted at sequence 1, updated at sequence 2, and deleted at sequence 3.
-
Seq 1 (insert). Row 101 written to
data-A;_last_updated_sequence_number = 1. -
Seq 2 (update). Row 101's new version written to
data-C, keeping_row_id = 101;_last_updated_sequence_number = 2; old position indata-Amarked in a deletion vector. -
Seq 3 (delete). Row 101's position in
data-Cmarked deleted; the logical row is gone.
Question. Show the lineage-field values for row 101 at each snapshot and how a change query reconstructs its history.
Input.
| Seq | Operation | Data file | _row_id |
_last_updated_seq |
|---|---|---|---|---|
| 1 | insert | data-A | 101 | 1 |
| 2 | update | data-C | 101 | 2 |
| 3 | delete | (vector on data-C) | 101 | 3 |
Code.
-- Change query: reconstruct the history of a single logical row by _row_id
-- Compare two snapshots (start_snapshot -> end_snapshot) via Iceberg changelog
CALL prod.system.create_changelog_view(
table => 'db.orders',
options => map('start-snapshot-id','30001','end-snapshot-id','30003'),
changelog_view => 'orders_changes'
);
SELECT _row_id, _change_type, _change_ordinal, status, total_cents
FROM orders_changes
WHERE _row_id = 101
ORDER BY _change_ordinal;
Step-by-step explanation.
- At sequence 1, row 101 is inserted with
_row_id = 101and_last_updated_sequence_number = 1. Any incremental consumer with watermark < 1 will pick it up. - At sequence 2, the update writes a new physical version of row 101 into a different data file but preserves
_row_id = 101. Its_last_updated_sequence_numberadvances to 2, and the old version is retired via a deletion vector. A consumer watermarked at 1 sees exactly one changed row (101) with the new values. - At sequence 3, the delete marks row 101's current position, and the change record for it is a delete with
_row_id = 101. Downstream systems keyed on_row_idretire their copy. - The changelog view joins the old and new versions by
_row_id, so the emitted change is precise: it says "row 101 went insert → update → delete" without ever consulting the business keyid. If the application had reusedid = 101for a different logical entity, lineage would still keep them distinct because each got its own_row_id. - This is why lineage beats business keys for CDC: identity is engine-maintained and immutable, so key churn, natural-key collisions, and re-inserts never confuse the change stream.
Output.
_change_ordinal |
_change_type |
_row_id |
status |
|---|---|---|---|
| 0 | INSERT | 101 | pending |
| 1 | UPDATE_AFTER | 101 | shipped |
| 2 | DELETE | 101 | (gone) |
Rule of thumb. When you need "what happened to this exact row over time," join on _row_id, not the business key. The _row_id survives updates, rewrites, and compaction; the business key does not survive re-use.
Worked example — incremental scan by sequence number
Detailed explanation. The most common lineage use is an incremental read: "process only rows changed since my last run." The _last_updated_sequence_number plus manifest-level sequence bounds make this a file-skipping operation.
- Watermark. Last processed sequence number = 4.
- New commits. Sequences 5 and 6 changed rows in two data files.
- Scan. Skip every manifest and file whose max sequence ≤ 4; read only changed rows.
Question. Implement an incremental consumer that reads only rows with _last_updated_sequence_number > 4.
Input.
| File | max seq in file | contains changed rows? |
|---|---|---|
| data-A | 3 | no (skip whole file) |
| data-D | 5 | yes |
| data-E | 6 | yes |
Code.
# Incremental consumer using Iceberg incremental scan (append + changelog)
from pyiceberg.catalog import load_catalog
tbl = load_catalog("prod").load_table("db.orders")
LAST_SEQ = 4 # durable watermark from the previous run
# Resolve snapshots strictly after our watermark and scan only their changes.
snaps = [s for s in tbl.snapshots() if s.sequence_number > LAST_SEQ]
start = min(s.snapshot_id for s in snaps)
end = max(s.snapshot_id for s in snaps)
scan = tbl.incremental_append_scan( # append/changelog incremental scan
from_snapshot_id_exclusive=None,
to_snapshot_id=end,
)
for batch in scan.to_arrow_batches():
process(batch) # only rows changed after LAST_SEQ
# Advance the watermark to the highest sequence number consumed
LAST_SEQ = max(s.sequence_number for s in snaps)
Step-by-step explanation.
- The consumer keeps a durable watermark — the highest
_last_updated_sequence_numberit has processed (here 4). It never needs an external CDC table; the watermark is a single integer. - Planning filters manifests by their sequence-number bounds.
data-A's max sequence is 3, so the entire manifest entry is skipped without opening the file — this is the file-skipping win that makes incremental scans cheap on huge tables. - Only
data-D(seq 5) anddata-E(seq 6) are opened, and within them only rows whose_last_updated_sequence_number > 4are emitted. The bulk of the table is never touched. - After the run, the watermark advances to 6. The next run repeats the same skip logic from 6, so steady-state cost is proportional to the change volume, not the table size.
- Because ids are stable, the downstream apply is idempotent: reprocessing a sequence range re-emits the same
_row_ids with the same values, so a consumer that dedupes by_row_idis safe to replay.
Output.
| File | Opened? | Rows emitted |
|---|---|---|
| data-A (seq 3) | no | 0 |
| data-D (seq 5) | yes | changed rows |
| data-E (seq 6) | yes | changed rows |
Rule of thumb. Incremental processing on v3 is "keep one integer watermark, scan sequences greater than it, skip files by their sequence bounds." Row lineage turns what used to need a bespoke CDC pipeline into a property of the table.
Data engineering interview question on row lineage
A senior interviewer might ask: "You maintain a downstream aggregate that currently full-refreshes off an Iceberg table every hour and it is getting too expensive. The table is on v3 with row lineage. Design an incremental refresh: how you track progress, how you scan only changed rows, how you apply inserts, updates, and deletes to the aggregate, and how you make replay idempotent."
Solution Using _row_id, _last_updated_sequence_number, and a sequence-number watermark
-- 1. Progress table: one durable watermark per consumer
CREATE TABLE IF NOT EXISTS meta.consumer_watermarks (
consumer_name STRING,
last_seq BIGINT
) USING iceberg;
-- 2. Build a changelog view over the unprocessed sequence range
CALL prod.system.create_changelog_view(
table => 'db.orders',
options => map('start-snapshot-id', :from_snap, 'end-snapshot-id', :to_snap),
changelog_view => 'orders_delta',
identifier_columns => array('_row_id')
);
-- 3. Apply the delta to the aggregate keyed by _row_id
MERGE INTO agg.orders_by_customer AS a
USING (
SELECT customer_id, _row_id, _change_type, total_cents
FROM orders_delta
) d
ON a.row_id = d._row_id
WHEN MATCHED AND d._change_type = 'DELETE' THEN DELETE
WHEN MATCHED AND d._change_type = 'UPDATE_AFTER' THEN UPDATE SET amount = d.total_cents
WHEN NOT MATCHED AND d._change_type = 'INSERT' THEN INSERT (row_id, customer_id, amount)
VALUES (d._row_id, d.customer_id, d.total_cents);
# 4. Driver: read watermark, resolve snapshot range, run the merge, advance watermark
from pyiceberg.catalog import load_catalog
tbl = load_catalog("prod").load_table("db.orders")
last_seq = read_watermark("orders_agg") # e.g. 4
snaps = [s for s in tbl.snapshots() if s.sequence_number > last_seq]
if snaps:
run_merge(from_snap=min(s.snapshot_id for s in snaps),
to_snap=max(s.snapshot_id for s in snaps))
write_watermark("orders_agg", max(s.sequence_number for s in snaps))
Step-by-step trace.
| Step | Mechanism | Result |
|---|---|---|
| Track progress | consumer_watermarks.last_seq |
one integer per consumer |
| Find changes | snapshots with sequence_number > last_seq
|
skip already-consumed commits |
| Scan | changelog view keyed by _row_id
|
only changed rows |
| Apply | MERGE by _row_id
|
insert/update/delete the aggregate |
| Idempotency | dedupe/key on _row_id
|
safe replay of a range |
| Advance | set last_seq to max consumed |
next run starts here |
After deployment, the hourly refresh reads only the rows changed in the last hour (identified by _last_updated_sequence_number and skipped at file granularity by sequence bounds), applies them to the aggregate by the stable _row_id, and advances a single-integer watermark; a re-run over the same range produces the same result because every change is keyed by an immutable id.
Output:
| Metric | Full refresh (before) | Incremental (after) |
|---|---|---|
| Rows scanned per run | whole table | only changed rows |
| Files opened | all | those with seq > watermark |
| Progress state | none / external | one integer per consumer |
| Delete handling | recompute | applied by _row_id
|
| Replay safety | recompute is idempotent | keyed dedupe is idempotent |
Why this works — concept by concept:
-
_row_idas stable identity — because the id survives updates, rewrites, and compaction, the downstream MERGE keys on it directly, so inserts, updates, and deletes apply to exactly the right aggregate row without relying on a business key that might churn. -
_last_updated_sequence_number— recording the commit that last touched a row turns "what changed since last run" into a numeric comparison, and manifest-level sequence bounds let the planner skip whole files whose max sequence is below the watermark. - Sequence-number watermark — a single durable integer per consumer replaces an external CDC table; advancing it after a successful apply gives exactly-once-effect processing when combined with keyed idempotency.
-
Changelog view — Iceberg materializes insert/update/delete change types over a snapshot range with
_row_idas the identifier, so the consumer receives a clean delta instead of diffing two full table states. - Cost — steady-state work is O(changed rows) for the scan plus O(changed rows) for the keyed MERGE, versus O(table) for a full refresh. Row lineage moves incremental processing from application code into the table format.
Data
Topic — data-processing
Incremental-processing and change-tracking problems
4. New types — variant, geometry, and binary
v3 adds binary-encoded semi-structured and geospatial types — variant, geometry, geography, nanosecond timestamps — plus default values and richer spec transforms
The mental model in one line: Iceberg v3 extends the type system with binary types that carry rich values without a side table — variant for binary-encoded semi-structured (JSON-like) data, geometry and geography for spatial values encoded as well-known binary, and nanosecond-precision timestamps (timestamp_ns / timestamptz_ns) — and it adds spec-level conveniences like default column values, an unknown type for forward compatibility, and multi-argument partition transforms, all while keeping the same four-layer table format tree and the same schema-evolution rules that never rewrite data. These additions are what let a lakehouse absorb event payloads, IoT/spatial data, and high-resolution timestamps natively instead of stuffing them into strings.
The new value types.
-
Variant. A single column type that stores semi-structured data (objects, arrays, scalars) in a compact binary encoding rather than as a text JSON string. Engines can extract fields (
variant_get) and, with shredding, store frequently-accessed sub-fields as typed columns for pushdown while keeping the full document. Variant is the answer to "we have event JSON and we do not want to pre-flatten it." -
Geometry and geography. Spatial types encoded as well-known binary (WKB).
geometryuses planar (Cartesian) coordinates;geographyuses spherical (lat/long on an ellipsoid) semantics for distance and containment. Both carry a coordinate reference system so spatial predicates are well-defined. -
Nanosecond timestamps.
timestamp_nsandtimestamptz_nsadd nanosecond precision alongside the existing microsecondtimestamp/timestamptz, for high-frequency and scientific workloads where microseconds lose information.
The spec-level conveniences.
-
Default column values. v3 lets a column declare an
initial-default(value for existing rows when the column is added) and awrite-default(value written when a writer omits the column). Adding a non-null column with a default is now a metadata-only change — no rewrite to backfill. -
The
unknowntype. A forward-compatibility placeholder: a reader that encounters a column of a type it does not understand can treat it asunknown(all nulls) rather than failing, which smooths spec evolution. - Multi-argument partition transforms. v3 generalizes partition transforms to take multiple arguments, enabling richer hidden-partitioning schemes than the single-column transforms of v2.
Schema evolution and type promotion — still no rewrite.
- Add / drop / rename / reorder columns are metadata-only, tracked by column field ids so data files never need touching.
-
Type promotion follows safe-widening rules (for example
int→long,float→double, and decimal precision increases). The new types slot into these rules; you cannot narrow a type in a way that would lose data. - Field ids, not names, bind schema to data. Renaming a column changes the name in metadata while the data files keep referencing the stable field id, which is why rename is free.
Common interview probes on the new types.
- "What is the variant type for?" — required answer: binary-encoded semi-structured data (JSON-like) stored natively, with optional shredding for pushdown.
- "Geometry vs geography?" — planar vs spherical coordinate semantics; both WKB-encoded with a CRS.
- "How can adding a NOT NULL column avoid a rewrite?" — v3 default values (
initial-defaultbackfills existing rows logically). - "Why do renames not rewrite data?" — schema binds by field id, not name.
Worked example — storing and querying variant event data
Detailed explanation. A clickstream table receives heterogeneous event payloads. Instead of a brittle wide schema or an opaque JSON string, store the payload as variant and extract fields at query time, optionally shredding hot fields.
-
Column.
payload variant. - Write. Parse JSON into the binary variant encoding at ingest.
-
Read. Extract
payload.event_typeandpayload.valuewith variant accessors.
Question. Create the table, ingest two differently-shaped events, and query a field that only some events have.
Input.
| event | payload JSON |
|---|---|
| A | {"event_type":"click","target":"buy","value":1} |
| B | {"event_type":"scroll","depth":80} |
Code.
-- 1. Table with a variant column (v3)
CREATE TABLE prod.db.clicks (
event_id BIGINT,
ts TIMESTAMP,
payload VARIANT
) USING iceberg TBLPROPERTIES ('format-version' = '3');
-- 2. Ingest heterogeneous events (JSON parsed into the binary variant encoding)
INSERT INTO prod.db.clicks VALUES
(1, TIMESTAMP '2026-09-05 10:00:00', parse_json('{"event_type":"click","target":"buy","value":1}')),
(2, TIMESTAMP '2026-09-05 10:00:01', parse_json('{"event_type":"scroll","depth":80}'));
-- 3. Extract fields; rows missing a field return NULL, not an error
SELECT event_id,
variant_get(payload, '$.event_type', 'string') AS event_type,
variant_get(payload, '$.value', 'int') AS value,
variant_get(payload, '$.depth', 'int') AS depth
FROM prod.db.clicks
ORDER BY event_id;
Step-by-step explanation.
- The
payload variantcolumn stores each event's structure in a self-describing binary encoding. Unlike a fixed schema, events A and B can have entirely different shapes in the same column. -
parse_jsonconverts the incoming text into the variant binary at write time, so reads never re-parse text — extraction walks the binary structure directly. -
variant_get(payload, '$.value', 'int')navigates the document by path and coerces to the requested type. Event B has novalue, so it returns NULL rather than failing — variant access is null-tolerant across heterogeneous rows. - With variant shredding, an engine can additionally materialize hot sub-fields (say
event_type) as a typed sub-column so filters likeWHERE event_type = 'click'push down to file statistics, while the full document remains available for rare fields. - Because the type is native, statistics, encoding, and predicate pushdown work far better than storing the same JSON as a
string, and downstream engines that understand variant get typed access instead of string parsing.
Output.
| event_id | event_type | value | depth |
|---|---|---|---|
| 1 | click | 1 | NULL |
| 2 | scroll | NULL | 80 |
Rule of thumb. Reach for variant when payload shape is heterogeneous or evolving and you still want typed, pushdown-friendly access. Store JSON as string only when you truly never query into it.
Worked example — geometry partitioning and a spatial predicate
Detailed explanation. A rides table stores pickup points. Store them as geometry (WKB) so spatial predicates and stats work, and hidden-partition by a spatial bucket for pruning.
-
Column.
pickup geometry. - Predicate. Rides within a bounding box.
- Pruning. Manifest bounds on the geometry let the planner skip files outside the box.
Question. Create the spatial table and query rides inside a bounding box, showing that non-overlapping files are pruned.
Input.
| ride_id | pickup (WKT) |
|---|---|
| 1 | POINT(-73.98 40.75) |
| 2 | POINT(-118.24 34.05) |
Code.
-- 1. Spatial table with a geometry column (planar coords, WKB-encoded)
CREATE TABLE prod.db.rides (
ride_id BIGINT,
pickup GEOMETRY,
fare DECIMAL(8,2)
) USING iceberg TBLPROPERTIES ('format-version' = '3');
INSERT INTO prod.db.rides VALUES
(1, ST_GeomFromText('POINT(-73.98 40.75)'), 24.50), -- NYC
(2, ST_GeomFromText('POINT(-118.24 34.05)'), 31.00); -- LA
-- 2. Rides inside a Manhattan bounding box; files outside are pruned by bounds
SELECT ride_id, fare
FROM prod.db.rides
WHERE ST_Intersects(
pickup,
ST_GeomFromText('POLYGON((-74.02 40.70,-73.93 40.70,-73.93 40.80,-74.02 40.80,-74.02 40.70))')
);
Step-by-step explanation.
-
pickup geometrystores each point as WKB with a coordinate reference system, so the value is a true spatial object, not two float columns the engine cannot reason about spatially. - On write, Iceberg records spatial lower/upper bounds (a bounding box) per data file in the manifest, just as it records min/max for scalars.
- The
ST_Intersectspredicate against the Manhattan polygon lets the planner compare each file's bounding box to the query box. The LA ride sits in a file whose bounds do not overlap Manhattan, so that file is pruned without being read. -
geometryuses planar math (fast, correct for local/projected coordinates); had we needed great-circle distances across the globe we would have chosengeographyfor spherical semantics. - The result is a spatial query that both returns correct geometry-aware results and prunes irrelevant files — the lakehouse equivalent of a spatial index at the file level.
Output.
| ride_id | fare | file pruned? |
|---|---|---|
| 1 | 24.50 | no (in box) |
| 2 | — | yes (LA file skipped) |
Rule of thumb. Use geometry for projected/planar coordinates and local analysis; use geography when distance and containment must be correct on the globe. Either way, store spatial data as the native type so file-level bounds can prune.
Worked example — adding a NOT NULL column with a default, no rewrite
Detailed explanation. You must add a currency column, defaulting existing rows to 'USD', to a billion-row table. In v3 this is metadata-only thanks to default values.
-
initial-default. Logical value for rows that predate the column. -
write-default. Value written when a writer omits the column. - No rewrite. Existing data files are untouched; reads synthesize the default.
Question. Add the defaulted column and show that old rows read back the default while the data files are never rewritten.
Input.
| Aspect | Value |
|---|---|
| New column | currency STRING NOT NULL |
| initial-default | 'USD' |
| write-default | 'USD' |
| Data rewrite | none |
Code.
-- Add a NOT NULL column with defaults — metadata-only in v3
ALTER TABLE prod.db.orders
ADD COLUMN currency STRING DEFAULT 'USD'; -- sets write-default (and initial-default)
-- Old rows (written before the column existed) read back the initial-default
SELECT id, total_cents, currency
FROM prod.db.orders
WHERE id IN (1, 2) -- rows that predate the ADD COLUMN
LIMIT 2;
-- New writes may omit currency; the write-default fills it
INSERT INTO prod.db.orders (id, customer_id, total_cents, status)
VALUES (9001, 7, 500, 'pending');
Step-by-step explanation.
-
ADD COLUMN currency STRING DEFAULT 'USD'records the column in the schema withinitial-defaultandwrite-defaultboth'USD'. It commits as a metadata-only change — no data file is read or rewritten. - Reading an old data file, the engine sees no
currencyfield (by field id) and synthesizes theinitial-default'USD'for those rows, so the NOT NULL contract holds logically without a physical backfill. - A new INSERT that omits
currencygets thewrite-default'USD'materialized into the new data file. A new INSERT that suppliescurrencystores the supplied value. - Because binding is by field id, a later rename of
currencywould again be metadata-only; the defaults travel with the field id. - The operation is O(1) regardless of table size — the entire billion-row backfill is expressed as a default in metadata rather than a rewrite job, which pre-v3 would have been a multi-hour
rewrite_data_files.
Output.
| id | total_cents | currency | source |
|---|---|---|---|
| 1 | 1500 | USD | initial-default (old row) |
| 2 | 900 | USD | initial-default (old row) |
| 9001 | 500 | USD | write-default (new row) |
Rule of thumb. In v3, "add a NOT NULL column and backfill a constant" is a metadata change, not a rewrite. Reserve rewrite_data_files for cases where the backfill value is computed per-row, not constant.
Data engineering interview question on the new v3 types
A senior interviewer might ask: "We ingest IoT telemetry: each device sends a JSON blob, a GPS fix, and a nanosecond timestamp. Today it all lands as strings in a v2 table and querying is painful. Redesign the schema on Iceberg v3 using the new types, explain how each type helps pruning or pushdown, and how you would add a defaulted region column to a huge existing table without a rewrite."
Solution Using variant, geometry, nanosecond timestamps, and default values
-- 1. v3 table using the new types instead of strings
CREATE TABLE prod.db.telemetry (
device_id BIGINT,
event_ts TIMESTAMP_NS, -- nanosecond precision, not string
location GEOGRAPHY, -- spherical GPS semantics, WKB
payload VARIANT, -- binary semi-structured, not JSON string
region STRING DEFAULT 'unknown' -- default value: metadata-only add
) USING iceberg
PARTITIONED BY (days(event_ts), bucket(32, device_id)) -- multi-field hidden partitioning
TBLPROPERTIES ('format-version' = '3');
-- 2. Query: field extraction + spatial predicate + nanosecond range, all pushdown-friendly
SELECT device_id,
variant_get(payload, '$.metric', 'double') AS metric
FROM prod.db.telemetry
WHERE event_ts BETWEEN TIMESTAMP_NS '2026-09-05 00:00:00.000000000'
AND TIMESTAMP_NS '2026-09-05 01:00:00.000000000'
AND ST_DWithin(location, ST_GeogFromText('POINT(-73.98 40.75)'), 5000) -- within 5km
AND variant_get(payload, '$.metric', 'double') > 100;
-- 3. Add a defaulted column to the existing huge table with NO rewrite
ALTER TABLE prod.db.telemetry ADD COLUMN ingest_source STRING DEFAULT 'kafka';
Step-by-step trace.
| Column | v2 (before) | v3 (after) | Benefit |
|---|---|---|---|
| timestamp | string | TIMESTAMP_NS |
range pruning + full precision |
| GPS | two floats / string | GEOGRAPHY |
spatial predicate + bounds pruning |
| JSON | string | VARIANT |
typed extraction + shredding pushdown |
| region | manual backfill | DEFAULT 'unknown' |
metadata-only add |
| partitioning | single transform |
days(...), bucket(...)
|
multi-field pruning |
After the redesign, the timestamp filter prunes files by nanosecond bounds, the geography predicate prunes by spatial bounds, the variant accessor reads typed fields (and pushes down shredded sub-columns), and the defaulted region and later ingest_source columns are added as pure metadata changes — no multi-hour rewrite on the billion-row table.
Output:
| Concern | Result |
|---|---|
| Timestamp precision | full nanoseconds preserved |
| Spatial query | correct spherical distance + file pruning |
| JSON access | typed variant_get, null-tolerant |
| New column | metadata-only default, no rewrite |
| Partition pruning | by day and device bucket |
Why this works — concept by concept:
-
Variant — storing JSON in Iceberg's binary variant encoding gives typed, null-tolerant field access and enables shredding hot sub-fields into columns for predicate pushdown, which a
stringcolumn can never do. - Geometry / geography — WKB-encoded spatial types carry a coordinate reference system, so predicates are well-defined and Iceberg records spatial bounds per file, turning spatial filters into file-pruning operations.
-
Nanosecond timestamps —
timestamp_ns/timestamptz_nspreserve sub-microsecond resolution and still support min/max bounds, so range queries prune files while keeping full precision that a string would lose. -
Default values —
initial-defaultlets existing rows read a value for a newly added column andwrite-defaultfills omitted writes, so adding even a NOT NULL column is a metadata change instead of a rewrite. - Cost — the type upgrades are write-time encoding choices with read-time pushdown payoffs; the defaulted-column add is O(1) metadata. Field-id binding keeps every add/rename/reorder free of data rewrites.
SQL
Topic — database
Schema-design and data-type modeling problems
5. Migration, compatibility, and engine support
Upgrading to v3 is a metadata format-version bump — but every reader and writer must understand v3 before you flip it
The mental model in one line: migrating a table to apache iceberg v3 is a metadata-only bump of the format-version integer from 2 to 3 that rewrites no data and is reversible in principle, but it is safe only when every engine that reads or writes the table has been upgraded to a v3-capable version — because a v2-only writer that commits to a v3 table, or a v2-only reader that cannot decode deletion vectors or the new types, is a correctness and availability risk — so the real work of migration is the compatibility audit, the phased rollout, and the validation, not the one-line ALTER. Treat the version bump as the last step of a migration, not the first.
What the version bump does and does not do.
-
Does. Sets
format-version = 3in the table metadata, unlocking deletion vectors, row lineage (enabled by default), and the new types for subsequent writes. It is a normal atomic commit — new metadata JSON, catalog CAS. -
Does not. Rewrite any existing data files, convert existing positional delete files to deletion vectors, or assign
_row_ids to already-written rows retroactively. Those happen lazily on the next rewrite/compaction, or via an explicit maintenance job. - Is reversible in principle, risky in practice. You can lower the version only if no v3-only feature has been used; once deletion vectors, lineage, or new-type columns exist, downgrading would strand data an older reader cannot interpret.
The compatibility audit — do this before the bump.
- Enumerate every reader and writer. Spark jobs, Flink jobs, Trino/Presto clusters, the warehouse's external-table reader (Snowflake, BigQuery, Athena, Redshift), dashboards, ad-hoc notebooks, and any home-grown PyIceberg/Java-API service.
- Check the Iceberg library version each uses. v3 features arrived across Iceberg library releases; an engine bundling an older Iceberg runtime will not understand deletion vectors, row lineage, or variant/geo types even if the engine itself is new.
- Classify each as read, write, or both. A writer on an old library is the dangerous case; a read-only consumer on an old library merely fails to see the table until upgraded.
The phased rollout.
- Phase 1 — upgrade engines, still v2. Roll v3-capable runtimes to every reader and writer while the table stays v2. Nothing changes behaviourally; you are just staging capability.
-
Phase 2 — bump one non-critical table. Flip
format-version = 3on a low-stakes table, run every consumer against it, and validate reads/writes, deletes, and any new-type columns. - Phase 3 — enable features deliberately. Turn on merge-on-read deletion vectors where you want cheap deletes, adopt new types where they help, and rely on default-on row lineage for incremental consumers.
- Phase 4 — bump the rest + validate. Roll to remaining tables with a validation query that counts rows, checks a deleted row is absent, and confirms a new-type column round-trips.
Common interview probes on migration.
- "Does upgrading to v3 rewrite data?" — required answer: no, it is a metadata-only version bump; conversions happen lazily.
- "What is the risk of the bump?" — a reader or writer on an old Iceberg library that cannot handle v3 content.
- "Can you downgrade?" — only before any v3-only feature is used.
- "What order do you roll it out?" — upgrade engines first, bump a test table, enable features, then bump the rest.
Worked example — the upgrade command and its metadata effect
Detailed explanation. Walk through flipping a table to v3 and confirming that only metadata changed.
-
Before.
format-version = 2, N data files, M positional delete files. -
Command.
ALTER TABLE ... SET TBLPROPERTIES ('format-version' = '3'). -
After.
format-version = 3, same N data files, same M delete files (not yet converted).
Question. Perform the upgrade and prove no data files were rewritten.
Input.
| Metric | Before | Expected after |
|---|---|---|
| format-version | 2 | 3 |
| data files | N | N (unchanged) |
| snapshots | S | S + 1 (the bump commit) |
| data rewritten | — | none |
Code.
-- 1. Record the pre-upgrade file inventory
SELECT count(*) AS data_files FROM prod.db.orders.files WHERE content = 0;
-- 2. The upgrade — a metadata-only commit
ALTER TABLE prod.db.orders SET TBLPROPERTIES ('format-version' = '3');
-- 3. Confirm the version and that data files are unchanged
SELECT count(*) AS data_files FROM prod.db.orders.files WHERE content = 0;
SELECT snapshot_id, operation, summary
FROM prod.db.orders.snapshots
ORDER BY committed_at DESC
LIMIT 2; -- newest commit is the version bump; no added/deleted data files
Step-by-step explanation.
- The pre-upgrade count of
content = 0files establishes the baseline data-file inventory. - The
ALTER TABLE ... SET TBLPROPERTIEScommit writes a new metadata JSON withformat-version = 3and swaps the catalog pointer. It is an ordinary atomic commit and completes in milliseconds. - The post-upgrade file count is identical — no data file was added, removed, or rewritten. The bump only changed the version integer and, because lineage is default-on for v3, recorded the row-id counter groundwork in metadata.
- The newest snapshot's
summaryshows zero added/deleted data files, proving the commit touched only metadata. Existing positional delete files remain as-is until a compaction converts them to deletion vectors. - From this point, new merge-on-read deletes produce deletion vectors, new rows participate in lineage, and new columns may use the new types — all without disturbing the historical data files.
Output.
| Check | Before | After |
|---|---|---|
| format-version | 2 | 3 |
| data-file count | N | N |
| newest op | (prior) | version-bump commit |
| added/deleted data files | — | 0 |
Rule of thumb. The version bump should be boring: same data files, one extra metadata commit. If a tool tries to rewrite data during the bump, stop — that is a separate compaction you should schedule deliberately, not couple to the upgrade.
Worked example — a compatibility gate before rollout
Detailed explanation. Before bumping any production table, gate the rollout on every engine reporting a v3-capable Iceberg runtime. Build a small check that fails loudly if any consumer is behind.
- Inventory. A registry of consumers and the Iceberg library version each ships.
- Threshold. A minimum library version that supports the v3 features you intend to use.
- Gate. Fail if any writer is below threshold; warn if any read-only consumer is.
Question. Implement the compatibility gate and show it blocking a rollout when a writer is on an old runtime.
Input.
| Consumer | Role | Iceberg lib | v3-capable? |
|---|---|---|---|
| spark-etl | write | 1.9 | yes |
| flink-cdc | write | 1.4 | no |
| trino-bi | read | 1.9 | yes |
| snowflake-ext | read | managed | yes |
Code.
# Compatibility gate: block the v3 bump unless every WRITER is v3-capable
MIN_V3_WRITER = (1, 8) # minimum Iceberg lib version for the v3 features we use
consumers = [
{"name": "spark-etl", "role": "write", "lib": (1, 9)},
{"name": "flink-cdc", "role": "write", "lib": (1, 4)},
{"name": "trino-bi", "role": "read", "lib": (1, 9)},
{"name": "snowflake-ext", "role": "read", "lib": (9, 9)}, # managed, treat as capable
]
blockers = [c for c in consumers if c["role"] == "write" and c["lib"] < MIN_V3_WRITER]
warnings = [c for c in consumers if c["role"] == "read" and c["lib"] < MIN_V3_WRITER]
if blockers:
names = ", ".join(c["name"] for c in blockers)
raise SystemExit(f"BLOCK v3 upgrade: writers below {MIN_V3_WRITER}: {names}")
for w in warnings:
print(f"WARN: read-only consumer {w['name']} may not see v3 tables")
print("gate passed")
Step-by-step explanation.
-
MIN_V3_WRITERencodes the minimum Iceberg library version that supports the v3 features you plan to enable (deletion vectors, lineage, new types). Pin it to the exact features you will use, not just "some v3." - The gate separates writers from readers. A writer below threshold is a hard blocker: it could commit v2-only metadata to a table others expect to be v3, or fail to write deletion vectors correctly.
-
flink-cdcon 1.4 is below threshold and is a writer, so the gate raises and aborts the rollout with a clear message naming the offender. No table is bumped while a stale writer exists. - A read-only consumer below threshold produces a warning, not a block — the worst case is that it cannot read v3 tables until upgraded, which is an availability issue for that consumer, not a corruption risk.
- Managed external readers (Snowflake, BigQuery, Athena) are treated as capable per their documented Iceberg support; you still validate them in phase 2 against a real v3 table before trusting them in production.
Output.
| Consumer | Decision | Reason |
|---|---|---|
| spark-etl | pass | writer, lib ≥ 1.8 |
| flink-cdc | BLOCK | writer, lib 1.4 < 1.8 |
| trino-bi | pass | reader, capable |
| snowflake-ext | pass (validate) | managed reader |
Rule of thumb. Gate the version bump on writers first — a stale writer is the only actor that can corrupt a v3 table. Stale readers just need upgrading; stale writers must be stopped before the flip.
Worked example — validation after the bump
Detailed explanation. After bumping a table, run a validation that proves reads, deletes, and new types behave. Build a check that a deleted row is absent, a round-tripped variant field matches, and row counts reconcile.
- Row count. Matches expectation post-delete.
- Delete visibility. A known-deleted key returns zero rows (deletion vector applied).
- New type round-trip. A variant field written equals the field read back.
Question. Write the post-bump validation and interpret a passing run.
Input.
| Check | Expectation |
|---|---|
| count after delete of id=104 | baseline − 1 |
| select id=104 | 0 rows |
| variant round-trip | written == read |
Code.
-- 1. Row count reconciles after a merge-on-read delete
SELECT count(*) AS live_rows FROM prod.db.orders; -- expect baseline - 1
-- 2. A deleted row is invisible (deletion vector applied on read)
SELECT count(*) AS should_be_zero FROM prod.db.orders WHERE id = 104;
-- 3. Row lineage is present and advancing
SELECT min(_last_updated_sequence_number) AS min_seq,
max(_last_updated_sequence_number) AS max_seq
FROM prod.db.orders;
-- 4. New-type round-trip
SELECT variant_get(payload, '$.event_type', 'string') AS et
FROM prod.db.clicks WHERE event_id = 1; -- expect 'click'
Step-by-step explanation.
- The live-row count must equal the pre-delete baseline minus the deletes applied through merge-on-read; a mismatch means a reader is not honouring deletion vectors.
- Selecting the known-deleted
id = 104must return zero rows, proving the deletion vector is applied at read time by every engine you validate — run this from Spark, Trino, and the warehouse reader. - The lineage columns must be present and populated, with sequence numbers advancing; this confirms row lineage is active on the v3 table and available to incremental consumers.
- Extracting a known variant field must return the written value (
'click'), proving the new type round-trips through write and read on this engine. - Run all four checks from each consumer in your inventory. A check that passes on Spark but fails on the warehouse reader localizes the incompatibility to that one engine before it reaches production dashboards.
Output.
| Check | Result | Verdict |
|---|---|---|
| live_rows | baseline − 1 | pass |
| id=104 rows | 0 | pass |
| lineage seq range | advancing | pass |
| variant round-trip | 'click' | pass |
Rule of thumb. Validate the three v3 behaviours — delete visibility, lineage presence, and new-type round-trip — from every engine in your inventory, not just the one that did the write. Cross-engine validation is the only way to catch a lagging reader before a user does.
Data engineering interview question on migration and engine support
A senior interviewer might ask: "You own 400 Iceberg v2 tables read by Spark, Flink, Trino, and Snowflake external tables, and you want the v3 deletion-vector and row-lineage benefits. Design the migration: how you audit compatibility, in what order you roll it out, how you avoid corrupting a table with a stale writer, how you validate, and what your rollback story is."
Solution Using an engine audit, a phased rollout, a compatibility gate, and cross-engine validation
v3 migration plan (400 tables, 4 engine families)
==================================================
PHASE 0 — inventory
- Registry: {table -> consumers}, {consumer -> role, iceberg_lib_version}
- Classify writers vs readers; pin MIN_V3_WRITER to features used
PHASE 1 — upgrade runtimes (tables stay v2)
- Roll v3-capable Iceberg runtimes to Spark, Flink, Trino
- Confirm Snowflake/BigQuery/Athena managed v3 support level
- Behaviour unchanged; capability staged
PHASE 2 — canary
- Bump ONE low-stakes table to format-version=3
- Run the compatibility gate (block if any WRITER < MIN_V3_WRITER)
- Cross-engine validation: delete visibility, lineage, new-type round-trip
PHASE 3 — enable features
- Set write.{delete,update,merge}.mode = merge-on-read where deletes are hot
- Adopt variant/geo/nanosecond types where they help
- Incremental consumers rely on default-on row lineage
PHASE 4 — fleet rollout
- Bump tables in waves; gate each wave; validate each wave
- Keep expire_snapshots retention long enough to roll back a wave
ROLLBACK
- Per table: rollback_to_snapshot to a pre-feature snapshot
- Only lower format-version if NO v3-only feature was written
-- Per-table wave: gate (external), bump, then validate
ALTER TABLE prod.db.orders SET TBLPROPERTIES ('format-version' = '3');
-- Roll back a wave if validation fails (to the snapshot before the bump)
CALL prod.system.rollback_to_snapshot(table => 'db.orders', snapshot_id => 30291);
Step-by-step trace.
| Phase | Action | Guard |
|---|---|---|
| 0 inventory | map tables↔consumers, versions | none |
| 1 runtimes | upgrade engines, stay v2 | behaviour unchanged |
| 2 canary | bump one table | compatibility gate + validation |
| 3 features | MoR deletes, new types | opt-in per table |
| 4 fleet | bump in waves | gate + validate each wave |
| rollback | rollback_to_snapshot |
retention window kept |
After running the plan, every engine is v3-capable before any production table is bumped; the compatibility gate blocks any wave that still has a stale writer; each wave is validated across Spark, Flink, Trino, and Snowflake before the next begins; and because the bump is metadata-only and snapshots are retained, any wave can be rolled back to its pre-bump snapshot without data loss.
Output:
| Risk | Mitigation | Result |
|---|---|---|
| Stale writer corrupts table | writer compatibility gate | blocked pre-bump |
| Reader can't see v3 | upgrade in phase 1 | staged capability |
| Feature regressions | canary + cross-engine validation | caught early |
| Bad wave |
rollback_to_snapshot + retention |
reversible |
| Data rewrite cost | metadata-only bump | none |
Why this works — concept by concept:
-
Metadata-only bump — because
format-version = 3is an atomic metadata commit that rewrites no data, the migration is cheap and the risk is entirely about engine capability, not data movement. - Writer compatibility gate — the only actor that can corrupt a v3 table is a writer on a stale Iceberg runtime, so gating each wave on writer versions removes the single genuine corruption path.
- Phased rollout — upgrading runtimes first, then canarying one table, then enabling features, then rolling the fleet keeps blast radius small and localizes any incompatibility to one wave.
- Cross-engine validation — checking delete visibility, lineage, and new-type round-trip from every consumer catches a lagging reader before a dashboard does, because a write that succeeds on Spark can still be unreadable on an old external reader.
- Cost — O(1) metadata per table bump, O(waves) validation runs, and a retention window sized to allow rollback. The migration's cost is coordination and validation, not compute, which is exactly where a senior plan should spend it.
Perf
Topic — optimization
Migration and compatibility optimization problems
Data
Topic — data-processing
Lakehouse rollout and pipeline-migration problems
Cheat sheet — Apache Iceberg v3 recipes
-
Metadata tree mnemonic. Catalog pointer → table
metadata.json(holds schema, partition specs, snapshot list,format-version) → manifest list (one Avro per snapshot) → manifests (Avro; list data + delete files with partition, stats,sequence_number) → data files (Parquet/ORC/Avro) + deletion vectors (Puffin). A commit is one atomic catalog compare-and-swap of the current-metadata pointer. -
Sequence-number rule. Every commit gets a monotonic sequence number; files inherit it. A delete (vector or file) applies to a data file only when
data_file.seq ≤ delete.seq. This single rule makes merge-on-read deterministic without timestamps. - Format-version meaning. v1 = analytic tables, no row-level deletes. v2 = merge-on-read via positional + equality delete files. v3 = deletion vectors, row lineage, variant/geometry/geography/nanosecond types, default values, multi-arg transforms. The integer is a contract every engine must honour.
-
Deletion-vector shape. One Roaring bitmap per data file, stored as a
deletion-vector-v1blob in a Puffin file, referenced from the manifest by offset/length; set bits = deleted row positions. New deletes update the same vector in place — no positional-delete-file swarm. Compact withrewrite_data_files; drop old snapshots withexpire_snapshots. -
CoW vs MoR dial.
write.delete.mode,write.update.mode,write.merge.modeeach choosecopy-on-write(rewrite files, fast reads) ormerge-on-read(mark + append, cheap writes). Pick MoR for hot deletes/updates; deletion vectors are the v3 MoR delete form. -
Row-lineage fields.
_row_id= table-unique stable identity (first-row-idof the data file + row position, inherited on rewrite);_last_updated_sequence_number= sequence number of the last commit that changed the row. Both default-on for v3. Join on_row_idfor CDC; filter by_last_updated_sequence_numberfor incremental scans. -
Incremental-scan recipe. Keep one durable integer watermark = highest sequence consumed; scan snapshots with
sequence_number > watermark; the planner skips files whose max sequence ≤ watermark; apply the changelog delta keyed by_row_id; advance the watermark. Replay is idempotent because ids are stable. -
Variant recipe. Store JSON-like payloads as
VARIANT(binary, self-describing), extract withvariant_get(col, '$.path', 'type')(null-tolerant across shapes), and shred hot sub-fields into typed columns for predicate pushdown. Usestringonly when you never query into the payload. -
Spatial recipe.
GEOMETRY= planar coordinates;GEOGRAPHY= spherical (globe-correct distance/containment); both WKB-encoded with a CRS. Iceberg records spatial bounds per file, soST_Intersects/ST_DWithinprune files. Pick geometry for local/projected, geography for global. -
Default-value recipe.
ADD COLUMN c T DEFAULT vsetsinitial-default(existing rows readv) andwrite-default(omitted writes getv) — a metadata-only add even for NOT NULL. Reserverewrite_data_filesfor per-row computed backfills, not constants. Binding is by field id, so rename/reorder are free too. -
Migration checklist. Inventory every reader/writer and its Iceberg library version → upgrade all runtimes while staying v2 → gate on writers ≥ your min v3 version → canary-bump one table → cross-engine validate (delete visibility, lineage, new-type round-trip) → enable MoR/new types → bump the fleet in waves → keep snapshot retention long enough to
rollback_to_snapshot. - Engine-support reality. v3 features rolled out across Iceberg library releases and engines lag the spec: verify the exact Iceberg runtime bundled by each engine (Spark, Flink, Trino, Dremio) and the documented v3 support level of managed readers (Snowflake, BigQuery/BigLake, Athena) before flipping. A new-looking engine on an old Iceberg runtime is the classic trap.
Frequently asked questions
What is new in Apache Iceberg v3?
Apache Iceberg v3 is the third format-version of the table spec, and it adds three headline capabilities on top of the same four-layer metadata tree: deletion vectors (one compact Roaring bitmap per data file, stored in a Puffin blob, that replaces v2's positional delete files for merge-on-read), row lineage (a stable _row_id and a _last_updated_sequence_number on every row, enabled by default, for incremental processing and CDC), and new binary types — variant for semi-structured data, geometry/geography for spatial data, and nanosecond-precision timestamps. v3 also adds default column values, an unknown type for forward compatibility, and multi-argument partition transforms. All of it is additive and none of it rewrites existing data on upgrade.
Deletion vectors vs positional delete files — what changed?
In v2, merge-on-read deletes were recorded as positional delete files, each naming (file_path, position) pairs; a busy delete/update workload produced a swarm of tiny delete files that every reader had to open and merge, so read latency climbed with delete volume. v3 replaces them with deletion vectors: exactly one Roaring bitmap per data file, stored as a deletion-vector-v1 blob in a Puffin sidecar, whose set bits are the deleted row positions. A reader consults a single bitmap per data file with O(1) position lookups instead of merging many files, and new deletes update the same vector in place rather than creating more files. The sequence-number rule (data.seq ≤ delete.seq) is unchanged; only the representation and its cost profile changed — read latency now stays flat as deletes accumulate, and compaction stops fighting a small-file swarm.
What is row lineage used for?
Row lineage gives every row a table-unique, immutable _row_id and a _last_updated_sequence_number that records the commit which last modified it. Because the _row_id survives updates, compaction, and rewrites, you can follow one logical row through insert → update → delete across snapshots and join old and new versions by identity rather than by a business key that might churn or be reused. That makes three things efficient: incremental processing (keep a single sequence-number watermark and scan only rows changed after it, skipping whole files by their sequence bounds), change data capture (emit precise per-row change events keyed by _row_id), and materialized-view maintenance (apply only the delta since the last refresh instead of recomputing). Lineage is enabled by default on v3 tables, so it is available to consumers without extra configuration.
What are the new binary types in Iceberg v3?
v3 adds several new value types, most of them binary-encoded. Variant stores semi-structured (JSON-like) data in a compact, self-describing binary encoding, with typed, null-tolerant field access via variant_get and optional shredding of hot sub-fields into columns for predicate pushdown. Geometry and geography store spatial values as well-known binary (WKB) with a coordinate reference system — geometry uses planar/Cartesian math, geography uses spherical (globe-correct) semantics — so spatial predicates work and Iceberg can record per-file spatial bounds for pruning. Nanosecond timestamps (timestamp_ns / timestamptz_ns) add sub-microsecond precision alongside the existing microsecond types. v3 also introduces default column values and an unknown type for forward compatibility. Storing these as native types (instead of strings) is what enables statistics, pushdown, and correct semantics.
How do I migrate a v2 table to v3?
Upgrading is a metadata-only bump of format-version from 2 to 3 (ALTER TABLE ... SET TBLPROPERTIES ('format-version' = '3')) that rewrites no data — existing data files, positional delete files, and rows are untouched, and conversions (to deletion vectors, retroactive lineage) happen lazily on the next compaction or explicit maintenance. The real work is the compatibility audit and rollout: inventory every reader and writer and the Iceberg library version each uses; upgrade all runtimes to a v3-capable version while the table stays v2; gate the bump on every writer being v3-capable (a stale writer is the only actor that can corrupt a v3 table); canary-bump one low-stakes table and validate delete visibility, lineage, and new-type round-trip from every engine; then bump the fleet in waves, keeping snapshot retention long enough to rollback_to_snapshot if a wave fails.
Which engines support Iceberg v3?
Support for v3 features rolled out incrementally across Apache Iceberg library releases, and engines lag the spec because each bundles a specific Iceberg runtime. Apache Spark, Flink, Trino, and Dremio gain v3 capabilities as they upgrade their bundled Iceberg version, and the managed Iceberg readers in Snowflake, BigQuery/BigLake, Athena, and similar services expose v3 support at their own documented pace. The reliable check is not "is the engine new?" but "which Iceberg runtime version does this engine actually bundle, and does that version implement the specific v3 features (deletion vectors, row lineage, variant/geo types) I intend to use?" Always verify the runtime version per engine and validate against a real v3 table before flipping production tables, because a modern-looking engine on an older Iceberg runtime is the classic compatibility trap.
Practice on PipeCode
- Drill the database practice library → for the table-format, schema-evolution, row-identity, and versioning problems that Iceberg v3 concepts map onto.
- Work the data-processing practice library → for merge-on-read, CDC-upsert, incremental-processing, and semi-structured/spatial pipeline scenarios.
- Tune scan cost on the optimization practice library → for read-amplification, file-pruning, and migration trade-off problems.
- Stack the fundamentals against PipeCode's broader 450+ data-engineering catalogue to anchor the v3 mental model — deletion vectors, row lineage, binary types — against real graded inputs.
Lock in Apache Iceberg v3 muscle memory
Docs explain the spec. PipeCode drills explain the decision — when merge-on-read with deletion vectors beats copy-on-write, when row lineage turns a full refresh into a delta, when a variant column beats a JSON string, and when a stale writer must block a v3 upgrade. Pipecode.ai is Leetcode for Data Engineering — format-first practice tuned for the production trade-offs data engineers actually face.
Practice data-processing problems →
Practice database problems →





Top comments (0)