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.