How to build real-time pipelines that survive failures, handle late data, and run efficiently in production.
Structured Streaming looks simple in demos. Read a stream, transform, write. Three lines of code and data flows in real time.
But production streaming is a different challenge entirely.
What happens when the cluster restarts? When data arrives three hours late? When a schema change breaks the pipeline at 2 AM? When costs spiral because the stream runs continuously on an expensive cluster?
This guide covers the patterns that separate reliable streaming pipelines from fragile ones. Every concept maps directly to what you practice in Chapter 11: Structured Streaming and connects to the broader pipeline architecture covered across the book.
Checkpointing is the single most important concept in production streaming. Without it, every restart means reprocessing everything from the beginning.
A checkpoint stores exactly where your stream left off: which files were read, which offsets were consumed, what state was accumulated. When the pipeline restarts, it picks up precisely where it stopped.
# Every streaming write MUST have a checkpoint location
query = (
df.writeStream
.format("delta")
.option("checkpointLocation", "/mnt/checkpoints/orders_stream")
.outputMode("append")
.toTable("silver.orders")
)The checkpoint location must be a durable path, not local disk. Use cloud storage (DBFS, S3, ADLS, GCS) so the checkpoint survives cluster termination.
One checkpoint per query. Every streaming write needs its own unique checkpoint path. Sharing checkpoints between different queries causes data corruption.
Never delete checkpoints casually. Deleting a checkpoint resets the stream. It will either reprocess everything or lose track of what was already processed. Only delete when you intentionally want to reprocess.
Organize checkpoint paths consistently. A clear naming convention prevents accidents.
# Consistent checkpoint path convention
checkpoint_base = "/mnt/checkpoints"
# Pattern: {base}/{layer}/{table_name}
orders_checkpoint = f"{checkpoint_base}/silver/orders"
metrics_checkpoint = f"{checkpoint_base}/gold/daily_metrics"This pattern is an extension of the namespace design principles covered in Unity Catalog Best Practices. The same organizational thinking applies to checkpoint paths.
A checkpoint directory contains three things:
Understanding this helps when debugging. If offsets exist but commits do not, the micro-batch failed mid-execution. The next restart will retry that batch automatically.
Real-world data arrives late. A mobile app event timestamped at 2:00 PM might reach your pipeline at 2:15 PM because of network delays. Without watermarks, your pipeline has two bad options: wait forever (impossible) or ignore late data (lossy).
Watermarks define how long the pipeline waits for late data before finalizing results.
from pyspark.sql.functions import col, window
# Allow data up to 30 minutes late
windowed_counts = (
events
.withWatermark("event_time", "30 minutes")
.groupBy(
window("event_time", "5 minutes"),
"event_type"
)
.count()
)This tells Spark: keep state for windows until 30 minutes after the window closes. Any data arriving later than 30 minutes is dropped.
The watermark duration is a trade-off between completeness and resource usage.
| Duration | Completeness | State Size | Use Case |
|---|---|---|---|
| 5 minutes | Low | Small | Real-time dashboards where speed matters more than precision |
| 30 minutes | Medium | Medium | Most production pipelines with moderate late data |
| 2 hours | High | Large | Financial or compliance systems where every record matters |
| 24 hours | Very high | Very large | Batch-like streaming where you aggregate daily |
A longer watermark means Spark holds more state in memory. For high-cardinality group keys, this can consume significant cluster resources.
This is a common source of confusion. In append output mode with watermarks, results are only emitted after the watermark passes. This means your downstream consumers see data with a delay equal to the watermark duration.
If your watermark is 30 minutes and your window is 5 minutes, results appear roughly 35 minutes after the events occurred. This is normal and expected. If you need faster output, use update mode instead, but understand that results may change as late data arrives.
The Structured Streaming chapter walks through these output modes step by step with practical examples.
By default, Structured Streaming processes data as fast as possible in continuous micro-batches. This is expensive. In most production scenarios, you do not need sub-second latency.
# Continuous processing (default) — expensive, sub-second latency
query = df.writeStream.trigger(processingTime="0 seconds")
# Fixed interval — process every 5 minutes
query = df.writeStream.trigger(processingTime="5 minutes")
# Available now — process all available data, then stop
query = df.writeStream.trigger(availableNow=True)
# Once — process one micro-batch, then stop (legacy)
query = df.writeStream.trigger(once=True)For most data engineering pipelines, availableNow=True is the right choice. It processes all accumulated data since the last run, then stops. You schedule it with Databricks Workflows to run at regular intervals.
# Scheduled streaming job — runs via Workflows every 15 minutes
(
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("/data/raw/events/")
.writeStream
.format("delta")
.trigger(availableNow=True)
.option("checkpointLocation", "/checkpoints/bronze/events")
.toTable("bronze.events")
)This pattern gives you streaming semantics (exactly-once processing, automatic schema detection) with batch economics (cluster runs only when there is work to do). On Serverless compute, this means you pay only for the minutes of actual processing. For more advanced orchestration, Databricks Lakeflow is GA and provides unified ingestion and transformation.
This is one of the production patterns detailed in our Workflows Best Practices article.
Auto Loader (cloudFiles format) is the recommended way to ingest files into Delta Lake. It automatically discovers new files, tracks which files have been processed, and handles schema evolution.
# Auto Loader with schema inference and evolution
raw_stream = (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.inferColumnTypes", "true")
.option("cloudFiles.schemaEvolutionMode", "addNewColumns")
.option("cloudFiles.schemaLocation", "/schemas/events")
.load("/data/landing/events/")
)
# Write to Bronze with checkpoint
(
raw_stream
.writeStream
.format("delta")
.trigger(availableNow=True)
.option("checkpointLocation", "/checkpoints/bronze/events")
.option("mergeSchema", "true")
.toTable("bronze.events")
)| Feature | spark.read (batch) | Auto Loader (streaming) |
|---|---|---|
| File tracking | None. You must track processed files yourself | Automatic via checkpoint |
| New file detection | Lists entire directory every run | Incremental discovery (notification or listing mode) |
| Schema evolution | Manual | Automatic with rescue column |
| Scalability | Slows with directory size | Scales to millions of files |
| Exactly-once | You must implement idempotency | Built in via checkpoint |
Auto Loader connects directly to the file format concepts in Chapter 10: File Formats and the data source patterns in Chapter 2: Data Sources. If you understand how Spark reads files, Auto Loader is the production extension of that knowledge. For managed pipelines, Lakeflow Declarative Pipelines (formerly Delta Live Tables) can automate this further.
Delta Lake is the natural sink for streaming pipelines. But there are patterns that matter for reliability.
# Bronze layer: append raw events
(
raw_stream.writeStream
.format("delta")
.outputMode("append")
.option("checkpointLocation", "/checkpoints/bronze/events")
.toTable("bronze.events")
)Append mode is ideal for Bronze. Every event is preserved exactly as it arrived, matching the Medallion Architecture principle that Bronze is append-only and immutable.
When you need MERGE, UPSERT, or multi-table writes, use foreachBatch.
def upsert_to_silver(batch_df, batch_id):
"""Merge streaming batch into Silver table."""
batch_df.createOrReplaceTempView("updates")
batch_df.sparkSession.sql("""
MERGE INTO silver.customers AS target
USING updates AS source
ON target.customer_id = source.customer_id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
""")
# Apply as streaming write
(
cleaned_stream.writeStream
.foreachBatch(upsert_to_silver)
.trigger(availableNow=True)
.option("checkpointLocation", "/checkpoints/silver/customers")
.start()
)This pattern enables the SCD (Slowly Changing Dimension) patterns and incremental processing strategies covered in earlier chapters, applied to streaming context.
A streaming pipeline that runs silently is a pipeline waiting to fail silently.
Every streaming query exposes progress metrics. Log them.
# Access the latest progress
query = df.writeStream.format("delta").start()
# Check progress
progress = query.lastProgress
if progress:
print(f"Input rows: {progress['numInputRows']}")
print(f"Processing time: {progress['batchDuration']}ms")
print(f"Input rows/sec: {progress['inputRowsPerSecond']}")If inputRowsPerSecond consistently exceeds processedRowsPerSecond, your pipeline is falling behind. You need more compute or optimized transformations.
These monitoring patterns extend the observability principles from Chapter 16: Debugging and Monitoring.
Build a simple metrics table for historical monitoring.
def log_streaming_metrics(query, pipeline_name):
"""Log streaming progress to Delta table."""
progress = query.lastProgress
if not progress:
return
metrics = spark.createDataFrame([{
"pipeline_name": pipeline_name,
"batch_id": progress["batchId"],
"input_rows": progress["numInputRows"],
"processing_time_ms": progress["batchDuration"],
"timestamp": progress["timestamp"]
}])
metrics.write.format("delta").mode("append") \
.saveAsTable("logs.streaming_metrics")Schema changes are inevitable. A new field appears in the source data, a column type changes, or a field is renamed.
Auto Loader handles this with rescue columns.
# Auto Loader with rescue column for schema mismatches
raw = (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", "/schemas/events")
.option("rescuedDataColumn", "_rescued_data")
.load("/data/events/")
)Any field that does not match the inferred schema is captured in _rescued_data instead of failing the pipeline. You can inspect rescued data later and evolve the schema intentionally.
This connects directly to the schema evolution strategies in Chapter 8: Schema Evolution, applied in a streaming context.
Before promoting a streaming pipeline to production, verify each of these.
Fault tolerance
Late data handling
Cost efficiency
Data quality
Monitoring
Structured Streaming is where several concepts from the book converge.
You need to understand data sources to read from files and message queues. You need Delta Lake for reliable writes. You need schema evolution for handling change. You need data quality for validation. You need Workflows for orchestration. And you need the Medallion Architecture to organize it all.
The streaming chapter in BricksNotes (Chapter 11) walks through each concept hands-on. This article gives you the production patterns to apply after you have built the foundation.
If you are preparing for the Databricks Data Engineer Associate certification, streaming checkpointing and trigger modes are frequently tested topics. Our certification guide maps each exam section to the chapters that cover it.
Structured Streaming is not about processing data fast. It is about processing data reliably, continuously, and efficiently. The best streaming pipelines are the ones nobody has to think about because they just work.