Roadmap To Be A Data Engineer / Lesson 16
Smaller and Slower
Gzipped CSV saves storage and makes every query slower. Parquet and codecs.
In June we made the nightly listings export 23.2% smaller. By July the job that reads it needed a bigger cluster. Both facts are true, and they are the same change — because “move the export to CSV” is not one decision, it is three.
All figures below were measured, not asserted, on a deterministic 1,000,000-row × 24-column synthetic
snapshot of a second-hand fashion marketplace’s listings table. The generator is in the appendix of the
artifact; it was extracted back out of the rendered page, re-run in a clean directory, and diffed
key-by-key against the published numbers — 0 differing deterministic keys. Sizes are 10⁶ bytes.
Timings are best-of-three on a 2-vCPU container.
1The scene
03:10. listings_snapshot lands in the lake: one file, 1,000,000 garments, 24 columns.
Until June it was Parquet, written with whatever codec the loader defaulted to. Then a pricing vendor came on board who could only read CSV, and rather than ship two files the drop became one gzipped CSV everybody could consume. One word in the writer. It was reviewed, and the evidence on the ticket was correct: the snapshot went from 108.11 MB to 83.08 MB — 23.2% less storage.
Six weeks later the nightly markdown job — which active listings have been sitting since before 15 March, and what are they worth on the books? — has gone from a job nobody noticed to a job you wait for. The ask on the table is to double the cluster.
Two explanations get offered: the catalogue grew, and the cluster is too small. Neither survives arithmetic.
2The storage dashboard was right
| Format | Bytes on disk | Bytes the job must read (4 of 24 cols) |
|---|---|---|
| Parquet + zstd | 60.62 MB | 12.63 MB |
| Parquet + gzip | 61.90 MB | — |
| CSV + gzip (what we ship today) | 83.08 MB | 83.08 MB |
| Avro + deflate | 85.25 MB | 85.25 MB |
| Parquet + snappy (what we shipped before) | 108.11 MB | 17.81 MB |
| CSV, uncompressed | 441.81 MB | 441.81 MB |
The ranking inverts. Parquet + snappy is the largest of the five compressed files on disk and the second smallest on read. The comparison in the ticket was against the uncompressed CSV that never existed in production; against the file that did exist, storage improved 23.2% and read volume got 6.58× worse.
And the option nobody put on the table: Parquet + zstd, 60.62 MB — 43.9% smaller than the snappy file the team started with, and 27.0% smaller than the gzipped CSV that replaced it. One word in the same writer, in the other direction, winning on both metrics.
3Dial one: layout — what the query has to touch
A Parquet file stores a column at a time. Rows are cut into row groups (10 here); inside each row group every column is its own column chunk, compressed separately; a footer (33,061 bytes) indexes all 240 chunks and carries each one’s min/max. The reader opens the footer, then fetches only what it needs.
The job names four columns → 40 chunks + footer = 12.63 MB, 20.8% of the file. The description
column is 37.4% of the file on its own and the job never touches a byte of it.
CSV has no footer and no index: price_cents is the twelfth field, and the only way to reach field 12 is
to parse fields 1–11. So the whole 441.81 MB is inflated and parsed to produce 23.20 MB of answer —
6.58× more bytes read, 19.04× more bytes decoded, every run, forever. This is Lesson 10’s meter
(bytes scanned is the billing unit); columnar layout is the half of that lever partition pruning didn’t
cover.
Avro does not help: it is also row-oriented, so the job reads all 85.25 MB. Avro’s virtue is the write path (append-only records, schema in the header, clean evolution). Choosing Avro for analytics speed is a category error, not a tuning mistake.
4Dial two: codec — whether more than one worker can read it
Count the gzip members in the file: exactly one. A gzip stream is a single continuous deflate window — byte 40,000,001 cannot be decoded without byte 40,000,000 — so the file has one reader, and no cluster size changes that. Measured directly: inflating it takes 1.148 s on one core (384.9 MB/s). That is a floor. A hundred workers do not cross it, and it is already 44.2× the time Parquet needs for the whole answer.
| File | 1 worker | 2 workers | what the 2nd worker bought |
|---|---|---|---|
| Parquet + snappy | 0.035 s | 0.019 s | ×1.80 |
| Parquet + zstd | 0.047 s | 0.026 s | ×1.82 |
| CSV, uncompressed | 0.879 s | 0.551 s | ×1.59 |
| CSV + gzip | 2.085 s | 1.764 s | ×1.18 |
So H2 (“double the cluster”) is refused by arithmetic before anyone provisions anything: you would double the bill to buy 18%.
Compression is two decisions wearing one name. How small is what the storage dashboard sees. How splittable decides your maximum parallelism — and it is a property of the file, not of the cluster. gzip and plain DEFLATE are not splittable;
bgzip, multi-frame zstd, bzip2 and LZ4 frames are. In Parquet the question doesn’t arise: the row group is the split unit.
5Dial three: schema — who decides what the bytes mean
Parquet begins and ends with the four bytes PAR1 and carries its column types in the footer. Avro begins
4f 62 6a 01 and the next thing in the file is the schema as JSON. Our CSV begins 6c 69 73 74 — the
letters list, the start of the header row. No magic number, no version, no types. The types are
whatever the reader guesses.
| Column | Parquet (the truth) | DuckDB, from the CSV | pandas, from the CSV |
|---|---|---|---|
| listing_id | BIGINT | BIGINT | int64 |
| ean | VARCHAR | VARCHAR | int64 |
| size_label | VARCHAR | VARCHAR | str |
| warehouse_bin | VARCHAR | VARCHAR | str |
| price_cents | INTEGER | BIGINT | int64 |
| listed_at | TIMESTAMP | TIMESTAMP | str |
| sold_at | TIMESTAMP | TIMESTAMP | str |
| is_test | BOOLEAN | BOOLEAN | bool |
| currency | VARCHAR | VARCHAR | str |
Both readers behave correctly and neither reports an error. DuckDB’s sniffer keeps ean as text; pandas
sees digits and returns int64. 419,996 EANs lose their leading zero (42.0% of the file), and the join to
the supplier catalogue quietly returns 580,004 of 1,000,000 rows instead of 1,000,000 of 1,000,000.
And here is why nobody noticed for six weeks: run the markdown job under both schemas and the answer is identical to the cent — 282,063 listings, €34,225,605.43. Prices and counts survive type-guessing, because a number meant to be a number stays a number. What breaks is the other kind of column — identifiers that happen to be written in digits: EANs, order references, postcodes, shoe sizes, German bank sort codes. Nothing on a revenue dashboard watches those.
A format that loses types will still get your money right. It loses your keys — and keys fail as missing rows in someone else’s report, not as an error in yours.
6Three dials, one word
| Dial | The question | What it decides | Measured here |
|---|---|---|---|
| Layout | rows or columns? | how much of the file a query must touch | 12.63 MB vs 83.08 MB — 6.58× |
| Codec | compressed how? | bytes on disk, and separately whether >1 worker can read it | 60.62–108.11 MB; a hard 1-reader ceiling on gzip |
| Schema | written down, or guessed? | whether two teams reading the same bytes get the same table | 1,000,000 vs 580,004 rows joining, from one file |
A seventh shape of failure: the bundled
The series has been collecting these — 11 the constant, 12 the drift, 13 the unwatched, 14 the reversible, 15 the unreproduced. This one is none of them. Someone was watching, the thing they watched did improve, and the evidence on the ticket was correct. The failure is that one name — “CSV” — covered three decisions, so approving the one with evidence silently approved the two without.
A review can only accept or reject the bundle, and only the named decision ever gets a number. It generalises well past file formats: “let’s move it to S3”, “let’s upgrade the client library”, “let’s just send it as JSON”. The countermeasure is cheap and entirely social: make the reviewer name the dials before quoting the number.
So what should the drop actually be?
The vendor’s requirement was real and never in dispute — it just never required changing the source of
truth. Keep Parquet + zstd as the table and generate the vendor’s csv.gz from it nightly as a derived
file: 13.48 s of writer time and 35.59 MB more total storage than the snappy file we had. In exchange the
job reads 12.63 MB instead of 83.08 MB, every night, forever. Two files is the answer, and storage is the
cheap axis to spend.
One more dial you only get with Parquet — if you use it
The footer stores min/max per chunk, so a reader can skip whole row groups. Ours skips none: the snapshot is
written in listing-id order, so the filter listed_at < 2026-03-15 touches 10 of 10 row groups. Sort
the table by listed_at before writing and the same filter touches 5 of 10. Statistics are only as
useful as the ordering they describe — Lesson 10’s pruning lever, one level down.
7Ask the team at standup
- For our biggest table — what fraction of its bytes does the job that reads it most often actually need? If nobody knows, that is the answer. Here it was 20.8%.
- Which of our files can be read by more than one worker? Anything gzipped or plain-DEFLATE is one reader per file, whatever the cluster. Ask what the codec is, not just whether it’s compressed.
- Where does the schema of that file live — in the file, or in whichever library opens it? Follow up: which columns are identifiers written in digits, and what happens to them?
- The last time we changed a format, what was the “before” file we compared against? A ratio against a file that never ran in production is a true number about nothing.
- Which of our tables is sorted on the column its queries filter on? Statistics without ordering skip nothing.
8Hands-on (25 minutes)
Run the generator from the artifact appendix (pip install pyarrow duckdb fastavro pandas numpy; ~4 min,
~1.1 GB of scratch). Your seconds will differ; your byte counts will not.
- Move the codec dial. Compare the three Parquet codecs on size and write time (1.59 s for zstd against 34.87 s for gzip). Decide your default, and why it is not “the smallest”.
- Move the layout dial. Add
"description"toJOB_COLSand re-run; watch Parquet’s advantage collapse. Then useparquet_col_bytesto find the one column you’d move to a separate table. - Break the schema dial on purpose. Read
out/listings.csvwith pandas and with DuckDB, compareean, then write the one argument that would have prevented it — and find where that argument is missing in your own pipelines. - Use the statistics. Point the job at
listings.sorted.zstd.parquetand re-rungroups_touched. - The real exercise. Take one file drop we actually own. Name its three dials out loud and say who chose each. If any answer is “it was the default”, you have found this week’s ticket.
9Takeaway
“Which file format?” is three questions — what a query has to touch, whether more than one worker can touch it, and who decides what the bytes mean. A change that answers only one of them will still be graded on that one, and pass.
Vocabulary
- columnar layout — values of one column stored together, so a reader can fetch a column without reading the rows (Parquet, ORC). CSV, JSON lines and Avro are row-oriented.
- row group / column chunk / footer — Parquet’s three units. Ours: 10 × 24 = 240 chunks, 33,061-byte footer.
- projection pushdown — reading only the columns the query names. Worth 6.58× here.
- predicate pushdown / row-group skipping — using footer min/max to skip chunks. Worthless without ordering.
- dictionary & run-length encoding — why
currencycosts almost nothing in Parquet and 4 bytes/row in CSV. - codec — none / snappy / gzip / zstd / LZ4 / brotli. Independent of layout; decides splittability.
- splittable — whether a file can be divided into independently decodable ranges. gzip: no. Parquet row groups: yes.
- compression ratio vs bytes read — the storage bill’s number vs the query bill’s number. Not the same, and they don’t always move together.
- schema-on-write vs schema-on-read — types fixed at write (Parquet, Avro) vs inferred at open (CSV, JSON). Schema-on-read means the schema is a property of the reader.
- type sniffing — a reader inspecting some rows to guess types. Data-dependent by construction, which is why two libraries disagree about one file.
Next: table formats and the lakehouse — Iceberg, Delta, and what happens when you need to delete one row from an immutable file.