A cloud data lake is a simple idea with a complicated implementation: keep all analytical data as open-format files in cheap object storage, and let any number of compute engines read and write it. Storage and compute scale and bill separately, no single vendor owns the data, and a new engine can be pointed at old data without a migration. That separation is the whole appeal.

The complications come from the fact that object stores are not filesystems, the three major clouds differ in ways that matter to correctness, and files alone cannot give you transactions. This article builds the lake from the bottom layer up in a cloud-neutral way, shows what each layer must guarantee, works through a concrete clickstream table with real numbers, and covers the costs and failures teams actually meet. For an AWS-specific build with Glue and Athena, the AWS data lake article goes deeper on that stack.

Advertisement

The five layers and what each one guarantees

Cloud data lake: five layers, one commit pointEnginesSpark, Trino, Flink, warehouse engines, PythonCatalogtable name to current metadata pointer (atomic swap)Table format metadatasnapshots, manifests, schema, partition specData filesParquet or ORC, 128 MB to 1 GB eachObject storeS3, GCS or ADLS Gen2Ingestionbatch and streamingGovernanceidentity, grants, auditReaders resolve the catalog pointer first; writers commit by swapping it.
A lake is layered. Object storage holds bytes, data files hold rows in columnar form, table-format metadata says which files make up a version of a table, the catalog says which metadata is current, and engines do the work.

Object store. Durable, cheap, effectively unlimited storage addressed by bucket and key. It guarantees that a completed write is durable and, on all three major clouds today, that a read after a completed write sees it. It does not provide multi-object transactions.

Data files. Columnar formats such as Parquet store each column separately with statistics per row group, so a query reading three columns out of fifty reads a fraction of the bytes, and min/max statistics let engines skip row groups entirely. The Parquet format article explains the layout.

Table-format metadata. Apache Iceberg, Delta Lake and Apache Hudi record which data files belong to each version of a table, along with schema, partitioning and per-file statistics. This is what turns a directory of files into a table with transactions.

Catalog. Maps a table name to its current metadata and performs the atomic swap that makes a commit visible. It is also where permissions attach.

Engines. Spark, Trino, Flink, cloud warehouse engines and single-node tools such as DuckDB. Engines are replaceable precisely because the layers below them are open.

The storage layer is not a filesystem, and the clouds differ

Tools built for HDFS assumed two things: listing a directory is cheap and consistent, and renaming a directory is atomic. Classic Hadoop and Hive jobs committed output by writing to a temporary directory and renaming it into place. On object storage, neither assumption holds uniformly, and the details vary by provider.

BehaviourAmazon S3Google Cloud StorageAzure ADLS Gen2
NamespaceFlat keys; prefixes only look like foldersFlat by default; hierarchical-namespace buckets availableHierarchical when the namespace feature is enabled
Directory renameNone; copy then delete, per objectCopy then delete in flat buckets; folder rename in hierarchical bucketsAtomic directory rename with hierarchical namespace
Read-after-writeStrong since December 2020StrongStrong
Conditional createSupported since 2024 (If-None-Match)Supported through generation preconditionsSupported through ETag conditions
Throughput scalingPer key prefix; published guidance is at least 3,500 writes and 5,500 reads per second per prefixScales with key-range load over timePer account and partition limits

Two consequences follow. First, any commit protocol that relies on rename is slow on S3, since renaming a directory of ten thousand files is ten thousand copies, and it is not atomic, so a crash leaves half the files moved. Second, a design that works on ADLS Gen2 because rename is atomic can silently lose that property when ported to S3. Modern table formats avoid rename entirely, which is the main reason they exist. The object storage architecture article covers request paths and key design in detail.

Advertisement

The commit problem and how table formats solve it

Imagine two jobs writing to the same table and a third reading it. With plain files in a directory, the reader may list the directory halfway through a write and see partial output; the two writers may each delete files the other depends on; and nobody can say what the table looked like yesterday. These are the problems a database transaction log solves, and table formats bring that log to object storage.

Iceberg's protocol shows the pattern clearly. A writer writes new data files to fresh, unique keys, so nothing existing is overwritten. It then writes manifest files listing those data files with their statistics, a manifest list for the new snapshot, and a new table metadata file. Finally it asks the catalog to change the table's current-metadata pointer from the version it started with to the new one, as a compare-and-swap. If another writer committed first, the swap fails, and the writer re-reads, checks whether the two changes conflict, and retries. Readers resolve the pointer once and read an immutable snapshot, so they never see partial writes.

Delta Lake keeps an ordered log of numbered JSON commit files next to the data and relies on the storage layer creating each commit file only if it does not exist. On S3, before conditional writes existed, that required a separate coordination service for multi-writer safety. S3's 2024 conditional writes make that service unnecessary in principle, but whether your Delta implementation uses them depends on which one and which version, so check before running concurrent writers. Hudi adds a timeline with record-level indexing aimed at upsert-heavy workloads. The shared idea is the same: data files are immutable, and one small atomic operation publishes a new version.

The payoff is concrete: snapshot isolation for readers, safe concurrent writers, time travel to earlier snapshots, schema evolution by column ID rather than position, and partition evolution without rewriting data. The Spark and Iceberg article walks through the write and read paths inside Spark.

The catalog is the control plane

Once commits go through the catalog, the catalog decides what the lake is. Every engine must use the same catalog for the same table, or two engines will commit to different pointers and diverge. This is the most important architectural decision in a multi-engine lake, and it is easy to get wrong by letting each team configure its own.

Options include the Hive Metastore, which many engines speak but which was not designed for table formats; cloud-native catalogs such as AWS Glue; and catalogs implementing the Iceberg REST catalog specification, such as Apache Polaris and several vendor services, which let any compliant engine commit through one HTTP API. Unity Catalog, open-sourced in 2024, is another option in Databricks-centred estates. The criteria that matter are engine coverage, whether commits are genuinely atomic, where access control is enforced, and whether the catalog can hand engines short-lived, scoped storage credentials so users never hold broad bucket access.

from pyspark.sql import SparkSession

spark = (SparkSession.builder
    .config("spark.sql.extensions",
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
    .config("spark.sql.catalog.lake", "org.apache.iceberg.spark.SparkCatalog")
    .config("spark.sql.catalog.lake.type", "rest")            # one catalog for every engine
    .config("spark.sql.catalog.lake.uri", "https://catalog.internal/api/catalog")
    .config("spark.sql.catalog.lake.warehouse", "analytics")
    .getOrCreate())

spark.sql('''
  CREATE TABLE IF NOT EXISTS lake.web.events (
    event_id STRING, user_id BIGINT, event_type STRING,
    event_ts TIMESTAMP, payload STRING)
  USING iceberg
  PARTITIONED BY (days(event_ts))
''')

The configuration above points Spark at a REST catalog. Trino, Flink and Python clients are configured against the same URI, so all of them see the same table and commit through the same compare-and-swap.

File layout: size, partitioning and compaction

Query cost in a lake is dominated by how many files and bytes an engine must open. Each file costs at least one request to open, plus metadata to plan, plus per-file overhead in the engine. Targets of roughly 128 MB to 1 GB per data file are common, with 256 to 512 MB a reasonable default for analytical tables.

Partitioning groups files so queries can skip them. Partition on columns that appear in most filters and have moderate cardinality: day, not second; region, not user ID. Iceberg's hidden partitioning, shown above with days(event_ts), derives the partition from the timestamp so queries filtering on event_ts prune automatically without users knowing a partition column exists. Over-partitioning is the classic mistake: a table partitioned by day and user ID produces millions of tiny files.

Streaming ingestion creates small files by construction, because each micro-batch writes at least one file per partition it touches. Compaction rewrites them into target-size files in the background, and snapshot expiry removes metadata and files no longer referenced by retained snapshots.

-- nightly, per table: compact small files, drop old snapshots, remove orphans
CALL lake.system.rewrite_data_files(
  table => 'web.events',
  options => map('target-file-size-bytes', '536870912'));

CALL lake.system.expire_snapshots(
  table => 'web.events',
  older_than => TIMESTAMP '2026-09-23 00:00:00',
  retain_last => 10);

CALL lake.system.remove_orphan_files(table => 'web.events');

Worked example: a clickstream table

A product ingests 2 TB of click events per day as compressed Parquet, streamed with one-minute micro-batches. Events are partitioned by day, and the writer uses 200 parallel tasks.

  • Without compaction. 200 tasks times 1,440 batches is 288,000 files per day, averaging about 7 MB each. A query over 30 days must plan over roughly 8.6 million files; planning alone can take minutes and the engine spends more time opening files than reading them.
  • With nightly compaction to 512 MB. 2 TB becomes about 4,000 files per day, or 120,000 for 30 days: seventy times fewer objects to open, list in metadata and track.
  • With a sort on user ID during compaction. Per-file min/max statistics on user_id become narrow, so a query for one user's sessions skips most files in each day, reading a few gigabytes instead of terabytes.
  • Retention. Expiring snapshots older than seven days bounds how far time travel reaches and lets the files replaced by compaction be deleted; without expiry, storage grows with every compaction.

The lesson is that ingestion freshness and read efficiency pull against each other. Choose the micro-batch interval from the latency the business needs, then schedule compaction to repair the file layout behind it.

Security and governance across clouds

Access control in a lake has two levels. Storage-level permissions, such as IAM policies, bucket policies and storage ACLs, decide who can read bytes. Catalog-level grants decide who can query a table, a column or a row. If users can read the bucket directly, catalog grants are advisory, because anyone can bypass them by reading Parquet files. The robust pattern is to deny direct bucket access to people and let the catalog or engine vend short-lived, table-scoped credentials.

Across clouds, keep one identity provider as the source of truth and federate it into each cloud through workload identity federation, rather than creating long-lived keys per cloud. Encrypt with provider-managed keys by default and customer-managed keys where regulation requires, and remember that table metadata files contain statistics, including min and max values, that may themselves be sensitive.

The cost model

Lake bills have four parts, and teams usually watch only the first. Storage is priced per gigabyte-month, with cheaper tiers for infrequently read data. Requests are priced per thousand operations; a job that opens eight million small files pays for eight million GETs every time it runs, which is why the compaction above saves money as well as time. Egress is charged when data leaves a region or a cloud; a query engine in one cloud scanning a lake in another pays for every byte scanned, which often outweighs storage costs. Compute is billed by the engine, often by bytes scanned in serverless engines such as BigQuery or by cluster time.

Two rules keep costs predictable: keep compute in the same region as the data it scans most, and replicate curated tables to a second cloud rather than querying across clouds repeatedly.

Failure modes

SymptomCauseFix
Query planning takes minutesMillions of small files or manifestsCompact data files and rewrite manifests
Two engines show different table contentsEngines configured against different catalogsOne catalog per table, enforced in config
Commit failures under concurrent writersOptimistic conflicts on the same partitionsPartition writers apart; retry with backoff
Storage bill grows after compactionSnapshots never expiredExpire snapshots and remove orphan files
Throttling errors on hot dataRequest rate concentrated on few prefixesLet the table format spread keys; back off and retry
Half-written output after job crashRename-based commit on object storageUse a table format, not directory renames
Catalog grants bypassedUsers can read the bucket directlyDeny direct access; vend scoped credentials
Large surprise network billCross-region or cross-cloud scansCo-locate compute; replicate curated tables

Trade-offs: lake, warehouse or both

A proprietary cloud warehouse gives you managed storage, automatic clustering, strong concurrency and one security model, at the cost of lock-in and paying warehouse prices for storage. A lake gives you open formats, many engines and cheap storage, at the cost of operating compaction, catalogs and permissions yourself. Many organisations now run both on the same files: warehouses increasingly read and write Iceberg tables directly, so the lake becomes the system of record and the warehouse one engine among several. Decide based on how many engines you truly need and whether you have the people to run table maintenance.

What to do next

  1. List every engine that reads or writes your lake and confirm they all use the same catalog for each table.
  2. Measure average file size for your five largest tables; schedule compaction for any table averaging under 64 MB.
  3. Enable snapshot expiry and orphan-file removal with a retention that matches your time-travel needs.
  4. Check whether any job still commits by renaming directories, and move it to a table format.
  5. Remove direct bucket read access for human users and move to catalog-vended credentials.
  6. Break last month's bill into storage, requests, egress and compute, and fix the largest line first.
Key takeaway: A cloud data lake is object storage plus open columnar files plus table-format metadata plus one catalog, with engines on top. Object stores are not filesystems and differ across S3, GCS and ADLS Gen2, so commit through a table format's atomic pointer swap, never directory renames. Keep one catalog per table, compact small files, expire snapshots, deny direct bucket access, and co-locate compute with data to control request and egress costs.