Skip to content
Apache Iceberg Tutorial: Tables, Snapshots, Time Travel

Click to use (opens in a new tab)

Apache Iceberg Tutorial: Tables, Snapshots, Time Travel

September 1, 2026 by Chat2DBChat2DB Team

Apache Iceberg turns a directory of Parquet files into something that behaves like a database table: atomic commits, consistent reads, schema evolution, and a full history you can query. This tutorial builds an Iceberg table from scratch, looks at the files it writes, and walks through the operations you will actually use — hidden partitioning, time travel, branching and compaction.

The mental model

An Iceberg table is three layers:

  1. Data files — Parquet (or ORC/Avro) files holding the rows. Iceberg does not care where they sit in the directory tree; file paths are recorded, not inferred.
  2. Metadata files — manifests listing data files with per-column statistics, manifest lists grouping manifests into a snapshot, and a table metadata file holding the schema, partition specs and snapshot history.
  3. A catalog — a small, transactional pointer that says "the current metadata file for db.events is v7.metadata.json".

Every write produces a new snapshot and a new metadata file; the catalog atomically swaps the pointer. Readers that started before the swap keep reading the old snapshot and see a consistent table. That is the whole trick.

The catalog is the one piece you must choose deliberately. Options:

  • REST catalog — the standard since Iceberg 1.x; implementations include Polaris, Lakekeeper, Gravitino and the vendor catalogs.
  • AWS Glue — convenient on AWS.
  • Nessie — adds Git-like multi-table branching.
  • JDBC — a table in Postgres or MySQL; fine for small deployments.
  • Hadoop/filesystem — no external service, but unsafe for concurrent writers on object storage. Do not use it in production.

Setting up

The quickest local setup is Spark with the Iceberg runtime and a JDBC catalog backed by SQLite:

spark-sql \
  --packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.9.0 \
  --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
  --conf spark.sql.catalog.demo=org.apache.iceberg.spark.SparkCatalog \
  --conf spark.sql.catalog.demo.type=jdbc \
  --conf spark.sql.catalog.demo.uri=jdbc:sqlite:/tmp/iceberg_catalog.db \
  --conf spark.sql.catalog.demo.warehouse=/tmp/iceberg_warehouse \
  --conf spark.sql.defaultCatalog=demo

If you prefer not to run Spark, DuckDB reads Iceberg tables directly:

INSTALL iceberg;
LOAD iceberg;
 
SELECT count(*) FROM iceberg_scan('/tmp/iceberg_warehouse/db/events');

Creating a table

CREATE NAMESPACE IF NOT EXISTS db;
 
CREATE TABLE db.events (
  event_id   bigint,
  event_ts   timestamp,
  user_id    bigint,
  event_type string,
  amount     decimal(12,2),
  country    string
)
USING iceberg
PARTITIONED BY (days(event_ts), bucket(8, user_id))
TBLPROPERTIES (
  'write.format.default'         = 'parquet',
  'write.parquet.compression-codec' = 'zstd',
  'write.target-file-size-bytes' = '134217728',
  'format-version'               = '2'
);

Two things to notice.

days(event_ts) is a partition transform, not a stored column. Iceberg computes the day from the timestamp itself. Available transforms are years, months, days, hours, bucket(N, col), truncate(W, col) and identity.

bucket(8, user_id) hashes user_id into 8 buckets. This spreads writes evenly and, critically, lets Iceberg prune on equality filters against user_id — it hashes your predicate value and reads only that bucket.

Insert some rows:

INSERT INTO db.events VALUES
  (1, TIMESTAMP '2026-09-01 08:12:00', 1001, 'purchase', 42.50, 'DE'),
  (2, TIMESTAMP '2026-09-01 09:03:00', 1002, 'view',      0.00, 'FR'),
  (3, TIMESTAMP '2026-09-02 11:44:00', 1001, 'purchase', 19.99, 'DE'),
  (4, TIMESTAMP '2026-09-02 12:10:00', 1003, 'refund',  -19.99, 'US');

What Iceberg wrote

/tmp/iceberg_warehouse/db/events/
  data/
    event_ts_day=2026-09-01/user_id_bucket=3/00000-0-abc.parquet
    event_ts_day=2026-09-02/user_id_bucket=5/00001-0-def.parquet
  metadata/
    00000-....metadata.json
    00001-....metadata.json
    snap-4592837462938-1-....avro
    a1b2c3-m0.avro

Directory names look Hive-like, but they are a convenience — Iceberg reads file paths from manifests, never from a directory listing. You can move files and rewrite the manifests without breaking anything.

Iceberg exposes its own metadata as queryable tables, which is the best way to understand what is going on:

-- Every snapshot ever created
SELECT snapshot_id, committed_at, operation, summary['added-records'] AS added
FROM db.events.snapshots;
 
-- Every data file, with statistics
SELECT file_path, record_count, file_size_in_bytes, partition
FROM db.events.files;
 
-- History, including which snapshot was current when
SELECT * FROM db.events.history;
 
-- Per-partition summary
SELECT partition, record_count, file_count FROM db.events.partitions;
 
-- Manifest-level view
SELECT path, added_snapshot_id, added_data_files_count FROM db.events.manifests;

db.events.files is worth exploring: it contains lower_bounds and upper_bounds per column, which is exactly what the planner uses to skip files. If a query is reading more files than it should, this table tells you why.

Hidden partitioning in practice

The query filters on the raw timestamp column — no derived partition column, no special knowledge required:

SELECT event_type, count(*), sum(amount)
FROM db.events
WHERE event_ts >= TIMESTAMP '2026-09-02 00:00:00'
  AND event_ts <  TIMESTAMP '2026-09-03 00:00:00'
GROUP BY event_type;

Iceberg applies days() to the predicate bounds and skips every manifest and file outside that day. Compare with Hive-style tables, where the same query scans everything unless the author remembers to filter on event_date too.

Equality filters hit the bucket transform as well:

SELECT * FROM db.events WHERE user_id = 1001;   -- reads 1 of 8 buckets

Schema and partition evolution

Both are metadata operations. No data is rewritten.

ALTER TABLE db.events ADD COLUMN device string;
ALTER TABLE db.events RENAME COLUMN country TO country_code;
ALTER TABLE db.events ALTER COLUMN amount TYPE decimal(16,2);
ALTER TABLE db.events DROP COLUMN device;

Columns are tracked by ID, so a rename is safe and a dropped-then-re-added name never picks up the old column's data. Type changes are allowed when they cannot lose information — int to bigint, float to double, widening decimal precision.

Partition evolution changes the layout going forward:

-- Traffic grew; switch from daily to hourly partitions
ALTER TABLE db.events REPLACE PARTITION FIELD days(event_ts) WITH hours(event_ts);
 
-- Add a second partition dimension
ALTER TABLE db.events ADD PARTITION FIELD country_code;

Existing files keep their old partitioning. Iceberg records each partition spec by ID and plans across specs, so a query spanning old and new data works without a rewrite.

Time travel and rollback

Every snapshot stays readable until you expire it.

-- By snapshot id
SELECT count(*) FROM db.events VERSION AS OF 4592837462938;
 
-- By timestamp
SELECT count(*) FROM db.events FOR TIMESTAMP AS OF TIMESTAMP '2026-09-01 23:59:59';
 
-- Incremental read: only rows added between two snapshots
SELECT * FROM db.events
FOR SYSTEM_VERSION AS OF 4592837462938;

If a bad job lands, roll back:

CALL demo.system.rollback_to_snapshot('db.events', 4592837462938);
 
-- Or, by wall-clock time
CALL demo.system.rollback_to_timestamp('db.events', TIMESTAMP '2026-09-01 12:00:00');

Rollback is itself a commit — it creates a new snapshot pointing at the old state — so the history of the mistake is preserved.

Branches and tags

Branches let you validate before publishing, the write-audit-publish pattern:

ALTER TABLE db.events CREATE BRANCH staging;
 
-- Write into the branch only
INSERT INTO db.events.branch_staging
SELECT * FROM raw_events WHERE ingest_date = '2026-09-01';
 
-- Validate
SELECT count(*) FROM db.events VERSION AS OF 'staging' WHERE amount IS NULL;
 
-- Publish atomically if the checks pass
CALL demo.system.fast_forward('db.events', 'main', 'staging');

Tags mark a snapshot for retention — useful for regulatory or reproducibility requirements:

ALTER TABLE db.events CREATE TAG month_end_2026_09 RETAIN 365 DAYS;

Row-level updates

Format version 2 supports MERGE, UPDATE and DELETE:

MERGE INTO db.events t
USING corrections s
ON t.event_id = s.event_id
WHEN MATCHED AND s.deleted THEN DELETE
WHEN MATCHED THEN UPDATE SET t.amount = s.amount
WHEN NOT MATCHED THEN INSERT *;

Choose how deletes are materialised:

ALTER TABLE db.events SET TBLPROPERTIES (
  'write.delete.mode' = 'merge-on-read',   -- fast writes, delete files applied at read
  'write.update.mode' = 'copy-on-write'    -- rewrite files, fastest reads
);

Streaming or frequent-update workloads want merge-on-read plus regular compaction. Batch-loaded, read-heavy tables want copy-on-write.

Maintenance — do not skip this

An Iceberg table that is never maintained accumulates small files and old snapshots until planning slows to a crawl. Three procedures, scheduled regularly:

-- 1. Compact small files into target-sized ones
CALL demo.system.rewrite_data_files(
  table => 'db.events',
  strategy => 'binpack',
  options => map('min-input-files', '5', 'target-file-size-bytes', '134217728')
);
 
-- 2. Expire snapshots older than the retention window (this is what frees storage)
CALL demo.system.expire_snapshots(
  table => 'db.events',
  older_than => TIMESTAMP '2026-08-01 00:00:00',
  retain_last => 10
);
 
-- 3. Remove files no snapshot references (from failed jobs)
CALL demo.system.remove_orphan_files(
  table => 'db.events',
  older_than => TIMESTAMP '2026-08-25 00:00:00'
);

Also compact metadata on tables with many commits:

CALL demo.system.rewrite_manifests('db.events');

Watch the file count as your health metric:

SELECT partition,
       count(*)                       AS files,
       sum(record_count)              AS rows,
       avg(file_size_in_bytes)::bigint AS avg_bytes
FROM db.events.files
GROUP BY partition
ORDER BY files DESC
LIMIT 20;

Partitions with hundreds of small files are the ones to compact first.

Querying from other engines

The point of Iceberg is that the table is not owned by one engine. The same table reads from Trino, DuckDB, Snowflake, Flink and ClickHouse via the catalog. If your day involves comparing an Iceberg table against the Postgres or MySQL source it came from, a client that connects to both in one window saves the constant context switch — Chat2DB (opens in a new tab) supports Trino, Spark SQL, Postgres, MySQL and 20+ more, with a browser version at app.chat2db.ai (opens in a new tab).

Summary

Iceberg's design is three layers — data files, a metadata tree of manifests and snapshots, and a catalog that atomically swaps the current pointer. That gives atomic commits and consistent reads without listing object storage. Hidden partitioning means queries filter on real columns and still prune; partition evolution changes the layout without a rewrite; snapshots give free time travel and rollback; branches let you validate before publishing. The one thing you must own operationally is maintenance: compact data files, rewrite manifests, expire snapshots and remove orphans on a schedule.