Storage Layers: Advanced
What you will be able to do
The Lakehouse: ACID on Object
Recognize the lakehouse as object storage plus a metadata layer and explain why the metadata is the actual product.
What the Metadata Layer Adds
| Capability | Plain Lake | Lakehouse Table Format |
|---|---|---|
| Atomic writes across files | Best-effort: a writer fails halfway leaves partial files | Atomic: either the snapshot commits or it does not |
| Snapshot isolation between readers and writers | Readers can see partial writes mid-operation | Readers see a consistent snapshot regardless of concurrent writes |
| Time travel to historical state | Possible only if old files were never overwritten | First-class: every snapshot is reachable by version or timestamp |
| Schema evolution without rewrites | Adding a column means rewriting partitions or splitting tables | Add, drop, rename at the metadata level; data files unchanged |
| Row-level updates and deletes | Possible only by rewriting whole partitions | Supported via merge-on-read or copy-on-write strategies |
| Concurrent writers | Last writer wins; data corruption likely | Optimistic concurrency: conflicting writers fail and retry |
How an Iceberg Table Is Laid Out
The Three Major Open Formats
Why the Metadata Is the Product
- ▸Adding a column to the schema
- ▸Renaming a column at the schema level
- ▸Tagging a snapshot for time travel reference
- ▸Setting a retention policy on old snapshots
What Plain Object Storage Cannot Do
- Choose an open table format for any analytical table that needs concurrent writers or schema evolution
- Treat the metadata layer as the operational surface, not just the data files
- Use hidden partitioning (Iceberg) or generated columns (Delta) so the partition logic survives schema evolution
- Edit data files directly; the metadata loses track of them
- Mix table-format-aware writers with plain object-store writers; the table state diverges
- Assume all three formats interoperate fully today; check engine compatibility before standardizing
Snapshot Isolation and Time Travel
Apply snapshot isolation and time travel to operate a table format under concurrent writers and reproduce historical state.
Snapshot Isolation in Practice
Time Travel
Operational Uses of Time Travel
| Use Case | What It Enables | What To Watch For |
|---|---|---|
| Debug a bad pipeline run | Compare today's output to yesterday's at the row level | Retention must extend back to the bad run |
| Reproduce an ML training run | Train against the exact snapshot the original model used | Pin the snapshot id at training time and store it with the model |
| Audit historical reporting | Reproduce the exact numbers a regulator saw last quarter | Snapshots required for compliance must not be expired by retention |
| Roll back an accidental DELETE | Restore the table to a snapshot before the destructive write | Roll back is itself a write; storage cost briefly doubles |
| Compare two versions of a transform | Run the new transform against an old snapshot and diff the result | Both snapshots must remain in retention for the comparison |
Optimistic Concurrency Control
- ▸Many writers update the same partition simultaneously
- ▸Long-running writes hold a stale parent snapshot for minutes
- ▸Streaming jobs commit at high frequency and each commit blocks the next
- ▸Compaction and ingestion compete for the same partition's metadata
Retention and Cost
Schema Evolution Without Rewrites
Apply schema evolution at the table-format level: add, drop, rename, and widen columns without rewriting data files.
Operations and Their Costs
| Operation | Cost in Plain Parquet Lake | Cost in Iceberg / Delta |
|---|---|---|
| Add a nullable column | Rewrite all partitions or split table | Metadata-only; data files unchanged |
| Rename a column | Rewrite all partitions | Metadata-only (Iceberg); Delta requires column-mapping mode |
| Drop a column | Rewrite all partitions | Metadata-only; data files retain bytes but readers ignore them |
| Reorder columns | Application-level concern only | Metadata-only; column ordering is logical |
| Change a column type (widening) | Rewrite all partitions | Metadata-only for safe widenings (int -> long, decimal precision up) |
| Change a column type (narrowing) | Lossy; usually new column instead | Disallowed without rewrite; force a new column for the new type |
Why Renaming Is Hard in Plain Parquet
Adding a Column
Dropping a Column
Type Widening Versus Type Change
- Add a nullable column
- Rename a column (with id-mapping table format)
- Drop a column
- Widen an integer or decimal type
- Narrow a type or change a type incompatibly
- Drop a column that downstream consumers still read
- Rename a column without coordinating with consumers
- Add a non-nullable column without a default value
Partition Evolution
- ▸Treat the schema as a contract; coordinate breaking changes with consumers
- ▸Add new columns nullable; populate them in a separate write
- ▸Use the expand-contract pattern for type changes that are not pure widenings
- ▸Audit schema changes via the snapshot history; every change is recorded
- Stable column ids embedded in the table schema
- Renames change only the schema mapping; data files unchanged
- Hidden partitioning lets partition spec evolve independently of the schema
- Engine neutrality is a design goal: same metadata works across Spark, Trino, Snowflake
- Column-mapping mode required to support physical-name renames
- Default mode keeps physical and logical names aligned
- Generated columns express derived partition values without separate spec
- Strongest tooling on Databricks; broad open-source support since 2023
Schema evolution is the feature that turns table format adoption from a nice-to-have into a load-bearing capability. A table that needs a column added every quarter is a table that needs a metadata-driven format.
The Small Files Problem
Diagnose the small files problem from file counts and sizes and run compaction or OPTIMIZE to restore query performance.
Why Small Files Hurt
| Cost | Per-File Overhead | Why Many Files Compounds |
|---|---|---|
| Object storage requests | Each open is a billable request | Thousands of files = thousands of billed requests per query |
| Metadata reads | Each Parquet footer has to be parsed | Footer reads dominate when files are small relative to footer size |
| Task scheduling | Each file is a separate task in Spark / Trino | Scheduler overhead exceeds actual work for small files |
| Compression efficiency | Per-file dictionary encoding is less effective on small batches | Total bytes on disk grow even though raw data is the same |
| Query planner cost | Manifest lists grow with file count | Iceberg / Delta metadata operations slow as manifests bloat |
The Numbers
Compaction Mechanics
Where the Small Files Come From
| Source | Why It Produces Small Files | Mitigation |
|---|---|---|
| Streaming micro-batches | Trigger interval forces a write per batch per partition | Schedule hourly or daily compaction; tune trigger interval up if SLA permits |
| Partition cardinality too high | Each partition gets a small slice of each batch | Reduce partition cardinality; cluster on the high-cardinality column instead |
| Concurrent writers without coordination | Each writer produces its own files for the same partition | Single ingestion pipeline per table; or coordinate via merge |
| CDC ingestion of high-update workloads | Each captured change becomes a small change file | Periodic merge-on-read compaction; switch to copy-on-write for read-heavy tables |
| Backfills with many small batches | Backfill granularity exceeds final partition granularity | Run backfill into a staging table and bulk-rewrite into the production table |
Compaction As an Operational Discipline
- ▸Set a target file size (typically 128 MB to 512 MB)
- ▸Schedule compaction at a cadence proportional to write rate (hourly for streams, daily for batch)
- ▸Monitor average and median file size per table; alert when median drops below 64 MB
- ▸Run vacuum or expire_snapshots to reclaim storage from rewritten files
- ▸Keep compaction in a separate compute pool from query workloads to avoid contention
Sorting and Z-Ordering During Compaction
Choosing Storage Across Workloads
Design a multi-layer storage architecture that combines an operational store, a lakehouse archive, and a warehouse mart for a workload with multiple concurrent access patterns.
The Three Workloads
| Workload | Access Pattern | Freshness Need |
|---|---|---|
| Regulatory archive | Bulk read of historical transactions for audit and compliance | Daily ingest is fine; reads happen quarterly |
| Customer-facing app: account history view | Single-row lookups by account_id, sub-50ms latency | Latest transaction must appear within seconds |
| Analytical BI: revenue and risk dashboards | Aggregations across millions to billions of rows by date and product | Daily by 7am for executive review; near-real-time for risk |
The Storage Layer Per Workload
How the Pipeline Glues Them Together
Choosing the Open Format
| Property | Why It Mattered Here | Format Choice |
|---|---|---|
| Multi-engine reads (Spark, Trino, Snowflake, BigQuery) | Compliance, BI, and risk teams use different engines | Iceberg has the broadest neutral engine support |
| Time travel for compliance | Auditors ask for historical state at a specified date | Iceberg's snapshot model exposes time travel cleanly via SQL |
| Schema evolution over seven years | Transaction shape changes every couple of years | Iceberg's id-based schema renames without rewriting |
| Concurrent CDC writers and analytical readers | Streaming writes from DynamoDB run alongside dashboard reads | Snapshot isolation handles the concurrency |
Connecting to Earlier Lessons
- Three separate copies of the transaction history
- Three separate ingestion pipelines
- Three different schema evolution policies
- Reconciliation across stores is a perpetual problem
- Operational copy in DynamoDB for the app
- Single Iceberg archive serving compliance, risk, and analytical preprocessing
- Snowflake mart derived from Iceberg for BI consumers
- Snapshot isolation guarantees the layers stay coherent
- ▸Application workloads with sub-50ms latency belong in an operational store, not a lake
- ▸Analytical workloads with aggregations belong in a columnar layer (warehouse or lakehouse)
- ▸Bulk historical retention belongs in object storage with a table format on top
- ▸Multi-engine analytical access pushes the choice toward Iceberg
- ▸Heavy mutation rate via CDC pushes the choice toward Hudi or copy-on-write strategies
What Goes Wrong Without This Discipline
- Map each workload to a storage shape based on access pattern, not historical preference
- Use one open table format consistently rather than mixing Iceberg, Delta, and Hudi without reason
- Treat the operational copy as derived from the lake archive, not the source of truth
- Force application traffic onto a warehouse with 100ms latency
- Force analytical scans onto an operational database with row storage
- Maintain three separate copies of transaction history when one Iceberg archive can serve archive, audit, and analytical scans
Each storage shape fits a job: the lake holds cheap raw files, the warehouse serves analytics, the operational DB serves the live app. Pick by how the data is read.
> A fintech platform serves three concurrent workloads against the same logical transactions: a customer-facing account view at sub-50ms latency, a daily executive BI dashboard, and a seven-year regulatory archive that auditors query quarterly. The current architecture loads everything into Snowflake, the app feels slow, audit costs are high, and a streaming risk service is being added at sub-hour freshness. The new staff data engineer is asked: 'What is the smallest architecture that makes all four workloads operable, and what does each piece actually do?'
Open table formats turn object storage into a transactional database without giving up the lake
- Category
- Pipeline Architecture
- Difficulty
- advanced
- Duration
- 38 minutes
- Challenges
- 0 hands-on challenges
Topics covered: The Lakehouse: ACID on Object, Snapshot Isolation and Time Travel, Schema Evolution Without Rewrites, The Small Files Problem, Choosing Storage Across Workloads
Lesson Sections
- The Lakehouse: ACID on Object (concepts: paTableFormats)
The lakehouse is a marketing term that names a real architectural shift. The shift is the addition of a metadata layer on top of files in object storage that provides the consistency guarantees a database has and a folder of files lacks. Iceberg, Delta Lake, and Apache Hudi are three implementations of the same idea. The data files are still Parquet (or ORC). The folders look mostly the same. The difference is a small set of metadata files that turn the directory of Parquet into a transactional
- Snapshot Isolation and Time Travel (concepts: paTableFormats)
Snapshot isolation is the consistency guarantee that turns a table format into a real table. A reader sees a consistent point-in-time view of the table, even when writers are landing new data concurrently. Time travel is the operational use of the same machinery: read the table as it existed at a previous snapshot. The two features come from the same underlying mechanism, which is that the table's state is defined by an immutable chain of snapshots and the metadata pointer that names which snaps
- Schema Evolution Without Rewrites (concepts: paSchemaEvolution)
A real production table changes shape over time. Producers add fields. Old fields get renamed. Columns become obsolete and need to be dropped. In a plain lake, every shape change requires rewriting partitions or splitting into a new table. In an open table format, additive and renaming changes happen at the metadata level and the data files stay where they are. The cost is bytes of metadata, not bytes of data. Operations and Their Costs Why Renaming Is Hard in Plain Parquet A plain Parquet file
- The Small Files Problem (concepts: paSmallFiles)
A 30-second streaming job writes a file every 30 seconds per partition. That is 2,880 files per partition per day. A daily ingestion that writes one large file per partition produces one. The two designs run the same SQL the same way, but the streaming version is dramatically slower because every read has to open thousands of tiny files instead of a few large ones. This is the small files problem, and it is the single most common operational headache in lake and lakehouse environments. The fix i
- Choosing Storage Across Workloads (concepts: paDataLake)
A real production system rarely has one workload. The example here is a financial services platform with three concurrent demands on the same logical data: a regulatory archive that must retain seven years of transactions, a customer-facing app that needs single-row lookups under 50 milliseconds, and an analytical BI workload that runs daily aggregations across the entire history. No single storage layer is correct for all three. The right answer is a multi-layer architecture in which each workl