skip to content

Why does Delta Lake write Parquet checkpoint files into _delta_log?

level: middleimportance: should knowfreq 52%

answer

  1. replaying every commit does not scale
  2. a summary of state, not of one change
  3. written every so many commits
  4. one small file says where the newest is
  5. expiring old commits depends on it

basics

~20 s

Checkpoints spare a Delta reader from replaying the whole JSON history: one Parquet file holds the complete set of live actions at a version, so the reader loads it plus the handful of commits written after it, found via the _last_checkpoint file.

solid answer

~40 s

Reconstructing a Delta table means replaying `add` and `remove` actions from `_delta_log`. On a table with 200,000 commits that would mean 200,000 small reads on object storage per query — unworkable. So every tenth commit by default (`delta.checkpointInterval`) Delta writes `<version>.checkpoint.parquet`: a columnar file containing the *entire* state at that version — all live `add` actions with their statistics, tombstones still inside the retention window, and the current `metaData`, `protocol` and `txn` actions. Alongside it, `_last_checkpoint` is a tiny JSON file naming the newest checkpoint's version and size, so a reader can go straight there instead of listing the directory. Startup becomes: read `_last_checkpoint`, load that Parquet, replay at most a few JSON commits on top. Checkpoints also gate log cleanup — commit files are only expired once a checkpoint covers them.

code

text · 11 lines
text
_delta_log/
  00000000000000000018.json
  00000000000000000019.json
  00000000000000000020.checkpoint.parquet
  00000000000000000020.json
  00000000000000000021.json
  00000000000000000022.json
  _last_checkpoint

_last_checkpoint contents:
{"version":20,"size":1743}

go deeper

for a junior

Know that Delta periodically writes a Parquet checkpoint summarising the table state so readers do not replay every JSON commit, and that _last_checkpoint points at the newest one.

for a middle

Explain exactly what a checkpoint contains — live adds with stats, retained tombstones, current metaData and protocol — and how a reader combines it with the newer JSON commits to build a snapshot.

for a senior

Recognise the production symptom: slow query planning on a high-commit-rate table. Check whether checkpoints are being written at all, whether log cleanup is running, and push back on fixing it by tuning the interval alone.

for a principal

Set the commit-rate and retention policy across the platform: checkpoint cost scales with live file count, time-travel depth is bounded by log retention, and both are contracts you make with downstream consumers.

## The problem checkpoints solve A Delta table's current state is the replay of every action ever committed to `_delta_log`. That is correct but not scalable. A streaming job committing every ten seconds writes 8,640 commit files a day; after a month the log holds a quarter of a million small JSON objects. Replaying them means a quarter of a million round trips to object storage before the first row is read, and object-store LIST calls are paginated and slow. Query planning would dominate query time. ## What a checkpoint is Every `delta.checkpointInterval` commits — 10 by default — the writer that produces that version also writes a checkpoint: a Parquet file named for the version, such as `00000000000000000010.checkpoint.parquet`, sitting in `_delta_log` next to the JSON files. It is not a delta; it is the whole state. Its rows are actions, with a column per action type, and it contains: - every `add` action still live at that version, with `path`, `partitionValues`, `size` and the `stats` used for file skipping; - `remove` tombstones that are still within the tombstone retention window, so a reader of an older version and the cleanup job both still see them; - the current `metaData` (schema, partition columns, properties) and `protocol` (required reader/writer versions and features); - `txn` actions recording the last committed version per streaming application id. Because it is Parquet, it is also columnar and compressed: a reader that only needs paths and partition values does not pay for the statistics columns. ## _last_checkpoint Finding the newest checkpoint by listing `_delta_log` would reintroduce the LIST cost the checkpoint was meant to remove. So Delta maintains `_delta_log/_last_checkpoint`, a small JSON file holding the version of the most recent checkpoint and its size in actions (plus part count for split checkpoints). A reader opens that single object, jumps directly to the named checkpoint, and then reads only the JSON commits with higher version numbers. `_last_checkpoint` is an optimisation, not a source of truth. If it is missing, stale or unreadable, a correct client falls back to listing the log and finding the highest checkpoint itself — slower, but never wrong. Nothing in it can make a reader see incorrect data. ## Very large tables A table with millions of live files produces a checkpoint too large to be one convenient object, so Delta supports splitting a checkpoint into multiple part files, with `_last_checkpoint` recording how many parts to expect. Newer Delta versions add a v2 checkpoint layout, gated by a table feature, where a top-level checkpoint file references sidecar files holding the bulk of the `add` actions — which makes writing a checkpoint on a huge table incremental rather than a full rewrite of the file list every time. ## Checkpoints gate log cleanup Delta keeps commit history for `delta.logRetentionDuration`, 30 days by default. Expired commit files are removed during metadata cleanup, but only when a checkpoint at or after them exists — otherwise deleting a JSON file would break the replay chain and make the table unreadable. This is the direct link between checkpoints and how far back you can time travel: once the commit files for version 40 are cleaned up, `VERSION AS OF 40` fails even though the checkpoint at version 50 is intact, because the exact state at 40 can no longer be reconstructed. ## Operational failure modes The common production symptom is a table whose queries have grown slow to *plan* rather than to scan. Two causes recur. First, checkpointing failing or falling behind — for example a writer that lacks permission to create the checkpoint object, so JSON commits accumulate unbounded and every reader replays thousands of them. Second, an extremely high commit rate: a job committing every few seconds outruns any checkpoint interval in terms of files created, and the honest fix is to commit less often (larger micro-batches) rather than to tune the interval down, since more checkpoints on a large table is itself expensive. The mirror-image mistake is lowering `delta.checkpointInterval` to 1 on a wide table, so every commit rewrites a full state file. On a table with hundreds of thousands of live files that turns a cheap append into an expensive metadata rewrite. ## What checkpoints are not They are not backups, and they are not data — a checkpoint holds file *references*, never rows. They are not required for correctness either: a reader that ignored every checkpoint and replayed all JSON from version 0 would compute exactly the same snapshot, just far more slowly. And they are not a substitute for compacting the data files themselves: a checkpoint makes reading a list of 400,000 tiny files fast, but the query still has to open 400,000 tiny files.

  • What happens if _last_checkpoint is deleted or corrupted?
    The table stays readable. `_last_checkpoint` is a pointer optimisation, so a client that cannot use it falls back to listing `_delta_log` and finding the highest checkpoint version itself. Planning gets slower on a large log, but no data is lost and no incorrect snapshot is produced.
  • A streaming job commits every 5 seconds and query planning has become slow — what would you change?
    Reduce the commit rate before tuning metadata. Larger micro-batch triggers mean fewer commit files for the same data. Confirm checkpoints are actually being written (a permissions failure silently stops them), verify log retention is expiring old commits, and separately compact the small data files, which checkpointing does not address.
  • Does a checkpoint at version 50 let you time travel to version 40 after log cleanup?
    No. A checkpoint reconstructs its own version only. Reaching version 40 requires the commit files up to 40, and once metadata cleanup has expired them past `delta.logRetentionDuration`, that version is unreachable even though later checkpoints are intact.

A checkpoint is the balance line printed at the top of a bank statement: rather than adding up every transaction since the account opened, you start from the stated balance and apply only the entries below it.

saying these in an interview costs you the question

  • Calling a checkpoint a backup of the table's data
  • Believing checkpoints are required for the table to be readable
  • Thinking a checkpoint at version 50 restores any earlier version
  • Assuming checkpointing compacts the small Parquet data files
  • Setting the checkpoint interval to 1 on a table with millions of files

context