Architecting Exactly-Once ETL with Spark Structured Streaming Micro-batches
Learn how to implement exactly-once ETL pipelines using Apache Spark Structured Streaming, focusing on micro-batch architecture, checkpointing, and operational monitoring.
29 Nov 2025, 01:25 UTC

The Latency vs. Reliability Trade-off
The primary challenge in real-time ETL is ensuring that data is neither lost nor duplicated during system crashes, while maintaining a predictable processing delay. In Apache Spark Structured Streaming, the micro-batch model solves this by treating a live stream as a series of small, discrete batches. The critical takeaway for architects is that exactly-once semantics are not automatic; they require a coordinated combination of a replayable source, a deterministic transformation, and an idempotent sink.
The Minimalist Architecture
To implement a reliable streaming pipeline, the smallest suitable design consists of three components: a replayable source (like Apache Kafka), a state-managed transformation layer, and a durable checkpoint store.
// Example configuration for a basic Structured Streaming pipeline
val streamDF = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host:9092")
.option("subscribe", "telemetry-data")
.load()
val processedDF = streamDF
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
.filter("value IS NOT NULL")
val query = processedDF.writeStream
.format("parquet")
.option("checkpointLocation", "s3a://my-bucket/checkpoints/etl-job-01")
.option("path", "s3a://my-bucket/output/telemetry")
.start()
Trust and Data Boundaries
The trust boundary begins at the source connector. Because streaming pipelines often run for weeks without restart, schema drift (unexpected changes in the source data structure) can crash the executor or corrupt the sink. Data must be validated or cast to a known schema immediately after the readStream call. Any record failing this validation should be routed to a "dead-letter queue" (a separate storage path) rather than allowing the entire micro-batch to fail and retry indefinitely.
Operational Guardrails
A streaming job is healthy only if the processing time for a batch is consistently lower than the trigger interval. If the batch duration exceeds the interval, the system develops processing lag, which can lead to memory exhaustion as the backlog grows.
Key Diagnostic Checks
- Batch Duration: Monitor via the Spark UI's "Structured Streaming" tab. If duration consistently climbs, you must either increase cluster resources or optimize the transformation logic.
- Checkpoint Latency: Checkpointing writes the current offset to storage. If using a slow distributed file system, the time spent writing the checkpoint can become the primary bottleneck, limiting throughput regardless of CPU power.
- Watermark Progress: For stateful operations (like windowed aggregations), ensure the
withWatermarkdelay is tuned. Without a watermark, Spark keeps all historical state in memory, eventually causing anOutOfMemoryError.
Failure Modes and Recovery
The system is designed to handle worker node failures automatically. However, the Driver node is a single point of failure. If the Driver crashes, the current batch state is lost. Recovery is only possible if the checkpointLocation is stored on a distributed file system (HDFS, S3, or ADLS) rather than local disk.
To verify recovery, you can manually terminate the Spark Driver during an active stream. Upon restarting the application with the same checkpointLocation, Spark will read the write-ahead logs (WAL) to determine the last successfully committed offset and resume processing from that exact point.
When to Pivot the Design
The micro-batch model typically provides latencies in the range of 100ms to several seconds. You should transition to Continuous Processing mode only if your business requirement demands sub-100ms latency. Note that moving to Continuous Processing sacrifices some of the strong exactly-once guarantees provided by the micro-batch model, as it relies on a different mechanism for fault tolerance.
Design Comparison Table
| Metric | Micro-batch Mode | Continuous Processing |
|---|---|---|
| Latency | 100ms+ | <10ms |
| Fault Tolerance | Exactly-once (with idempotent sink) | At-least-once |
| Throughput | High (optimized for batches) | Lower (per-record overhead) |
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.