Checkpointing for Fault‑Tolerant Stateful Processing in Apache Spark Structured Streaming
Configure a reliable checkpoint directory to preserve state and offsets, enabling exactly‑once processing after driver or executor failures in Spark Structured Streaming.
21 Aug 2025, 19:05 UTC

Problem: State Loss in Stateful Streaming
When a Structured Streaming job uses stateful operators (e.g., mapGroupsWithState, windowed aggregations), the intermediate state lives in memory. If the driver or an executor crashes, that state is lost and the job would either restart from scratch or risk duplicate processing upon recovery.
Minimal Viable Design
A fault‑tolerant stateful design needs three elements:
- A reliable distributed file system that guarantees strong consistency for writes (HDFS, S3 with versioning, ADLS).
- A unique checkpoint directory where Spark writes write‑ahead logs, offset ranges and the state store.
- A stateful operator in the streaming query.
Configuration Example
The following Scala snippet shows a simple stateful aggregation that writes to the console and persists its checkpoint to HDFS.
import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder.appName("CheckpointedAggregation").getOrCreate()
val lines = spark.readStream
.format("socket")
.option("host", "localhost")
.option("port", 9999)
.load()
val words = lines.as[String].flatMap(_.split(" "))
val wordCounts = words.groupBy("value").count() // stateful aggregation
val query = wordCounts.writeStream
.outputMode("complete")
.option("checkpointLocation", "hdfs:///tmp/spark-checkpoints/wordcount")
.format("console")
.start()
query.awaitTermination()
Important notes:
- The
checkpointLocationmust be unique for each streaming query; sharing it between queries leads to metadata collisions. - Output mode
completeis required for aggregations that emit the full updated state. - The underlying DFS must support atomic renames or strong consistency to avoid partial checkpoint files.
Data Boundaries and Trust
Two boundaries are enforced by the checkpoint:
- Offset boundary: The exact offset of the source (e.g., Kafka or socket) is recorded. After a failure Spark resumes from that offset, ignoring any newer data that arrived while the job was down.
- State boundary: The state store (often RocksDB or HDFS‑backed) is written at the end of each micro‑batch. The checkpoint interval determines how often this persists.
Operational Checks and Failure Modes
Monitoring
- Track the size of the checkpoint directory; unbounded growth can occur with high‑cardinality state unless a state timeout is configured.
- Verify that the DFS reports successful writes (e.g., HDFS sync or S3 versioning) for each checkpoint.
Failure Scenarios
| Failure Event | Recovery Behavior | Risk |
|---|---|---|
| Executor crash | Driver reschedules tasks; state is restored from the latest checkpoint. | Brief latency increase while tasks are relaunched. |
| Driver crash | A new driver reads the checkpoint location, restores offsets and state, then resumes. | Potential duplicate processing of the micro‑batch that was in progress when the driver died. |
| Storage unavailable | Checkpoint write fails; the streaming query stops. | Job downtime until the DFS is reachable again. |
| Checkpoint corruption | Job fails to start because metadata cannot be read. | Requires manual inspection; deleting the checkpoint allows a fresh start but loses accumulated state. |
Verification Steps
- Start the stateful streaming query and let it process several batches of data.
- While the query is running, terminate the driver process (e.g.,
kill -9 <pid>) or simulate an executor loss. - Restart the query using the exact same
checkpointLocation. - Observe that the output continues from the last processed offset and that the aggregated counts reflect the pre‑failure values (no reset to zero and no duplicate counts).
Rollback Procedure
Because checkpointing writes to an external directory, resetting the job’s progress requires removing that directory. Deleting checkpointLocation clears all stored state and offsets, causing the job to reprocess all available data from the source. This action should be performed only when a clean restart is intended.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.