CertSafari
    Databricks Certified Associate Developer for Apache Spark· Lessons

    Domain 5 · Lesson 25/32

    Structured Streaming Fault Tolerance: Checkpoints and Exactly-Once

    Explain the Structured Streaming engine in Spark, including its functions, programming model, micro-batch processing, exactly-once semantics, and fault tolerance mechanisms.

    11 min read
    3.12% of exam
    5 sources
    Published 3 Oct 2026
    Docs as of 30 Sep 2026

    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 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.

    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?

    Sources1234

    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.

    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?

    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)

    Sources543

    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.

    Common changes between runs that reuse the same checkpoint
    ChangeRestart 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 maxOffsetsPerTriggerAllowed
    Change the trigger intervalAllowed
    Change the number or type of input sourcesNot allowed by default; source naming lets you add or reorder sources
    Change the subscribed Kafka topicGenerally 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 sinkAllowed; Kafka sees only the new data
    Change the output sink from a Kafka sink to a file sinkNot allowed

    Checkpoint 6 of 6· Check yourself

    Which change can be made to a running streaming aggregation and then restarted from the same checkpoint?

    Sources1

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 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. 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. 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.

      Covered in How exactly-once is achieved, and where it stops

    Sources

    Every claim above is drawn from one of these pages, quoted as it was written on the date shown.

    1. 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. 2.
      “end-to-end fault tolerance with exactly-once processing guarantees using familiar Spark APIs”
      ↩︎ Checkpoints and write-ahead logs
    3. 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. 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. 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

    Spotted a mistake, or was something unclear? Tell us.