Spark Structured Streaming
Production-ready streaming pipelines with Spark Structured Streaming. This skill provides navigation to detailed patterns and best practices.
Quick Start
from pyspark.sql.functions import col, from_json
# Basic Kafka to Delta streaming
df = (spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "broker:9092")
.option("subscribe", "topic")
.load()
.select(from_json(col("value").cast("string"), schema).alias("data"))
.select("data.*")
)
df.writeStream \
.format("delta") \
.outputMode("append") \
.option("checkpointLocation", "/Volumes/catalog/checkpoints/stream") \
.trigger(processingTime="30 seconds") \
.start("/delta/target_table")
Core Patterns
| Pattern |
Description |
Reference |
| Kafka Streaming |
Kafka to Delta, Kafka to Kafka, Real-Time Mode |
See references/kafka-streaming.md |
| Real-Time Mode (RTM) |
Sub-second E2E latency — cluster setup, slot math, supported ops (incl. stream-stream inner join on DBR 18+), transformWithState, observability, error classes, delivery semantics |
See references/real-time-mode.md |
| Lakebase Sink |
Write streaming records into Lakebase Postgres with transactional upserts. Native format("postgresql") sink (DBR 18.3+) and manual foreach sink as a fallback |
See references/lakebase-sink-python.md |
| Stream Joins |
Stream-stream joins, stream-static joins |
See references/stream-stream-joins.md, references/stream-static-joins.md |
| Multi-Sink Writes |
Write to multiple tables, parallel merges |
See references/multi-sink-writes.md |
| Merge Operations |
MERGE performance, parallel merges, optimizations |
See references/merge-operations.md |
Configuration
| Topic |
Description |
Reference |
| Checkpoints |
Checkpoint management and best practices |
See references/checkpoint-best-practices.md |
| Stateful Operations |
Watermarks, state stores, RocksDB configuration |
See references/stateful-operations.md |
| Trigger & Cost |
Trigger selection, cost optimization, RTM |
See references/trigger-and-cost-optimization.md |
Best Practices
| Topic |
Description |
Reference |
| Production Checklist |
Comprehensive best practices |
See references/streaming-best-practices.md |
Production Checklist
1---2name: databricks-spark-structured-streaming3description: Comprehensive guide to Spark Structured Streaming for production workloads. Use when building streaming pipelines, working with Kafka ingestion, implementing Real-Time Mode (RTM), configuring triggers (processingTime, availableNow), handling stateful operations with watermarks, optimizing checkpoints, performing stream-stream or stream-static joins, writing to multiple sinks, or tuning streaming cost and performance.4---5
6# Spark Structured Streaming
7
8Production-ready streaming pipelines with Spark Structured Streaming. This skill provides navigation to detailed patterns and best practices.
9
10## Quick Start
11
12```python
13from pyspark.sql.functions import col, from_json
14
15# Basic Kafka to Delta streaming
16df = (spark
17 .readStream
18 .format("kafka")
19 .option("kafka.bootstrap.servers", "broker:9092")
20 .option("subscribe", "topic")
21 .load()
22 .select(from_json(col("value").cast("string"), schema).alias("data"))
23 .select("data.*")
24)
25
26df.writeStream \
27 .format("delta") \
28 .outputMode("append") \
29 .option("checkpointLocation", "/Volumes/catalog/checkpoints/stream") \
30 .trigger(processingTime="30 seconds") \
31 .start("/delta/target_table")
32```
33
34## Core Patterns
35
36| Pattern | Description | Reference |
37|---------|-------------|-----------|
38| **Kafka Streaming** | Kafka to Delta, Kafka to Kafka, Real-Time Mode | See [references/kafka-streaming.md](references/kafka-streaming.md) |
39| **Real-Time Mode (RTM)** | Sub-second E2E latency — cluster setup, slot math, supported ops (incl. stream-stream inner join on DBR 18+), `transformWithState`, observability, error classes, delivery semantics | See [references/real-time-mode.md](references/real-time-mode.md) |
40| **Lakebase Sink** | Write streaming records into Lakebase Postgres with transactional upserts. Native `format("postgresql")` sink (DBR 18.3+) and manual `foreach` sink as a fallback | See [references/lakebase-sink-python.md](references/lakebase-sink-python.md) |
41| **Stream Joins** | Stream-stream joins, stream-static joins | See [references/stream-stream-joins.md](references/stream-stream-joins.md), [references/stream-static-joins.md](references/stream-static-joins.md) |
42| **Multi-Sink Writes** | Write to multiple tables, parallel merges | See [references/multi-sink-writes.md](references/multi-sink-writes.md) |
43| **Merge Operations** | MERGE performance, parallel merges, optimizations | See [references/merge-operations.md](references/merge-operations.md) |
44
45## Configuration
46
47| Topic | Description | Reference |
48|-------|-------------|-----------|
49| **Checkpoints** | Checkpoint management and best practices | See [references/checkpoint-best-practices.md](references/checkpoint-best-practices.md) |
50| **Stateful Operations** | Watermarks, state stores, RocksDB configuration | See [references/stateful-operations.md](references/stateful-operations.md) |
51| **Trigger & Cost** | Trigger selection, cost optimization, RTM | See [references/trigger-and-cost-optimization.md](references/trigger-and-cost-optimization.md) |
52
53## Best Practices
54
55| Topic | Description | Reference |
56|-------|-------------|-----------|
57| **Production Checklist** | Comprehensive best practices | See [references/streaming-best-practices.md](references/streaming-best-practices.md) |
58
59## Production Checklist
60
61- [ ] Checkpoint location is persistent (UC volumes, not DBFS)
62- [ ] Unique checkpoint per stream
63- [ ] Fixed-size cluster (no autoscaling for streaming)
64- [ ] Monitoring configured (input rate, lag, batch duration)
65- [ ] Exactly-once verified (txnVersion/txnAppId)
66- [ ] Watermark configured for stateful operations
67- [ ] Left joins for stream-static (not inner)