File Formats & the Lakehouse
Your nightly revenue export is a CSV. It was fine when it was 40 MB. Now it's 900 GB, it takes an hour to download, and the analyst's "quick question" — sum revenue by region — reads the entire file to touch one column. Someone suggests Parquet. Someone else says "just put it in the data lake." A third person says "Delta Lake," which sounds like a place you'd go fishing.
This post settles all three. You'll learn how Parquet actually stores data (column chunks, pages, encodings — not hand-waving), measure real compression differences between SNAPPY, GZIP, and no compression, watch dictionary encoding kick in on a low-cardinality column and fall back on a unique one, build a hive-partitioned directory tree, and then run a real Delta Lake table — create, append, schema evolution, time travel — with plain Java and no Spark engine. By the end you'll know what a lakehouse is made of, and why "the file format" turned out to be one of the highest-leverage decisions in data engineering.
Lab honesty, up front. Everything below that shows a command or an output ran on this lab machine (Temurin JDK 17, parquet-mr 1.13.1, delta-standalone 3.0.0 fetched from Maven Central): the codec size comparison, the footer encoding inspection, the hive directory layout, the Parquet-level schema evolution, and the full Delta lifecycle — create, append, add-column, time travel. Two things did not run and are taught from documented design, clearly labeled: Delta MERGE (a Spark SQL execution feature — delta-standalone has no MERGE engine) and optimistic-concurrency conflict detection (our commits were sequential). Iceberg and Hudi get the honest one-paragraph treatment. Nothing below is presented as a run unless it says so.
Row stores read rows; analytics reads columns
Post 1 stated this precisely and it's worth restating exactly: a row-oriented store (Postgres, MySQL) lays out a whole record contiguously, so a query that touches one column of a wide row still reads the other columns' bytes off the disk (or page cache) to get to it. A columnar store lays out each column contiguously, so a query over one column reads roughly that column's bytes and skips the rest. That's the whole idea — no magic, just layout.
Make it concrete with a toy illustration (not a benchmark — a layout sketch). A table with four columns — id, region, trace, amount — and the query SELECT sum(amount):
Rule of thumb, not a law: columnar wins when queries touch few columns of wide tables — the classic analytics shape. It loses its edge when you need whole rows (point lookups, SELECT * over narrow tables), because reassembling rows from scattered columns costs more than reading them contiguously. Write the query pattern first; the layout follows.
Columnar also compresses better, for a structural reason: values within one column look alike (dates near each other, a handful of region codes, sorted IDs), and compression feeds on similarity. That observation is the entire reason encodings exist — which is where Parquet gets interesting.
Anatomy of a Parquet file
A Parquet file is a small database, not a dumb byte stream. Four levels, top to bottom:
- File. Starts and ends with the magic bytes
PAR1. Self-describing: the schema travels with the data. - Row group. A horizontal slice of rows (default target ~128 MB). The unit of parallel reads — and the unit the writer buffers before flushing.
- Column chunk. One column's data within one row group. This is the physical reason column pruning works: to read
amount, the reader seeks to theamountcolumn chunks and never touches the rest. - Page. The unit of encoding and compression inside a column chunk (default ~1 MB). Each page is encoded, then compressed, independently.
- Footer. Written last, at the end of the file: the schema, the offset of every row group and column chunk, the encodings used — and, crucially, min/max statistics per column chunk. A query with
WHERE amount > 1000can skip an entire column chunk whose max is 500 without reading it. That footer is why Parquet readers open the end of the file first.
Lab 1: writing Parquet from Java — SNAPPY vs GZIP vs none
The lab uses parquet-mr 1.13.1 directly (the same library Spark uses under the hood) with an Avro schema — no Spark engine needed. The dataset: 60,000 identical rows (id, region drawn from 5 values, unique trace string, amount), written three times with different compression codecs:
Schema schema = new Schema.Parser().parse(
"{\"type\":\"record\",\"name\":\"Order\",\"fields\":[" +
"{\"name\":\"id\",\"type\":\"long\"}," +
"{\"name\":\"region\",\"type\":\"string\"}," +
"{\"name\":\"trace\",\"type\":\"string\"}," +
"{\"name\":\"amount\",\"type\":\"double\"}]}");
try (ParquetWriter<GenericRecord> w =
AvroParquetWriter.<GenericRecord>builder(new Path(file))
.withSchema(schema)
.withConf(conf)
.withCompressionCodec(CompressionCodecName.SNAPPY) // GZIP / UNCOMPRESSED
.withDictionaryEncoding(true)
.build()) {
for (long i = 0; i < 60_000; i++) {
GenericRecord r = new GenericData.Record(schema);
r.put("id", i);
r.put("region", REGIONS[rnd.nextInt(REGIONS.length)]);
r.put("trace", "trace-" + String.format("%08d", i));
r.put("amount", rnd.nextDouble() * 1000);
w.write(r);
}
}
Real output, same 60k rows, three codecs (file sizes read from disk after the write):
LAB1 codec=SNAPPY bytes=1044056 file=orders-snappy.parquet
LAB1 codec=GZIP bytes=714360 file=orders-gzip.parquet
LAB1 codec=UNCOMPRESSED bytes=2064352 file=orders-uncompressed.parquet
Read those numbers as measurements of this dataset, not universal constants: on these 60k rows, GZIP produced a file about 68% the size of SNAPPY's and about 35% the size of uncompressed; SNAPPY roughly halved the uncompressed size. The rule of thumb the industry converged on — and the reason SNAPPY is the common default in engines like Spark — is a tradeoff, not a size contest: GZIP typically compresses further but costs more CPU on every read, while SNAPPY is engineered for fast decompression. For data you scan repeatedly, decompression speed usually dominates the decision; for cold archival data, maximum compression wins. Measure on your own data before treating that sentence as law.
Lab 2: encodings — the dictionary kicking in (and giving up)
Compression squeezes bytes; encodings change how values are represented before compression. The important one is dictionary encoding: within a column chunk, each distinct value is stored once in a dictionary, and each row stores a small integer key into it. Five region codes repeated 200,000 times become five strings plus 200,000 tiny keys. The keys themselves are bit-packed, so the per-row cost shrinks to a few bits.
But dictionaries have a budget. If a column has too many distinct values — a unique trace ID per row — the dictionary would be as big as the data, so the writer falls back to PLAIN encoding (values stored as-is) once the dictionary outgrows its page budget (default 1 MB). Both behaviors are visible in one file's footer. I wrote 200,000 rows — enough unique trace values to blow the dictionary budget — then read the footer with ParquetFileReader.readFooter:
ParquetMetadata meta = ParquetFileReader.readFooter(conf, new Path(encFile));
for (BlockMetaData block : meta.getBlocks())
for (ColumnChunkMetaData col : block.getColumns())
System.out.println("column=" + col.getPath()
+ " encodings=" + col.getEncodings()
+ " codec=" + col.getCodec());
LAB2 footer of orders-enc.parquet: rowGroups=1
LAB2 column=[id] encodings=[BIT_PACKED, PLAIN] codec=SNAPPY
LAB2 column=[region] encodings=[BIT_PACKED, PLAIN_DICTIONARY] codec=SNAPPY
LAB2 column=[trace] encodings=[BIT_PACKED, PLAIN] codec=SNAPPY
LAB2 column=[amount] encodings=[BIT_PACKED, PLAIN] codec=SNAPPY
There's the whole story in four lines. region — five distinct values — got PLAIN_DICTIONARY: the dictionary path was used (this parquet-mr version reports it under the older PLAIN_DICTIONARY name; newer versions report RLE_DICTIONARY). trace — 200,000 unique values — shows PLAIN: the dictionary exceeded its budget and the writer fell back, exactly as designed. The numeric columns stayed PLAIN too: with 200k unique values there was nothing for a dictionary to compress.
Rule of thumb, not a law: dictionary encoding pays off on low-cardinality columns (status codes, regions, categories) and is neutral-to-harmful on unique IDs — the writer already knows when to quit, but you're still paying the attempt. This is also why "sorted data compresses better" is real advice: sorting by a low-cardinality column groups equal values into the same pages, which makes both dictionary and run-length encoding dramatically more effective.
Lab 3: partitioning — directories as an index
Encodings and column pruning reduce what you read inside a file. Partitioning reduces which files you read at all — by encoding a column's value into the directory path. The convention is hive-style: column=value directories.
I wrote three small Parquet files into this layout (1,500 rows each, real ls -R):
$ ls -R orders | grep -v '\.crc'
orders:
year=2024
year=2025
orders/year=2024:
month=12
orders/year=2024/month=12:
part-00000.snappy.parquet
orders/year=2025:
month=01
month=02
orders/year=2025/month=01:
part-00000.snappy.parquet
orders/year=2025/month=02:
part-00000.snappy.parquet
Now the query WHERE year=2025 AND month=01 never opens a file listing problem: the engine resolves the filter against the directory names and reads exactly one directory — partition pruning. The partition columns (year, month) aren't stored inside the Parquet files at all; they're recovered from the path. That's the elegance and the trap: it's a free coarse index, but only on the columns you partitioned by, and only at directory granularity.
Two rules of thumb, both frequently violated:
- Partition by what you filter on, at a granularity that keeps files reasonably sized. Partitioning by
yearwhen every query filters bydayprunes nothing; partitioning byuser_idwith millions of users creates millions of tiny directories. - Beware the small-files problem. Each file costs a listing, an open, and a footer read. Ten thousand 12 KB files can easily spend more time on per-file overhead than on data — a pattern that shows up constantly in streaming ingestion (post 6's territory). The standard remedy is compaction: periodically rewrite many small files into fewer large ones, roughly aligned with the storage's preferred block or object size (hundreds of MB is the usual rule-of-thumb neighborhood).
Principle: the cheapest I/O is the I/O you never do. Columnar layout skips the columns you don't need. Encodings shrink the bytes you do read. Partition pruning skips the files you don't need. Every optimization in this post is the same idea applied at a different level of the stack — and it's the lens to use when the next format or engine claims a speedup: ask which I/O it eliminated.
The lakehouse: what Parquet alone doesn't give you
Parquet files on object storage (S3, GCS, Azure Blob) are the modern data lake: cheap, durable, open-format storage that any engine can read. But files alone are missing everything a database gives you:
- Atomicity. A job writing 500 files that crashes at file 499 leaves a half-visible dataset. Readers see garbage mid-job.
- Isolation. Two writers appending to the same directory interleave unpredictably.
- Schema enforcement. Nothing stops a job from writing strings into yesterday's integer column.
- History. Overwrite a file and yesterday's data is gone — no audit, no rollback.
The lakehouse answer: keep the cheap open storage and the Parquet files, and add a table format — a metadata layer that turns a directory of files into a table with transactions. The modern formula, stated exactly: object storage + Parquet + a table format (Delta Lake, Apache Iceberg, or Apache Hudi). Same idea, different trade-offs; Delta is the one we run today.
Notice the layering: engines come and go at the top, but the table on object storage outlives any one of them — that's the point of the open format. (You'll also encounter Trino as the ad-hoc SQL engine over the lake, Hive as the original SQL-on-Hadoop layer underneath many warehouses, and Airflow / Dagster orchestrating the pipelines that write these tables — the ecosystem is wider than any one diagram.)
Lab 4: Delta Lake without Spark — the transaction log is just files
Delta Lake is usually taught as a Spark library, but the core idea needs no engine: a Delta table is Parquet data files plus a _delta_log directory of JSON files, one per table version, each recording the actions that version committed. I used delta-standalone 3.0.0 (the JVM client from Maven Central, no Spark) to run a full lifecycle: create the table (v0), append (v1), add a column (v2), then time-travel back to v0.
The commit code is explicit about what a "transaction" is — a list of actions written atomically as one JSON log file:
DeltaLog log = DeltaLog.forTable(conf, tableDir);
// v0: create the table — protocol + metadata + data files, one atomic commit
List<Action> v0 = new ArrayList<>();
v0.add(new Protocol(1, 2));
v0.add(new Metadata(UUID.randomUUID().toString(), "orders_demo",
"delta-standalone lab table", new Format(),
Collections.emptyList(), Collections.emptyMap(),
Optional.of(System.currentTimeMillis()), DELTA_V1)); // schema: id long, name string
v0.add(new AddFile("part-00000-v0.snappy.parquet", /* path, size, ... */));
log.startTransaction().commit(v0, new Operation(Operation.Name.WRITE), "lab");
// v1: append — just an AddFile action, no metadata change
// v2: schema evolution — a new Metadata action (schema gains `email`) + AddFile
The resulting table directory shows the design with unusual clarity — the log really is just numbered JSON files:
$ ls delta-table delta-table/_delta_log | grep -v '\.crc'
delta-table/:
_delta_log
part-00000-v0.snappy.parquet
part-00001-v1.snappy.parquet
part-00002-v2.snappy.parquet
delta-table/_delta_log/:
00000000000000000000.json
00000000000000000001.json
00000000000000000002.json
And each JSON file is a list of actions — here is v0 (create) and v1 (append) verbatim:
$ cat delta-table/_delta_log/00000000000000000000.json
{"commitInfo":{"timestamp":1791345671028,"operation":"WRITE","isolationLevel":"Serializable",[...]}}}
{"metaData":{"id":"fe5f899e-dd90-4350-9df9-9610d0cb0ad7","name":"orders_demo",
"format":{"provider":"parquet","options":{}},
"schemaString":"{\"type\":\"struct\",\"fields\":[{...id: long...},{...name: string...}]}",[...]}}
{"protocol":{"minReaderVersion":1,"minWriterVersion":2}}
{"add":{"path":"part-00000-v0.snappy.parquet","partitionValues":{},"size":637,
"modificationTime":1791345668583,"dataChange":true}}
$ cat delta-table/_delta_log/00000000000000000001.json
{"commitInfo":{"timestamp":1791345680515,"operation":"WRITE","readVersion":0,[...]}}
{"add":{"path":"part-00001-v1.snappy.parquet","partitionValues":{},"size":619,
"modificationTime":1791345680475,"dataChange":true}}
Read that v1 entry carefully: the append is one action — "add this file" — committed against readVersion: 0. That's the whole ACID trick, and it's worth stating precisely: a Delta commit is atomic because the version-N JSON file either exists or it doesn't — readers resolve "the table" to the latest complete version. Writers use optimistic concurrency: each commit declares the version it read, and if another writer committed first, the loser retries against the new version. (Our lab's commits were sequential, so the conflict path is documented design, not a lab run — but the log entries above show exactly the mechanism it operates on.)
Lab 5: time travel and schema evolution, for real
Because every version's file set is recorded, "the table as of last Tuesday" is just "replay the log up to version N." The lab read the latest version and then time-traveled to v0:
DELTA after v2 (latest): version=2 schema=[id, name, email]
DELTA data file: part-00001-v1.snappy.parquet
DELTA data file: part-00002-v2.snappy.parquet
DELTA data file: part-00000-v0.snappy.parquet
DELTA row: id=4 name=dave email=null
DELTA row: id=5 name=erin email=null
DELTA row: id=6 name=fred email=fred@example.com
DELTA row: id=7 name=gina email=null
DELTA row: id=1 name=alice email=null
DELTA row: id=2 name=bob email=null
DELTA row: id=3 name=carol email=null
DELTA time travel to version 0: version=0 schema=[id, name]
DELTA v0 row: id=1 name=alice
DELTA v0 row: id=2 name=bob
DELTA v0 row: id=3 name=carol
DELTA time travel to version 0: rows=3
Two things to notice. First, time travel works: version 0 still resolves to exactly the 3 rows and the 2-column schema it had at commit time — no snapshots to manage, no backup to restore; old data files are simply never deleted by new commits (a background VACUUM reclaims them later, past the retention window). Second, schema evolution is a metadata action: v2's commit added a metaData action whose schema gained the nullable email column, and rows written before it read back with email=null — the old Parquet files were never rewritten. (The same add-a-nullable-column evolution at the raw Parquet level also ran in Lab 2's companion check: reading a v1 file with an evolved Avro schema filled the new column with its default.)
What about MERGE (MERGE INTO — upsert: update matched rows, insert new ones)? It's the operation that makes Delta feel like a database, and it runs inside Spark SQL's Delta integration — delta-standalone has no execution engine for it, so it did not run in this lab. Conceptually it's still the log doing the work: a MERGE reads the target version, writes new Parquet files for the changed rows, and commits add actions for the new files plus remove actions for the files they replace — one atomic version bump, so readers never see a half-merged table. [NOT RUN — taught from the documented design.]
Iceberg and Hudi: the honest paragraph
Delta isn't the only table format, and this post won't pretend the others got equal lab time. Apache Iceberg and Apache Hudi solve the same problem — ACID transactions over Parquet (or ORC) files on object storage — with different trade-offs in how they organize metadata: Iceberg tracks the table state in snapshot metadata files with manifest lists (designed for huge tables and partition evolution — changing partition schemes without rewriting data), while Hudi was built around streaming upserts with a timeline of instants and offers copy-on-write vs merge-on-read table types (different read/write amplification trade-offs for update-heavy workloads). Same idea — immutable data files plus a metadata layer that makes them transactional — different engineering bets. If your platform standardizes on one, learn its metadata layout the way this post taught Delta's log; the Parquet and partitioning fundamentals transfer unchanged.
Cheat sheet
- Row stores lay out whole records contiguously; columnar stores lay out whole columns contiguously. Rule of thumb, not a law: columnar wins when queries touch few columns of wide tables; it loses its edge on whole-row point lookups.
- Parquet anatomy: file → row groups (parallel-read unit) → column chunks (one column per row group — the unit column pruning skips) → pages (the unit of encoding + compression) → footer (schema, offsets, encodings, min/max stats). Readers open the end of the file first.
- Measured in this lab (60k identical rows): SNAPPY 1,044,056 bytes, GZIP 714,360 bytes, uncompressed 2,064,352 bytes. GZIP usually compresses further; SNAPPY is engineered for faster decompression and is the common default in engines like Spark. Measure on your own data.
- Dictionary encoding replaces repeated values with small integer keys — it pays off on low-cardinality columns. The writer falls back to PLAIN when the dictionary outgrows its page budget (observed in the lab footer:
region→PLAIN_DICTIONARY, uniquetrace→PLAIN). - Hive-style partitioning (
year=2025/month=01) is a free coarse index: partition pruning reads only the directories the filter matches. Partition by what you filter on; watch the small-files problem and compact toward reasonably sized files (hundreds of MB is the usual rule-of-thumb neighborhood). - The cheapest I/O is the I/O you never do. Column pruning, encodings, footer min/max skipping, partition pruning — every win in this post is skipped work at a different level.
- Lakehouse = object storage (S3/GCS/Blob) + Parquet + a table format (Delta/Iceberg/Hudi). Open formats outlive any single engine.
- A Delta table is Parquet files plus
_delta_log/*.json— one JSON file per version, each a list of actions (commitInfo,metaData,protocol,add/remove). Atomicity comes from versioned log files; writers use optimistic concurrency against a read version. - Time travel = replay the log to version N (ran in the lab: v0 → 3 rows, 2-column schema). Schema evolution = a new
metaDataaction; old files are never rewritten, new columns read as null on old rows. - MERGE (upsert) commits
add+removeactions in one atomic version bump — taught from documented design, not run in this lab (needs the Spark SQL engine). - Iceberg (snapshot metadata + manifest lists, partition evolution) and Hudi (streaming upserts, copy-on-write vs merge-on-read) are the same idea with different trade-offs.
- Adjacent tools you'll meet around the lake: Trino (ad-hoc SQL), Hive (legacy SQL layer), Airflow/Dagster (pipeline orchestration).
What's next
You now know how data rests: columnar files, pruned partitions, transactional tables. But so far everything has been still — bounded datasets, written once, read many times. The revenue dashboard from post 1 wants updates every minute, and the small-files warning in the partitioning section was really about a deeper problem: what happens when data never stops arriving?
The next post, Streaming: Spark Structured Streaming, puts the lakehouse in motion: event time versus processing time, watermarks for late data, windowed aggregations, and exactly-once sinks — including how streaming writes land in Delta tables without the small-file explosion.
Field check before you move on: (1) Re-run Lab 1 on your own machine with a dataset of your choosing (a CSV you have lying around works) — which codec wins on size, and does the ranking match this post's? Write down why your data might differ. (2) Take the Lab 2 footer output and explain, in two sentences, why sorting the input by region before writing would shrink the file further. (3) Design a partition scheme for an orders table queried two ways — "last 7 days of all orders" and "all of user X's history" — and argue which query your scheme favors and what the other one pays. (4) Open the _delta_log JSON from the lab and point to the exact line that makes v1 an append rather than an overwrite.
Continue: Java Learning Roadmap 2026
Comments
Post a Comment