Apache Iceberg: The Open Table Format Reshaping Data Lakes

By Marcus Chen • • 9 min read

Why Data Lakes Needed a Better Table Format

The Hive table format served data engineering well for a decade, but its limitations became increasingly painful as data volumes grew beyond petabyte scale. Hive tracks partitions through a central metastore that knows only directory paths. It has no concept of individual files, no transaction semantics, and no way to safely handle concurrent writes. Every engineer who has debugged a partially-written partition at 2 AM understands the cost of these limitations.

Apache Iceberg was created at Netflix to solve these problems directly. Rather than bolting transactional semantics onto an existing format, Iceberg redesigned the metadata layer from scratch. The result is a table format that provides ACID transactions, schema evolution, partition evolution, and time travel, all while remaining completely engine-agnostic.

The core insight behind Iceberg is that a table format should track every data file individually rather than relying on directory listings. This single architectural decision eliminates an entire class of correctness bugs that plague Hive-style data lakes, from phantom reads of partially-written partitions to silent data loss during failed compaction jobs.

ICEBERG METADATA ARCHITECTURE Catalog (metadata pointer) Metadata File (JSON/Avro) Snapshot S1 (v1) Snapshot S2 (v2) Snapshot S3 (current) Manifest List Manifest List Manifest List Each manifest list points to manifest files → individual data files (Parquet/ORC/Avro)

The Three-Layer Metadata Architecture

Iceberg organizes metadata into three distinct layers, each serving a specific purpose in the transactional model. Understanding this layered structure is essential for tuning performance and debugging issues in production deployments.

Catalog Layer

The catalog holds a single pointer to the current metadata file for each table. This pointer is updated atomically when a transaction commits. Iceberg supports multiple catalog implementations including Hive Metastore, AWS Glue, Nessie, and REST catalogs. The choice of catalog determines how concurrent writers coordinate, which directly impacts write throughput in multi-writer scenarios.

Metadata File Layer

Each metadata file contains the table schema, partition spec, sort order, a list of all snapshots, and a pointer to the current snapshot. When you alter a table schema or change its partition layout, Iceberg writes a new metadata file rather than modifying the existing one. This immutability is what makes schema evolution safe and reversible.

Manifest Layer

Manifest files track individual data files along with their column-level statistics: min/max values, null counts, and distinct value counts. These statistics enable Iceberg to skip entire data files during query planning without reading the files themselves. On a 10TB table with 50,000 data files, manifest-level pruning routinely eliminates 95% of files before any I/O occurs against the actual data.

The manifest list for each snapshot records which manifest files belong to that snapshot and includes partition-level summary statistics. This two-level structure allows the query planner to eliminate irrelevant manifests before examining individual file entries, reducing planning time from seconds to milliseconds even on tables with millions of files.

Snapshot Isolation and Concurrency Control

Every write operation in Iceberg creates a new snapshot. Readers always see a consistent, complete view of the table because they are bound to a specific snapshot at query start time. This snapshot isolation model means readers never block writers and writers never block readers, a critical property for data platforms serving both batch ETL and interactive queries simultaneously.

Iceberg uses optimistic concurrency control for writes. When two writers attempt to commit simultaneously, one succeeds and the other retries with a fresh view of the table state. The retry logic checks whether the conflicting commit invalidated any assumptions made during planning. For append-only workloads, conflicts never occur because appends are always compatible with each other.

-- Time travel query: read the table as it was yesterday
SELECT count(*), avg(revenue)
FROM analytics.orders
FOR SYSTEM_TIME AS OF TIMESTAMP '2026-09-30 00:00:00';

-- Query a specific snapshot by ID
SELECT *
FROM analytics.orders
FOR SYSTEM_VERSION AS OF 847293561038475;

Conflict detection becomes more nuanced for operations that delete or overwrite data. Iceberg tracks which files each operation reads and writes. If a concurrent commit modified any file that the current operation read, the operation must retry. This granular tracking avoids false conflicts that would serialize writes unnecessarily in systems with coarser locking granularity.

Partition Evolution Without Rewrites

Traditional partitioning in data lakes is a permanent decision. Changing the partition scheme of a Hive table requires rewriting every data file, an operation that can take days on large tables and requires significant coordination to avoid disrupting downstream consumers.

Iceberg eliminates this constraint through partition evolution. When you change a table's partition spec, existing data stays in its original layout while new data is written with the updated scheme. The metadata layer tracks which partition spec applies to each data file, allowing the query planner to route predicates correctly regardless of when the data was written.

-- Original partition: daily
ALTER TABLE analytics.events
SET PARTITION SPEC (day(event_timestamp));

-- Later: change to hourly for higher cardinality
ALTER TABLE analytics.events
SET PARTITION SPEC (hour(event_timestamp));

-- Even later: add a second partition dimension
ALTER TABLE analytics.events
SET PARTITION SPEC (hour(event_timestamp), bucket(16, user_id));

This capability is transformative for teams managing rapidly evolving data platforms. A table that starts with daily partitioning can transparently switch to hourly partitioning as query patterns demand finer granularity. The transition requires no data migration, no downstream pipeline changes, and no query rewrites. Related approaches for handling evolving data complement partition evolution in a complete data management strategy.

Hidden Partitioning

One of Iceberg's most underappreciated features is hidden partitioning. In Hive, users must know the partition scheme and include partition columns in their queries to benefit from partition pruning. Iceberg decouples the logical query from the physical layout through partition transforms.

When a table is partitioned by day(event_timestamp), users query using WHERE event_timestamp > '2026-09-01' without knowing or caring about the physical partition structure. Iceberg automatically translates the predicate to the partition level, pruning irrelevant partitions during planning. This eliminates an entire class of performance bugs caused by users forgetting to include partition filters in their queries.

Available transforms include year, month, day, hour, bucket, and truncate. The bucket transform applies a consistent hash function to distribute data evenly across a fixed number of partitions, which is particularly useful for high-cardinality columns like user IDs where range partitioning would create skewed partitions. For more on choosing the right file formats within these partitions, consider the tradeoffs between Parquet, ORC, and Avro.

Production Maintenance Operations

Running Iceberg tables in production requires regular maintenance to control metadata growth and data file proliferation. Three operations form the core of Iceberg table hygiene:

  • expire_snapshots removes snapshots older than a specified retention period, freeing the data files that are no longer referenced by any retained snapshot. Set retention to at least 3-7 days to allow long-running queries and lineage tracking jobs to complete.
  • rewrite_data_files compacts small files into larger ones, targeting the optimal file size of 256MB-512MB for Parquet. Small file accumulation is the most common performance degradation pattern in streaming-to-Iceberg architectures.
  • rewrite_manifests merges manifest files to reduce the number of metadata reads during query planning. This operation is less frequently needed but becomes important for tables that receive thousands of small commits per day.
-- Expire snapshots older than 7 days
CALL system.expire_snapshots(
  table => 'analytics.events',
  older_than => TIMESTAMP '2026-09-24 00:00:00',
  retain_last => 100
);

-- Compact small files targeting 512MB
CALL system.rewrite_data_files(
  table => 'analytics.events',
  options => map('target-file-size-bytes', '536870912')
);

Scheduling these operations requires balancing maintenance overhead against query performance. Most production deployments run expire_snapshots daily, rewrite_data_files hourly for streaming tables and daily for batch tables, and rewrite_manifests weekly. Monitor the pipeline freshness SLAs to ensure maintenance windows do not interfere with data delivery commitments.

ICEBERG MAINTENANCE LIFECYCLE Streaming Writes Many small files rewrite_data_files Compact → 256-512MB expire_snapshots Remove old versions Optimized Table Schedule: rewrite_data_files: hourly (streaming) / daily (batch) Retention: expire_snapshots: 7 days | rewrite_manifests: weekly Key Metrics to Monitor File count per partition Avg file size (target: 256-512MB) Snapshot count / metadata size

Engine Compatibility and Ecosystem

Iceberg's engine-agnostic design is its strongest competitive advantage. The same table can be written by Spark Structured Streaming, queried by Trino, and maintained by Airflow orchestrated jobs, all without any compatibility concerns. This interoperability eliminates vendor lock-in and allows teams to choose the best engine for each workload.

EngineReadWriteDDLMaintenance
Spark 3.xFullFullFullFull
Flink 1.16+FullFullPartialLimited
Trino/PrestoFullFullFullFull
DremioFullFullFullFull
SnowflakeFullExternalLimitedN/A
DuckDBFullExperimentalLimitedN/A

The REST catalog specification, finalized in 2024, provides a standardized API for catalog operations across all engines. Teams using a REST catalog can swap compute engines without reconfiguring catalog connectivity, further reducing the operational cost of a multi-engine data platform. For teams evaluating local analytics options, DuckDB's Iceberg support enables laptop-scale exploration of production tables without provisioning cluster resources.

Migration from Hive Tables

Migrating existing Hive tables to Iceberg does not require rewriting data files. The migrate procedure creates Iceberg metadata pointing to existing data files in place, converting a Hive table to an Iceberg table in minutes regardless of table size. After migration, all subsequent operations benefit from Iceberg's transactional semantics.

-- In-place migration: no data movement
CALL system.migrate('analytics.legacy_events');

-- Verify migration
DESCRIBE EXTENDED analytics.legacy_events;
-- Output includes: provider = iceberg, format-version = 2

The migration is reversible. If issues arise, the rollback_to_snapshot procedure restores the table to its pre-migration state. This reversibility makes it safe to migrate production tables incrementally, starting with lower-risk analytical tables and progressing to mission-critical datasets as confidence grows. Teams managing data quality across the migration can run parallel validation queries against both the Hive and Iceberg versions during the transition period.

Key Takeaways

Apache Iceberg solves the fundamental reliability problems of data lakes without sacrificing the openness and flexibility that made data lakes attractive in the first place. Its three-layer metadata architecture provides ACID transactions, its snapshot isolation model enables safe concurrent access, and its partition evolution capability eliminates the most painful migration scenario in data engineering.

For teams building new data platforms, Iceberg should be the default table format. For teams running existing Hive-based lakes, the in-place migration path removes the primary objection to adoption. The ecosystem has reached the maturity point where every major compute engine provides production-grade Iceberg support, making engine lock-in a solved problem for organizations willing to adopt an open table format as their foundation.