What you will be able to do
- Describe what a checkpoint directory stores and how checkpoints and write-ahead logs let a query resume
- Explain how offsets, commits and a transactional sink together give exactly-once processing, and where the guarantee drops to at-least-once
- Judge which query changes are safe when restarting from an existing checkpoint
1.Checkpoints and write-ahead logs
A streaming query runs as a long series of micro-batches, and sooner or later one of them will fail partway through. Structured Streaming promises "end-to-end fault tolerance with exactly-once processing guarantees using familiar Spark APIs". The mechanism behind that promise is the checkpoint: "Checkpoints and write-ahead logs work together to provide processing guarantees for Structured Streaming workloads." The checkpoint is the stream's identity. It records which source data has been read, which micro-batches have reached the sink, and the state that stateful operators have built up.
Checkpoint 1 of 6· Fill the gap
Which option has to be set on this streaming write before the query runs?
(df.writeStream
.option(" ? ", "/Volumes/catalog/schema/volume/path")
.toTable("catalog.schema.table")
)The docs say you must specify the checkpointLocation option before you run a streaming query. path sets the output location, and cloudFiles.schemaLocation is an Auto Loader schema option.
The directory has four parts. Offsets are the source offsets processed in each micro-batch, so a restarted query can resume exactly where it stopped without reprocessing data. Commits record which micro-batches have been committed to the sink. State exists only for stateful queries: it holds metadata about the stateful operator, the state schema, and the checkpointed contents of the state store. Structured Streaming automatically checkpoints that state to fault-tolerant storage and restores it after a restart. Metadata holds the unique query ID, and configuration settings are stored in the offset log.
Because the directory is the query's identity, two rules follow. First, "Each query must have a different checkpoint location. Multiple queries should never share the same location." Second, a checkpoint is the only memory the query has of its progress.
Checkpoint 2 of 6· Match them up
Match each part of a checkpoint directory to what it records
Tap a term, then the definition that fits it.
Offsets record what was read, commits record what reached the sink, state holds what stateful operators remember, and metadata identifies the query.
“Commits: A record of which micro-batches have been committed to the sink, enabling exactly-once semantics.”Source: docs.databricks.com
The write-ahead part is about order. Before a micro-batch runs, the source offsets it will read are planned and recorded in the checkpoint's offset log. Once the batch has been written to the sink, the commit log records it as done. If the query dies in between, the restart finds a batch that was planned but not committed and processes it: "When a query restarts, the micro-batch planned during the previous run processes." That is how the offsets let the query resume exactly where it left off.
Two things have to be true for that replay to be safe. The source must still hold the data at the recorded offsets. For a Delta table source, the streaming query must run at least once within the table's retention window, or it fails. The sink must also tolerate a batch being written again: Delta's transaction log does, and for other sinks you make the write idempotent.
One trap concerns sinks that work without an explicit location. The notebook display() output and the memory sink create a temporary checkpoint if you leave the option out, but these temporary locations do not guarantee fault tolerance or data consistency, and they might not get cleaned up properly. Databricks recommends always setting a checkpoint location, even for these sinks. If the target is a Delta table, you can keep the checkpoint next to the data in a directory such as <table-name>/_checkpoints, because VACUUM skips directories whose names begin with _.
Checkpoint 3 of 6· Exam question
A data engineer starts the following streaming query without specifying a `.trigger(...)` clause: ```python query = (streamingDF.writeStream .format("console") .outputMode("append") .start()) ``` Which best describes how Spark schedules micro-batches for this query?
Correct answer: A — Spark uses the default micro-batch trigger, immediately starting the next micro-batch as soon as the previous one finishes processing
- A. This is correct: when no trigger is specified, Spark uses the default micro-batch trigger, which starts a new batch as soon as the prior one finishes, giving the lowest latency the micro-batch engine can offer without explicit tuning.
- B. This is incorrect because continuous processing is an experimental, opt-in execution mode requiring `.trigger(continuous=...)`; omitting a trigger does not switch the engine into that mode.
- C. This is incorrect because the absence of a trigger clause does not pause the query; Spark falls back to its default trigger behavior rather than waiting for configuration.
- D. This is incorrect because processing a single batch and stopping is the behavior of a one-time trigger such as `Trigger.Once()`, not the default trigger that is used when no trigger is specified.
2.How exactly-once is achieved, and where it stops
Two guarantees are worth separating. At-least-once processing "guarantees every record is processed, but a failure and retry might process some records more than once, which risks duplicates." Exactly-once means every record affects the result as if it had been processed exactly once, even when there are retries.
The checkpoint supports the stronger guarantee from the read side. The offset log fixes which input belongs to each micro-batch, and the commit log records which of those micro-batches reached the sink. After a failure, the restarted query replays the batch it had already planned: "When a query restarts, the micro-batch planned during the previous run processes." If you changed configuration between runs, the changes apply only to the first newly planned batch.
Replay only works if the source can supply the same data again. The recorded offsets are useful only while the source still holds the data behind them. A Delta table source, for example, must be read by the query at least once within its retention window, or the query fails. This is why a replayable source is one of the conditions for exactly-once.
The sink has to do its part too. With Delta Lake, "The Delta Lake transaction log guarantees exactly-once processing, even when there are other streams or batch queries running concurrently against the table." For streaming tables in Lakeflow pipelines, which combine checkpoints with Delta Lake's transactional writes, each micro-batch commits its source offsets and its output together. A retried batch there either fully succeeds or is fully rolled back and retried, and is never partially applied twice. This is why Databricks recommends configuring streaming jobs to restart automatically on failure: the restart itself is safe.
The checkpoint controls what the engine reads and records which batches it considers committed. A retried batch can still write again to a sink that cannot make that write transactional or idempotent. Delta's transaction log closes that gap. For an arbitrary sink, an idempotent write is what makes a repeated batch harmless.
That gap is where the guarantee drops. "Some operations like foreachBatch provide at-least-once rather than exactly-once guarantees." With these operations it is your job to make the processing pipeline idempotent, meaning that writing the same batch twice gives the same result as writing it once. Databricks gives the same advice for low-latency sinks such as message buses and operational databases: design sink operations to be idempotent so that downstream consumers can handle duplicates.
Checkpoint 4 of 6· Check yourself
A query uses foreachBatch to write each micro-batch to an external database. The job fails mid-batch and restarts from its checkpoint. What should the engineer assume?
foreachBatch is at-least-once. The planned batch is replayed on restart, so an external write may happen twice unless it is idempotent.
“Some operations like foreachBatch provide at-least-once rather than exactly-once guarantees.”Source: docs.databricks.com
Checkpoint 5 of 6· Exam question
A team is designing a Structured Streaming pipeline that must guarantee exactly-once end-to-end processing semantics, from source to sink, even if the driver restarts mid-query. Which two conditions are both required to achieve this guarantee? (Choose 2 answers)(Select 2)
Correct answers: A, B — The source must be replayable, so the engine can re-read the same data range identified by a recorded offset after a failure; The sink must be idempotent or support deterministic writes, so reapplying the same micro-batch output does not create duplicate records
- A. This is correct: a replayable source lets Structured Streaming re-read exactly the offset range of a failed or incomplete micro-batch after recovery, which is a prerequisite for reproducing identical output.
- B. This is correct: because the engine can re-execute a micro-batch after a failure, the sink must tolerate that reapplication without duplicating records, which is why idempotent or deterministic sink writes are required.
- C. This is incorrect because exactly-once semantics do not depend on cluster topology; multi-executor clusters with shuffles are fully compatible with exactly-once guarantees when checkpointing and idempotent sinks are used.
- D. This is incorrect because output mode choice is unrelated to exactly-once guarantees; append and update modes can also achieve exactly-once semantics when paired with a replayable source and an idempotent sink.
- E. This is incorrect because `cache()` persists data in executor memory for reuse within a single run and provides no durability across a driver restart, so it does not contribute to exactly-once recovery.
3.Restarting from a checkpoint after changing the query
Recovery assumes the restarted query still matches what the checkpoint describes, which limits what you can edit between runs. The docs use two terms. Allowed means you can make the change, though whether its effect is well-defined depends on the query. Not allowed means the restarted query is likely to fail with unpredictable errors.
Sources are the strictest case. Changing the number or type of input sources is not allowed by default, "because Structured Streaming identifies sources by their position in the query plan." The default can be lifted: if you turn on source naming (source evolution), you can reorder existing sources and add new sources without starting from a fresh checkpoint. Stateful operators come next. State is restored on the assumption that its schema has not changed, so any additions, deletions or schema changes to stateful operations, such as the grouping keys of an aggregation or the columns of a dropDuplicates, are not allowed between restarts.
| Change | Restart from same checkpoint? |
|---|---|
| Add or remove a filter, e.g. sdf.selectExpr("a") to sdf.where(...).selectExpr("a").filter(...) | Allowed |
| Add, remove or modify a rate limit such as maxOffsetsPerTrigger | Allowed |
| Change the trigger interval | Allowed |
| Change the number or type of input sources | Not allowed by default; source naming lets you add or reorder sources |
| Change the subscribed Kafka topic | Generally not allowed; results are unpredictable |
| Change grouping keys or aggregates of sdf.groupBy("a").agg(...) | Not allowed |
| Change the output sink from a file sink to a Kafka sink | Allowed; Kafka sees only the new data |
| Change the output sink from a Kafka sink to a file sink | Not allowed |
Checkpoint 6 of 6· Check yourself
Which change can be made to a running streaming aggregation and then restarted from the same checkpoint?
Adding or removing filters is generally safe. Changing grouping keys alters the state schema, changing the topic is generally not allowed, and adding a source changes the sources' positions in the plan.
“Changes that are generally safe include adding or removing filters, changing rate limits, trigger intervals”Source: docs.databricks.com
Sources1
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.Two streaming queries writing to related targets can share one checkpoint directory to save storage.Why is that wrong?
The checkpoint is a query's identity: its offsets, commits, state and query ID. Every query needs its own location.
Covered in Checkpoints and write-ahead logs
2.The memory sink and display() are fault tolerant because they create a checkpoint automatically.Why is that wrong?
The checkpoint they create automatically is temporary and guarantees nothing. Set a checkpoint location explicitly.
Covered in Checkpoints and write-ahead logs
3.Every Structured Streaming sink is exactly-once, including foreachBatch, because the engine checkpoints offsets and commits.Why is that wrong?
foreachBatch and similar operations are only at-least-once. Exactly-once output needs a transactional sink such as Delta, or idempotent write logic.
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“Checkpoints and write-ahead logs work together to provide processing guarantees for Structured Streaming workloads.”
↩︎ Checkpoints and write-ahead logs“Offsets: The source offsets processed in each micro-batch.”
↩︎ Checkpoints and write-ahead logs“Structured Streaming automatically checkpoints the state data to fault-tolerant storage”
↩︎ Checkpoints and write-ahead logs“This is not allowed by default because Structured Streaming identifies sources by their position in the query plan.”
↩︎ Restarting from a checkpoint after changing the query“If you turn on source naming, you can reorder existing sources and add new sources without starting from a fresh checkpoint.”
↩︎ Restarting from a checkpoint after changing the query“any changes (that is, additions, deletions, or schema modifications) to the stateful operations of a streaming query are not allowed between restarts”
↩︎ Restarting from a checkpoint after changing the query“File sink to Kafka sink is allowed. Kafka will see only the new data. Kafka sink to file sink is not allowed.”
↩︎ Restarting from a checkpoint after changing the query“Each query must have a different checkpoint location. Multiple queries should never share the same location.”
↩︎ Exam trap 1“These temporary checkpoint locations do not ensure any fault tolerance or data consistency guarantees”
↩︎ Exam trap 2“Commits: A record of which micro-batches have been committed to the sink, enabling exactly-once semantics.”
↩︎ Checkpoint“When you delete the files in a checkpoint directory or change to a new checkpoint location, the next run of the query begins fresh.”
↩︎ Prediction“Changes that are generally safe include adding or removing filters, changing rate limits, trigger intervals”
↩︎ Checkpoint - 2.
“end-to-end fault tolerance with exactly-once processing guarantees using familiar Spark APIs”
↩︎ Checkpoints and write-ahead logs - 3.
“When a query restarts, the micro-batch planned during the previous run processes.”
↩︎ Checkpoints and write-ahead logs“When a query restarts, the micro-batch planned during the previous run processes.”
↩︎ How exactly-once is achieved, and where it stops“Design sink operations to be idempotent so that downstream consumers handle duplicates and late-arriving data.”
↩︎ How exactly-once is achieved, and where it stops“Databricks recommends that you always configure streaming jobs to automatically restart on failure.”
↩︎ How exactly-once is achieved, and where it stops“For these operations, make sure that your processing pipeline is idempotent.”
↩︎ Exam trap 3“Some operations like foreachBatch provide at-least-once rather than exactly-once guarantees.”
↩︎ Checkpoint - 4.
“the streaming query must run at least one time within the source table's retention window”
↩︎ Checkpoints and write-ahead logs“The Delta Lake transaction log guarantees exactly-once processing, even when there are other streams or batch queries running concurrently against the table.”
↩︎ How exactly-once is achieved, and where it stops“the streaming query must run at least one time within the source table's retention window”
↩︎ How exactly-once is achieved, and where it stops - 5.
“At-least-once processing guarantees every record is processed, but a failure and retry might process some records more than once, which risks duplicates.”
↩︎ How exactly-once is achieved, and where it stops“Streaming tables use Structured Streaming checkpoints combined with Delta Lake's transactional writes”
↩︎ How exactly-once is achieved, and where it stops“each micro-batch commits its source offsets and its output together”
↩︎ How exactly-once is achieved, and where it stops