DEV Community

Cover image for Apache Iceberg v3: Deletion Vectors, Row Lineage & Binary Types
Gowtham Potureddi
Gowtham Potureddi

Posted on

Apache Iceberg v3: Deletion Vectors, Row Lineage & Binary Types

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.

PipeCode blog header for Apache Iceberg v3 — bold white headline 'Apache Iceberg v3' over a hero composition of an Iceberg metadata tree with deletion-vector, row-lineage, and binary-type badge medallions around a central purple 'v3' seal, on a dark gradient.

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


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 the format-version integer — 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) or bucket(16, user_id); Iceberg records the partition value in each manifest entry and prunes on it. Queries never reference the partition column explicitly — no WHERE dt = '2026-09-05' folklore.
  • Partition evolution changes the spec without rewriting data. Adding hours(event_ts) to a table that was partitioned by days(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;
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. orders.snapshots returns one row per commit. The sequence_number column is the monotonic counter; manifest_list is the Avro file for that snapshot. operation is append, overwrite, delete, or replace — a one-word summary of what the commit did.
  2. orders.manifests lists the manifest files of the current snapshot. The added/existing/deleted counts let you see at a glance whether the last commit added files (append) or rewrote them (compaction/overwrite).
  3. orders.files is the leaf layer. content = 0 is a data file, content = 1 is a positional delete file, content = 2 is an equality delete file. The partition, lower_bounds, and upper_bounds columns are exactly what the planner uses to skip files.
  4. orders.history maps wall-clock time to snapshot ids. This is the lookup you use for SELECT ... FOR SYSTEM_TIME AS OF '2026-09-01' — the engine resolves the timestamp to a snapshot id, then reads that snapshot's manifest list.
  5. 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 snapshot 30291, 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 (snapshot 30292, 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'
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. 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.
  2. It writes a new manifest listing data-0006.parquet with its partition tuple, record count, and column bounds. Then it writes a new manifest list that references the existing manifests plus the new manifest.
  3. It writes v8.metadata.json with a new snapshot 30292, parent 30291, and sequence number 6 (one more than the parent's 5).
  4. 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.
  5. 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.json on 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;
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. REPLACE PARTITION FIELD mutates 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.
  2. 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.
  3. The files metadata table exposes spec_id per file. Grouping by it shows the split: millions of legacy files under spec 0, new files under spec 1.
  4. 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.
  5. If you later want history at hourly granularity too, you run a rewrite_data_files compaction 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
Enter fullscreen mode Exit fullscreen mode

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-version integer 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

Practice →

Data Topic — data-processing Data-processing problems on lakehouse pipelines

Practice →


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."

Iconographic deletion-vectors diagram — a data file card, a compact bitmap grid marking deleted row positions inside a Puffin blob, and a reader merging only live rows, replacing a stack of tiny positional delete files.

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, and write.merge.mode choose copy-on-write or merge-on-read independently. 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 — "in data-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 = 42 is 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-v1 blobs 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()))
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. The deletion vector for data-0002.parquet is 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.
  2. 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.
  3. 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.
  4. 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.
  5. Because the vector is addressed by data-file path, a reader that opens data-0002.parquet fetches 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)
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. 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).
  2. 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.
  3. Only after the deletion skip does the engine apply the WHERE filter, so deleted rows never reach predicate evaluation, projection, or aggregation. The count reflects only live, matching rows.
  4. Compaction cost drops correspondingly. To garbage-collect deletes, rewrite_data_files reads 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.
  5. 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);
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. 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).
  2. 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.
  3. 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.
  4. The commit produces a new snapshot referencing data-0007.parquet (new versions) and the updated deletion vector for data-0002.parquet (retiring old versions). Sequence numbers ensure the vector, committed at seq K, applies to data-0002 which was committed at seq < K.
  5. 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 in data-0002 via 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');
Enter fullscreen mode Exit fullscreen mode
# 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")]))
Enter fullscreen mode Exit fullscreen mode

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-read tells 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_files materializes vector deletes into clean data files and expire_snapshots drops 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

Practice →

Perf Topic — optimization Scan-cost and read-amplification optimization

Practice →


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.

Iconographic row-lineage diagram — three data rows each carrying a _row_id tag and a _last_updated_sequence_number chip, tracked across three snapshot columns as the same row_id is updated and deleted over sequence numbers.

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_id even 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 its first-row-id. Row k in that file has _row_id = first_row_id + k unless 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_id across 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_id assigned?" — first-row-id of 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_id 0,1,2,3.
  • File B. first-row-id = 4, 3 rows. Rows get _row_id 4,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
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. The writer reserves a contiguous id block per data file. data-A reserves [0,4); its rows get ids 0–3 by adding position to first-row-id = 0.
  2. data-B reserves the next block [4,7); its rows get 4–6. The blocks never overlap because the counter only advances.
  3. The final counter value 7 is persisted as next-row-id in table metadata. The next data file to be committed will start at first-row-id = 7.
  4. 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.
  5. 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 in data-A marked in a deletion vector.
  • Seq 3 (delete). Row 101's position in data-C marked 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;
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. At sequence 1, row 101 is inserted with _row_id = 101 and _last_updated_sequence_number = 1. Any incremental consumer with watermark < 1 will pick it up.
  2. 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_number advances 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.
  3. 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_id retire their copy.
  4. 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 key id. If the application had reused id = 101 for a different logical entity, lineage would still keep them distinct because each got its own _row_id.
  5. 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)
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. The consumer keeps a durable watermark — the highest _last_updated_sequence_number it has processed (here 4). It never needs an external CDC table; the watermark is a single integer.
  2. 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.
  3. Only data-D (seq 5) and data-E (seq 6) are opened, and within them only rows whose _last_updated_sequence_number > 4 are emitted. The bulk of the table is never touched.
  4. 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.
  5. 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_id is 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);
Enter fullscreen mode Exit fullscreen mode
# 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))
Enter fullscreen mode Exit fullscreen mode

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_id as 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_id as 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

Practice →

SQL Topic — database Database problems on row identity and versioning

Practice →


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.

Iconographic Iceberg format-layering diagram — catalog pointer to a table metadata JSON branching to a manifest list, manifests, and data files, with a snapshot-isolation lens and a side panel of new v3 type glyphs (variant, geometry, nanosecond timestamp, defaults).

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). geometry uses planar (Cartesian) coordinates; geography uses 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_ns and timestamptz_ns add nanosecond precision alongside the existing microsecond timestamp / 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 a write-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 unknown type. A forward-compatibility placeholder: a reader that encounters a column of a type it does not understand can treat it as unknown (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-default backfills 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_type and payload.value with 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;
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. The payload variant column 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.
  2. parse_json converts the incoming text into the variant binary at write time, so reads never re-parse text — extraction walks the binary structure directly.
  3. variant_get(payload, '$.value', 'int') navigates the document by path and coerces to the requested type. Event B has no value, so it returns NULL rather than failing — variant access is null-tolerant across heterogeneous rows.
  4. With variant shredding, an engine can additionally materialize hot sub-fields (say event_type) as a typed sub-column so filters like WHERE event_type = 'click' push down to file statistics, while the full document remains available for rare fields.
  5. 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))')
       );
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. pickup geometry stores 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.
  2. 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.
  3. The ST_Intersects predicate 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.
  4. geometry uses planar math (fast, correct for local/projected coordinates); had we needed great-circle distances across the globe we would have chosen geography for spherical semantics.
  5. 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');
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. ADD COLUMN currency STRING DEFAULT 'USD' records the column in the schema with initial-default and write-default both 'USD'. It commits as a metadata-only change — no data file is read or rewritten.
  2. Reading an old data file, the engine sees no currency field (by field id) and synthesizes the initial-default 'USD' for those rows, so the NOT NULL contract holds logically without a physical backfill.
  3. A new INSERT that omits currency gets the write-default 'USD' materialized into the new data file. A new INSERT that supplies currency stores the supplied value.
  4. Because binding is by field id, a later rename of currency would again be metadata-only; the defaults travel with the field id.
  5. 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';
Enter fullscreen mode Exit fullscreen mode

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 string column 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_ns preserve 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-default lets existing rows read a value for a newly added column and write-default fills 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

Practice →

Data Topic — data-processing Semi-structured and spatial data-processing problems

Practice →


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.

Iconographic migration diagram — a v2 table card upgrading to a v3 table card via a format-version bump arrow, a compatibility gate checking reader/writer support, and an engine-support matrix of check and cross marks for common query engines.

What the version bump does and does not do.

  • Does. Sets format-version = 3 in 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 = 3 on 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
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. The pre-upgrade count of content = 0 files establishes the baseline data-file inventory.
  2. The ALTER TABLE ... SET TBLPROPERTIES commit writes a new metadata JSON with format-version = 3 and swaps the catalog pointer. It is an ordinary atomic commit and completes in milliseconds.
  3. 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.
  4. The newest snapshot's summary shows 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.
  5. 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")
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. MIN_V3_WRITER encodes 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."
  2. 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.
  3. flink-cdc on 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.
  4. 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.
  5. 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'
Enter fullscreen mode Exit fullscreen mode

Step-by-step explanation.

  1. 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.
  2. Selecting the known-deleted id = 104 must 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.
  3. 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.
  4. Extracting a known variant field must return the written value ('click'), proving the new type round-trips through write and read on this engine.
  5. 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
Enter fullscreen mode Exit fullscreen mode
-- 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);
Enter fullscreen mode Exit fullscreen mode

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 = 3 is 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

Practice →

Data
Topic — data-processing
Lakehouse rollout and pipeline-migration problems

Practice →


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-v1 blob 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 with rewrite_data_files; drop old snapshots with expire_snapshots.
  • CoW vs MoR dial. write.delete.mode, write.update.mode, write.merge.mode each choose copy-on-write (rewrite files, fast reads) or merge-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-id of 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_id for CDC; filter by _last_updated_sequence_number for 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 with variant_get(col, '$.path', 'type') (null-tolerant across shapes), and shred hot sub-fields into typed columns for predicate pushdown. Use string only 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, so ST_Intersects / ST_DWithin prune files. Pick geometry for local/projected, geography for global.
  • Default-value recipe. ADD COLUMN c T DEFAULT v sets initial-default (existing rows read v) and write-default (omitted writes get v) — a metadata-only add even for NOT NULL. Reserve rewrite_data_files for 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)