A lakehouse is a data lake that behaves like a database table. The data still sits as open files in cheap object storage or HDFS, readable by any engine, but a thin metadata layer on top adds what a warehouse used to provide: atomic commits, consistent snapshots for readers, schema evolution, row-level updates and deletes, and time travel. Spark, Trino, Flink and Python can all read and write the same tables without exporting copies between systems.
The idea is simple, but the details decide whether a lakehouse is fast and cheap or slow and fragile. This article explains the architecture from first principles: what each layer guarantees, how a commit can be atomic on storage that has no transactions, why the catalog matters more than it looks, and what ingestion and maintenance you must run. It stays neutral between the three open table formats. Format-specific detail lives in Spark with Iceberg, Delta Lake architecture and Spark with Hudi.
What problem the lakehouse solves
The classic Hadoop data lake stored tables as directories of files registered in the Hive metastore. That design has three structural problems. First, a table is whatever files happen to be in the directory when you list it, so a reader that lists while a job is writing sees half a result, and a failed job leaves garbage that readers pick up. Second, listing is the planning step, and listing millions of files on object storage is slow and, on cloud stores, billed per request. Third, changing one row means rewriting a whole partition, and two jobs doing so silently overwrite each other.
The lakehouse keeps open storage and fixes the semantics by changing one thing: a table is no longer a directory, it is an explicit list of files recorded in metadata, and the current version of that list is chosen by a single atomic pointer.
The layers and what each one owns
Storage (S3, GCS, Azure Blob or HDFS) provides durable, immutable objects. It offers no transactions across objects, and that is acceptable because the layers above never modify a file in place.
File format. Parquet stores data by column in row groups with min and max statistics, so a query reads only the columns it needs and skips row groups that cannot match.
Table format. Iceberg, Delta Lake and Hudi each define a metadata structure that records which data files make up each version of the table, the schema, the partitioning and per-file statistics. Iceberg uses a tree: a table metadata file points to a manifest list per snapshot, which points to manifests, which list data files with their statistics. Delta Lake uses an ordered transaction log of JSON commit files with periodic Parquet checkpoints. Hudi uses a timeline of instants and organises files into file groups, with copy-on-write and merge-on-read table types.
Catalog. The catalog maps a table name to the location of its current metadata and performs the atomic swap that makes a commit visible. Options include the Hive metastore, AWS Glue, a REST catalog service such as Apache Polaris or Unity Catalog, and Nessie.
How a commit is atomic without transactions
This is the core trick. A writer never changes what readers can see until the final step. It writes new data files to fresh, unique paths; at that moment no metadata references them, so they are invisible. It then writes new metadata files describing the next version: the previous file list, plus the new files, minus any files it replaced. Finally it asks the catalog to move the table's pointer from the version it started from to the new metadata, and only if the pointer still holds the version it started from. That compare-and-swap is the commit.
In pseudocode, every format's writer does something like this:
def commit(catalog, table, base, added_files, removed_files, read_filter):
while True:
current = catalog.load_pointer(table)
if current != base:
# Someone committed since we started. Did they touch what we depend on?
for change in changes_between(base, current):
if overlaps(change.removed, removed_files) or overlaps(change.added, read_filter):
raise CommitConflict(change)
new_meta = write_metadata(parent=current, add=added_files, remove=removed_files)
if catalog.compare_and_swap(table, expected=current, new=new_meta):
return new_meta # visible to every reader from now on
base = current # lost the race: re-validate against the newer versionTwo consequences follow. Readers get snapshot isolation: a query loads the pointer once and reads that version's files, whatever commits happen meanwhile. And concurrency is optimistic: two jobs appending to the same table both succeed after a retry, while two jobs that delete or rewrite the same files conflict and one fails. Design pipelines so that writers that might overlap own disjoint partitions, or serialise them.
The atomicity is only as good as the swap. On HDFS a rename can serve; on object storage the swap must come from a catalog service or a storage primitive that offers a conditional write. Without a real compare-and-swap, two writers can both believe they committed, and one commit disappears.
Choosing the catalog
The catalog outlives any engine choice. Ask which engines must commit, where access control should live, and whether you need more than a pointer swap.
| Catalog | Good fit | Watch out for |
|---|---|---|
| Hive metastore | Existing Hadoop estates; every engine speaks it | A database to operate; locks for commits; coarse access control |
| AWS Glue | AWS-only estates | Service limits; cloud lock-in for the metadata |
| REST catalog (Polaris, Unity and others) | Many engines and languages; credentials vended per table | A service to run or buy; check each engine's REST support |
| Nessie | Git-like branches and tags across tables | Another service; branching semantics teams must learn |
For Iceberg, the REST catalog protocol has become the common interface: the engine talks HTTP to the service, which owns the commit and can hand out short-lived storage credentials scoped to one table. Configuring Spark looks like this:
spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
spark.sql.catalog.lake=org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.lake.type=rest
spark.sql.catalog.lake.uri=https://catalog.internal.example/api/catalog
spark.sql.catalog.lake.warehouse=analyticsOther engines point at the same URI and share one commit path. The Iceberg and Trino stack walks through catalog choice from the SQL engine's side.
Designing tables: layout decides cost
Three design decisions control almost all read performance. Partitioning should be coarse: partition by day, or by a low-cardinality key such as region, so that partitions hold hundreds of megabytes to gigabytes. Iceberg's hidden partitioning derives partitions from columns, such as days(updated_at), so queries filter on the timestamp and still prune. Partitioning by a high-cardinality column like customer_id creates millions of tiny files, the classic small-files problem in a new costume.
Clustering within partitions is the second lever: sort or Z-order files by the columns you filter on, so the per-file min and max statistics become narrow and the planner can skip most files. File size is the third: target roughly 256 MB to 1 GB per data file. Small files multiply metadata, planning time and request costs; very large files reduce parallelism and make row-level rewrites expensive.
CREATE TABLE lake.sales.orders (
order_id BIGINT,
customer_id BIGINT,
status STRING,
amount DECIMAL(12,2),
updated_at TIMESTAMP)
USING iceberg
PARTITIONED BY (days(updated_at))
TBLPROPERTIES (
'write.target-file-size-bytes'='536870912',
'write.delete.mode'='merge-on-read',
'write.update.mode'='merge-on-read',
'write.merge.mode'='merge-on-read');
Ingestion: appends, upserts and change data capture
Append-only events are easy: each micro-batch commits an append, and appends rarely conflict. Mutable entities such as orders arrive as change data capture and must be applied as upserts and deletes.
There are two ways to apply a change to an immutable file. Copy-on-write rewrites every data file containing a changed row: reads stay fast, but a small update to a large file costs a large rewrite. Merge-on-read writes small delete files (or log files in Hudi) recording which rows are gone, plus new files with the new versions, and readers merge them at query time: writes are cheap, reads get slower as delete files pile up until compaction folds them in. High-frequency CDC almost always wants merge-on-read plus scheduled compaction; tables updated nightly can use copy-on-write.
MERGE INTO lake.sales.orders t
USING (
SELECT * FROM (
SELECT c.*, row_number() OVER (PARTITION BY order_id ORDER BY change_lsn DESC) AS rn
FROM staged_changes c)
WHERE rn = 1
) s
ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'D' THEN DELETE
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED AND s.op != 'D' THEN INSERT *;The deduplication step matters: a micro-batch can contain several changes for one key, and MERGE fails or misbehaves when one target row matches more than one source row. Keep the source log sequence number in the table so a replayed batch can be detected, and make the job idempotent: if it crashes after committing but before recording its offset, the rerun must not apply the batch twice.
Maintenance is part of the architecture
Tables accumulate three kinds of debris, each needing a scheduled job. Small files and delete files come from streaming writes and merge-on-read; compaction rewrites them into target-sized files. Old snapshots keep every replaced file alive for time travel; snapshot expiry drops snapshots past a retention window so their unreferenced files can be deleted. Orphan files come from failed writers that wrote data and never committed; an orphan cleanup deletes files no snapshot references, and must only touch files older than the longest-running write, or it will delete a live job's uncommitted output. In Iceberg on Spark these are stored procedures:
CALL lake.system.rewrite_data_files(table => 'sales.orders',
where => 'updated_at >= current_date - INTERVAL 2 DAYS');
CALL lake.system.rewrite_manifests('sales.orders');
CALL lake.system.expire_snapshots(table => 'sales.orders',
older_than => TIMESTAMP '2026-09-25 00:00:00', retain_last => 50);
CALL lake.system.remove_orphan_files(table => 'sales.orders',
older_than => TIMESTAMP '2026-09-29 00:00:00');Delta Lake has OPTIMIZE and VACUUM; Hudi has compaction and cleaning services. Whatever the format, schedule maintenance per table from the write rate, and remember that compaction is itself a commit that can conflict with writers: compact partitions the streaming job is no longer writing.
Worked example: sizing a CDC table
Suppose an orders table receives about 200 GB of changes a day, applied by a streaming MERGE every minute with 8 writer tasks. That is 1,440 commits a day, each writing about 139 MB, so roughly 17 MB per file and about 11,500 new data files a day before delete files. Left alone, a month produces over 340,000 files and 43,000 snapshots, query planning slows to seconds, and storage grows with every rewritten row.
The plan: partition by day of updated_at, so each day's activity is a bounded partition. Compact every hour, only on partitions older than the current hour, targeting 512 MB files: 200 GB becomes about 400 files a day. Expire snapshots older than seven days, which bounds time travel to a week and lets replaced files be deleted. Run orphan cleanup daily with a three-day age threshold. Then check the numbers in the metadata tables: files per partition, average file size and delete files per data file. If the averages drift down, the compaction schedule is losing to the write rate.
Also ask whether one-minute commits are needed: five-minute commits cut snapshots and small files fivefold.
Migrating from Hive tables
There are two ways to migrate a Hive table. An in-place migration creates table-format metadata that points at the existing Parquet files without copying them; it is fast and cheap but inherits the old layout, small files and all. A rewrite with CREATE TABLE AS SELECT copies the data into a new table with a better partition spec and file sizes; it costs a full read and write but fixes the layout once. A common path is in-place first, so consumers can switch quickly, then compaction and partition evolution to fix the layout over time.
Cut over writers before readers, never let directory-style writers touch a migrated table, since files they drop are invisible to the metadata, and keep the old table read-only until every consumer has moved. The Hive side of this is covered in Hive and Iceberg tables.
Failure modes and trade-offs
- Lost commits: a catalog or store without a true compare-and-swap lets two writers both commit; one disappears. Use a supported catalog.
- Metadata explosion: frequent commits without manifest rewriting and snapshot expiry make planning slower than the query itself.
- Orphan cleanup deleting live data: an age threshold shorter than the longest write job removes uncommitted files that are about to be committed.
- Conflict storms: compaction and a streaming MERGE rewriting the same partitions keep failing each other's commits. Separate them in time or by partition.
- Engine feature skew: test every reader before enabling a new table feature on a shared table.
- The trade-off itself: a lakehouse gives open storage and many engines at the price of operating maintenance that a warehouse hides.
What to do next
- Pick one catalog that every engine you run can commit through, and confirm it performs an atomic swap.
- Inventory your largest tables: files per partition, average file size and snapshot count.
- Set a partition spec that yields partitions of hundreds of megabytes or more, and a target file size.
- Choose copy-on-write or merge-on-read per table from its update frequency.
- Schedule compaction, snapshot expiry and orphan cleanup per table, with safe age thresholds.
- Make every ingestion job idempotent and deduplicate each batch before MERGE.
- Migrate one Hive table end to end, writers first, before planning the rest.