Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Streaming checkpoint location

PySpark · Streaming & Newer Spark

Streaming checkpoint location

Hardpyspark-86
structured-streamingcheckpointstate-storeexactly-once

Question

What is stored in a streaming checkpoint, and what happens if you delete it or change the query?

Solution

The checkpoint directory is where a streaming query records its progress: which source offsets it has read, which batches it has committed, and the state of stateful operations such as aggregations and joins. It is how a restarted query knows where to continue.

What is inside

checkpoint/
  offsets/     what range of source data each batch was planned to read
  commits/     which batches finished and were written to the sink
  state/       running aggregates, dedup keys, join buffers
  metadata     the query id

Restart behaviour

If the query stops, whether it crashed or you stopped it, and you start it again with the same checkpoint, it reads the last committed batch and resumes at the next offsets. Together with an idempotent or transactional sink, this gives end-to-end exactly-once results. Rows are neither skipped nor duplicated.

Deleting the checkpoint

Spark forgets everything, and treats the query as new. What happens next depends on the source. With Kafka and startingOffsets=earliest, it reprocesses the whole retained topic, which can duplicate data in a sink that is not idempotent. With latest, it skips whatever arrived while the query was down, so you lose data. State is also lost, so a running count starts from zero. Never delete a checkpoint casually in production.

Changing the query

Many changes are fine: adding a filter, changing a sink path, adjusting the trigger. Others are not compatible with the saved state:

  • Changing the grouping keys or the aggregation of a stateful query.
  • Changing the number of shuffle partitions of a stateful query (state is split by partition).
  • Changing the source or sink type in some cases.

Spark then fails at startup with an error about incompatible state. The fix is a new checkpoint location, and a plan for reprocessing or backfilling the data that the old query handled.

Rules

Use one checkpoint location per query, never shared. Keep it on reliable storage such as S3 or HDFS, not on local disk of a node that can vanish. Plan for its growth: state stores can become large, so use watermarks to let old state expire.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext