Skip to content
Reliable Data Engineering
Lesson
Open in the interactive app

Storage and Open Table Formats

Most lakehouse performance and cost problems are file layout problems. Know how bytes sit on disk and you can answer half of all “this query is slow” questions.


1. Row vs columnar formats

Row (CSV, JSON, Avro)Columnar (Parquet, ORC)
LayoutAll columns of a row togetherAll values of a column together
Best forWriting events, full-row reads, streaming (Avro)Analytics: read few columns over many rows
CompressionWeak (mixed types)Strong (similar values, dictionary/RLE)
Predicate pushdownNoYes (min/max stats, dictionary, bloom filters)

Rule: Avro/Protobuf on the wire (Kafka), Parquet at rest (lake).

2. Parquet internals

flowchart TB
    F["Parquet file (e.g. 512 MB)"] --> RG1["Row group 1 (~128 MB)"]
    F --> RG2["Row group 2"]
    F --> FT["Footer: schema + per row-group,<br/>per column stats (min, max, null count)"]
    RG1 --> C1["Column chunk: user_id"]
    RG1 --> C2["Column chunk: country"]
    RG1 --> C3["Column chunk: amount"]
    C2 --> P1["Pages (~1 MB)<br/>dictionary + RLE / bit-packed values"]

How a query like SELECT sum(amount) FROM t WHERE country = 'DE' reads less:

  1. Column pruning: reads only country and amount chunks.
  2. Row-group skipping: footer stats say a row group has country min=‘AT’, max=‘CH’ → skip it.
  3. Dictionary filtering: if ‘DE’ isn’t in the dictionary page, skip.
  4. Encoding: low-cardinality columns become tiny with dictionary + run-length encoding.

That’s why sorting/clustering data by commonly filtered columns makes min/max stats selective and queries 10–100× faster.


3. Why table formats exist

Plain Parquet folders have no atomic commits (readers see half-written jobs), no updates/deletes without rewriting partitions, no schema enforcement, and slow file listing on object storage.

Table formats add a metadata layer that lists exactly which files form each table version.

Delta Lake

flowchart LR
    subgraph log["_delta_log/"]
        J0["00000.json<br/>add f1, f2 · schema"]
        J1["00001.json<br/>add f3"]
        J2["00002.json<br/>remove f1 · add f4 (MERGE)"]
        CP["00010.checkpoint.parquet<br/>(state snapshot every N commits)"]
    end
    J0 --> J1 --> J2 --> CP
    J2 --> V["Version 2 = {f2, f3, f4}"]

Apache Iceberg

flowchart TB
    CAT["Catalog: table → current metadata pointer"] --> MD["metadata.json<br/>schema, partition spec, snapshots"]
    MD --> SN["Snapshot S3"]
    SN --> ML["Manifest list"]
    ML --> M1["Manifest A<br/>files + column stats"]
    ML --> M2["Manifest B"]
    M1 --> D1[data-1.parquet]
    M1 --> D2[data-2.parquet]
    M2 --> D3[data-3.parquet]

Apache Hudi

Comparison

DeltaIcebergHudi
MetadataOrdered JSON log + checkpointsSnapshot → manifest list → manifestsTimeline + file groups
StrengthSpark/Databricks integration, simplicity, CDF, liquid clusteringEngine neutrality, huge tables, hidden partitioning, partition evolutionStreaming upserts, MoR, indexing
Row-level deletesDeletion vectorsPosition/equality deletes (v2), DVs (v3)Native (MoR)
InteropUniForm exposes Iceberg metadataREST catalog standardXTable for interop

Interview answer: “All three give ACID, time travel and schema evolution on object storage. I’d choose based on the engine ecosystem: Delta if we’re Databricks-centric, Iceberg if multiple engines (Snowflake, Trino, Flink) must read and write the same tables. Interop layers like UniForm and XTable are narrowing the gap.”


4. Partitioning

Partitioning = physical directories by column value → whole-directory pruning.

Guidelines:

5. The small-files problem

Causes: streaming with frequent triggers, over-partitioning, many parallel writers, high spark.sql.shuffle.partitions on small data.

Symptoms: slow queries (per-file open + task overhead), slow planning, driver OOM listing files, high object-storage request costs.

Fixes:

FixHow
CompactionOPTIMIZE table (bin-packing to ~1 GB) on a schedule
Optimized writesdelta.autoOptimize.optimizeWrite = true (shuffle before write to produce right-sized files)
Auto compactiondelta.autoOptimize.autoCompact = true (small compaction after writes)
Fewer output partitionscoalesce(n) / AQE coalescing (spark.sql.adaptive.coalescePartitions.enabled)
Longer triggers1-minute instead of 1-second micro-batches if latency allows
Coarser partitioningmonth instead of day/hour

6. Data skipping: Z-order vs liquid clustering

flowchart LR
    subgraph Before["Random layout"]
        A1["file1: ids 1–9000"]
        A2["file2: ids 3–8800"]
        A3["file3: ids 7–9100"]
    end
    subgraph After["Clustered by id"]
        B1["file1: ids 1–3000"]
        B2["file2: ids 3001–6000"]
        B3["file3: ids 6001–9100"]
    end
    Before -->|"OPTIMIZE ... ZORDER BY (id)<br/>or CLUSTER BY (id)"| After

WHERE id = 4242 → before: all 3 files overlap → read all. After: only file2’s range matches → read one.

Hive partitioningZ-ORDERLiquid clustering
MechanismDirectoriesSpace-filling curve rewrite during OPTIMIZEIncremental clustering by keys, managed by the table
Columns1–2 low-cardinalitySeveral, high-cardinality OK1–4 keys, high-cardinality OK
Change keysRewrite tableRe-OPTIMIZE (full rewrite)ALTER TABLE ... CLUSTER BY (...), no full rewrite
Incrementaln/aNo (rewrites clustered data)Yes
RecommendationOnly for very large tables with a clear date filterLegacyDefault for new Delta tables
CREATE TABLE silver.events (...) CLUSTER BY (event_date, customer_id);
OPTIMIZE silver.events;   -- incrementally clusters new data

7. Maintenance operations cheat-sheet

OperationPurposeGotcha
OPTIMIZECompact + clusterCosts compute; schedule off-peak or rely on predictive optimization
VACUUMDelete files no longer referenced (older than retention)Breaks time travel beyond retention; never set retention below the longest-running reader/stream
ANALYZE TABLE ... COMPUTE STATISTICSBetter join planning
Checkpoint / log cleanupKeep log smallAutomatic in Delta
Iceberg expire_snapshots, rewrite_data_files, rewrite_manifestsEquivalent maintenanceOrphan file cleanup separate

8. Interview questions

A query on a 5 TB Delta table filtering on customer_id takes 4 minutes. How do you speed it up?

Check the plan/metrics: files scanned vs total (is skipping happening?). Likely data isn’t clustered on customer_id, so min/max stats overlap in every file. Apply liquid clustering (or Z-order) on customer_id, run OPTIMIZE, verify files pruned. Also check the small files count, whether stats are collected on that column (first 32 columns by default; move it or configure dataSkippingNumIndexedCols/dataSkippingStatsColumns), and whether the filter is wrapped in a function (prevents pushdown).

How does Delta guarantee atomicity on S3, which has no rename?

Commits are new files named by version (000…042.json); a commit succeeds only if no other writer created that version first. That needs a put-if-absent primitive: S3 now supports conditional writes, and historically a DynamoDB-based LogStore or the Databricks commit service coordinated it. Readers only consider committed log entries, so partially written data files are invisible.

Two jobs MERGE into the same Delta table at the same time. What happens?

Optimistic concurrency: both read version N; the first commits N+1; the second detects a conflict and checks whether the files it read were modified. If they touched disjoint partitions/files (and predicates prove it), it retries automatically and succeeds; otherwise it fails with ConcurrentAppendException/ConcurrentUpdateException. Mitigations: partition/cluster so writers touch disjoint data, include the partition predicate in the MERGE condition, row-level concurrency with deletion vectors, or serialise writers.

When would you pick Merge-on-Read over Copy-on-Write?

High-frequency updates/upserts with latency-sensitive writes (CDC every minute) where read latency can tolerate merging logs, with compaction scheduled. Copy-on-write for read-heavy tables with infrequent batch updates.