Roadmap To Be A Data Engineer / Lesson 17

Lesson 17 Lakehouse Fundamentals §2 Storage (+§8, §6) About 15 min read

Both Copies Were There

A delete job makes a total go up. Why table formats like Iceberg exist.

On the night of 2 September a job whose only purpose is to delete rows made the year-to-date GMV figure go up by €1,023,333.61. That direction is not a coincidence — it is the signature.


103:12, close week

Finance runs the same query every morning of close week: FY2026 year-to-date GMV, 1 January – 31 August 2026. A closed period — the last sale in it happened days ago.

when figure
1 Sep, 07:40 €59,469,941.60 the figure in the close pack
2 Sep, 03:12 €60,493,275.21 same query, same window, +1.721 %
2 Sep, 09:00 €59,468,906.13 re-run, and stable ever since

Three explanations were on the table by 09:30:

  1. Late-arriving sales — ruled out: the window is closed, and every sale lands in its own day’s partition the night it happens.
  2. An FX re-run — ruled out: the marketplace prices and settles in EUR.
  3. The GDPR erasure job (Lesson 14), automated on 1 June, running nightly at 03:00 — ruled out in about four seconds, because a job that removes rows cannot make a total go up.

That third rejection is the lesson. A total that rises while rows are being deleted is not evidence that the delete is innocent — it is evidence that the delete is running right now, and something is counting rows twice.

2The direction is the tell

That night the erasure removed 10 verified subjects, 34 sales rows, €2,363.50 of GMV, of which €1,035.47 falls inside the reporting window.

value
correct move −€1,035.47 34 rows genuinely removed
observed move +€1,023,333.61 988× as large, opposite sign
worst answer available +€3,404,614.95 3,288× the legitimate change

Nothing about the report, the query or the period changed. Only the files under the table changed, and only for nine seconds.

3A directory is not a table

fct_sales is Parquet on S3, Hive-partitioned by sale_date: one file per day, 1,096 files, 3,467,768 rows, 108.16 MB. There is no metadata layer. The table is the answer to “what objects are under this prefix, right now?” — a question with a different answer every millisecond.

To erase 34 rows the job must rewrite every file that contains one, and it does the careful thing:

  1. GET the partition file
  2. PUT a replacement without the subject’s rows, under a new key
  3. verify the replacement is readable and has the expected row count
  4. DELETE the original

Step 3 is what a good engineer adds. Step 3 is also the bug: between the PUT and the DELETE, both objects sit under the prefix, both match the table, and every reader reads both.

This is not an S3 consistency problem. S3 has been strongly read-after-write consistent since December 2020 and behaved perfectly. Each PUT and each DELETE is individually atomic. What is missing is a way to make two object operations one operation. That is a transaction, and a prefix does not have one.

The two obvious ways out both fail:

  • Overwrite the same key. S3 makes that atomic — but you have destroyed the original before you know the replacement is good. You have swapped a correctness window for an unrecoverable one.
  • Delete first, then write. Now the window under-counts. The number goes down instead of up, and the arithmetic that would have exonerated the delete job now wrongly convicts it.

There is no ordering of two independent object operations that has no window.

4Sixty-five states, twenty-seven answers

The job performs 32 writes and then 32 deletes, so the table passes through 65 states. One unchanged query has 65 possible answers that night, and they are computable.

states 65
wrong states 54 (83.1 %)
distinct wrong answers 27
time returning a wrong answer 75.18 % of the 9.67 s window

Two details matter more than the headline:

  • 14 of the 32 rewritten files fall inside this report’s window and 18 do not — so a different report, over a different period, sees a different subset and gets its own private wrong answer.
  • Nobody designed the sequence. It is the order the erasure registry happened to return subjects in; re-sorting the loop by date instead of by customer id produces a completely different set of wrong answers with identical arithmetic behind them.

5Why nothing caught it — the transient

At the peak the table holds 1,128 files instead of 1,096 and 3,576,535 rows instead of 3,467,768 — +3.137 %. Against the last 90 days of daily row growth (0.1122 % ± 0.0108 %) that is z = +278.8. It is not a subtle anomaly. It is the loudest thing that has ever happened to this table.

And nothing will ever see it, because the window is 9.67 seconds — 0.0112 % of a day. A single random read has a 0.0099 % chance of landing inside it, so it takes about 10,079 reads for even odds.

reads / night nights per corrupted read corrupted reads / year expected since 1 June
10 1,008 0.36 0.09
50 202 1.81 0.47
200 50 7.24 1.87
1000 10 36.21 9.33

At 200 reads a night — modest for a warehouse serving dbt models, dashboards and analysts — that is one wrong answer every 50 nights and 1.87 since the automation went live. Exactly one of them was noticed.

Two structural facts make it worse.

  • Freshness monitors are blind here by construction. A subject must be inactive for 31 days before the request is accepted, and identity verification takes 30 more (Lesson 14). So no partition younger than 61 days can ever be touched — and measured across all 30 nights, the youngest partition the job touched was exactly 61 days old. The 60 newest partitions, where every freshness and volume check looks, are never involved.
  • The more useful the table, the more often it is wrong. Exposure scales with reads. This is the only failure mode in the series where being popular makes you less correct.

The eighth failure axis

An event (06–10) is the shape every test is built for. A constant (11) hides because nothing changes. A drift (12) hides because each step is inside the band. An unwatched failure (13) hides on the far side of an instrumented boundary. A reversible one (14) is undone by the next scheduled run. An unreproduced one (15) needs production-shaped data. A bundled one (16) rides inside a change whose stated metric improved.

A transient is none of those: the table is correct before, correct after, and wrong only in between. Every check is a sample. Sensitivity does not help — z = +278.8 did not help. Only sampling rate would, and no sampling rate reaches 0.0112 %.

The corollary is the expensive one: a transient becomes permanent the moment a consumer materialises it. A human who re-runs the query gets the right answer and closes the ticket. A scheduled job that reads the table and writes what it read freezes the wrong number into a table nobody will ever re-derive.

6The bill for a copy-on-write delete

To remove 34 rows the job read, filtered and wrote back 108,801 rows across 32 of 1,096 files — a row amplification of 3,200× — and 3.38 MB of writes.

Over the 30 nights to 2 September: 869 rows erased (0.0251 % of the table), 851 file rewrites over 419 distinct files, 92.37 MB written = 85.64 % of the table’s rows / 85.39 % of its bytes, every month.

Nobody is billed for this in a form anyone reads (Lesson 10). It appears as S3 PUT requests on an invoice line and as a nightly job that takes a little longer each month.

7What a table format changes, and what it does not

Iceberg and Delta both do one structural thing, and everything people like about them follows from it: the table stops being a prefix and becomes a pointer to an explicit list of files. A write stages new objects (invisible, because nothing points at them), writes a new manifest, and then performs a single compare-and-swap on the pointer. A reader resolves the pointer once and reads exactly the files it names.

  1. The window is zero by construction. A reader sees snapshot 0006 or snapshot 0007. There is no third state to land in.
  2. The question becomes answerable. SELECT * FROM fct_sales.history names the commit; FOR SYSTEM_VERSION AS OF 7 re-runs finance’s query against exactly what it saw. On 2 September the team could not reconstruct what the 03:12 read had read, because nothing recorded it. Reproducibility is a property of the storage layer, not of the query.
  3. Deletes and updates become SQL instead of a hand-rolled rewrite loop that each team writes once, badly.

The default that survives the migration

In Iceberg’s TableProperties.java, write.delete.mode, write.update.mode and write.merge.mode all default to copy-on-write. Migrate the table, leave the defaults, and every one of those 92.37 MB a month is still rewritten — you will have bought atomicity and kept the amplification, and the migration ticket will say “moved to Iceberg” and be closed. (Lesson 16’s bundled failure, in a new coat.)

Set write.delete.mode = merge-on-read and the night’s delete becomes a 1,553-byte position-delete file instead of 3.38 MB — 2,179× less. Over 30 nights: 45,016 bytes against 92,365,787 — 2,052×.

But merge-on-read moves the cost to the read side and it accumulates: after 30 nights 419 of 1,096 data files (38.2 %) carry at least one delete, and every scan has to apply them. That is what compaction is for — rewriting those 419 files once a month costs 44.61 MB.

30 nights, 869 erasures bytes written
copy-on-write (the default) 92.37 MB
merge-on-read: delete files only 45.02 KB
merge-on-read + one monthly compaction 44.65 MB (2.07× less)

…and on a schedule you choose rather than at 03:00 in the middle of everybody’s queries. (Iceberg v3 replaces v2’s position-delete files with binary deletion vectors in Puffin files; AWS measured a 73.6 % smaller delete file and 28.5 % faster full-table reads on their benchmark.)

Copy-on-write is right when deletes are rare, bulk and partition-shaped, and reads are hot. Merge-on-read is right when deletes are small, frequent and scattered — which is exactly the shape of a GDPR erasure queue, and exactly the case the default gets wrong.

8Time travel is a liability

Snapshot 0006 is still there. It still lists the pre-erasure files. Those files are still in the bucket. A subject you have erased is one AS OF away.

The retention properties are real, but they are inputs to a job, not a background process:

  • Iceberg history.expire.max-snapshot-age-ms → 5 days; history.expire.min-snapshots-to-keep → 1. Only expire_snapshots applies them, and it runs when you schedule it.
  • Delta VACUUM retains 7 days by default and likewise runs when called.
  • Iceberg history.expire.max-ref-age-ms → Long.MAX_VALUE: any branch or tag someone created (CREATE TAG audit_2026_q2) pins its files forever.

Here: erasure automation went live 1 June 2026. In the 94 nights since, 950 subjects and 2,589 rows have been removed from the live table. With daily appends plus nightly erasure commits and no maintenance job, that is 1,190 unexpired snapshots, and the rows erased on 1 June have been readable for 93 days.

This is Lesson 14’s reversible failure in better clothes: the correction is real, and a mechanism nobody scheduled preserves the thing it corrected. Erasure on a lakehouse is not one commit. It is a commit, plus an expiry, plus an orphan sweep — and only the first of the three happens by itself.

9Ask the team

  1. When we delete or update a row in the lake, what does a query running at that instant see? Walk me through the two object operations and the window between them.
  2. Can we re-run last Tuesday’s finance query against exactly the files it read? If not, what do we say the next time a closed number moves?
  3. What is write.delete.mode set to on our tables — and did anyone choose it, or is it the default?
  4. Who schedules expire_snapshots / VACUUM, and does our DPO think “erased” means gone from the current snapshot or gone from the bucket?
  5. Which of our jobs write what they read? Those are the ones that turn a nine-second glitch into a permanent row.

10Hands-on (~30 min)

  1. Run the generator in the artifact’s appendix. It builds the 1,096-file table in a few seconds and prints all 21 figures used above.
  2. Change the processing order — sort the affected partitions by date instead of by customer id — and re-run. The set of wrong answers changes completely; the arithmetic producing them does not.
  3. Change the reporting window to H1 2026 and to the trailing twelve months. Each window has its own peak error, because each contains a different subset of the 32 rewritten files.
  4. Then the real one: find your largest Hive-partitioned table, find the job that rewrites partitions in it, and time that job. Multiply by your reads per night. That product is roughly how many wrong answers you shipped last year.

11Takeaway

A table is not a set of files. It is a statement about which files count, at a moment. Object storage gives you atomic objects and no way to change two of them at once; a table format supplies the missing sentence — one pointer, swapped once. Time travel, reproducible queries, row-level deletes and safe concurrent writers are not four features, they are four consequences of having that sentence. The costs it does not remove are the ones you have to configure: copy-on-write is still the default, and expiry is still a job someone has to schedule.

12Vocabulary

table format · manifest / manifest list · snapshot · atomic commit (compare-and-swap on the catalog pointer) · snapshot isolation · time travel · copy-on-write vs merge-on-read · position deletes / deletion vectors (Puffin) · compaction · expire_snapshots / VACUUM · orphan files · write amplification · the transient (a failure that is correct at rest and wrong only in between)


Figures computed from a 3,467,768-row deterministic synthetic marketplace table (no RNG; every value from modular arithmetic on an index). The generator ships inside the artifact; it was extracted from the rendered page, run in a clean directory and diffed key-by-key against the published figures — 21 keys, 0 differences. Iceberg defaults verified against core/src/main/java/org/apache/iceberg/TableProperties.java.

Back to top