Databricks Spark Structured Streaming

by databrickse77e37e8a4daNo license345 starsListed Oct 8, 2026Updated Oct 8, 2026Repository updated today

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.

AI-generated overview

Guide to building production Spark Structured Streaming pipelines on Databricks, covering Kafka, joins, state, checkpoints and cost.

What it does
This skill provides a navigational guide to Spark Structured Streaming patterns for production workloads on Databricks. It points to reference documents on Kafka ingestion, Real-Time Mode, Lakebase sinks, stream-stream and stream-static joins, multi-sink writes, merge operations, checkpoints, stateful operations, triggers and cost optimization. It also includes a quick-start code snippet and a production checklist covering checkpoint persistence, cluster sizing, monitoring, exactly-once semantics and watermarks.
When to use it
Use it when building or tuning streaming pipelines, such as ingesting from Kafka into Delta, implementing Real-Time Mode, configuring triggers, handling stateful operations with watermarks, or optimizing streaming cost and performance.
Requirements
Requires the Databricks CLI (>= v1.0.0) and a Databricks environment with Spark Structured Streaming. Ships no scripts; it is instructions plus reference markdown files and image assets.

Spark Structured Streaming

Production-ready streaming pipelines with Spark Structured Streaming. This skill provides navigation to detailed patterns and best practices.

Quick Start

python
from pyspark.sql.functions import col, from_json
# Basic Kafka to Delta streamingdf = (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

PatternDescriptionReference
Kafka StreamingKafka to Delta, Kafka to Kafka, Real-Time ModeSee references/kafka-streaming.md [blocked]
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 semanticsSee references/real-time-mode.md [blocked]
Lakebase SinkWrite streaming records into Lakebase Postgres with transactional upserts. Native format("postgresql") sink (DBR 18.3+) and manual foreach sink as a fallbackSee references/lakebase-sink-python.md [blocked]
Stream JoinsStream-stream joins, stream-static joinsSee references/stream-stream-joins.md [blocked], references/stream-static-joins.md [blocked]
Multi-Sink WritesWrite to multiple tables, parallel mergesSee references/multi-sink-writes.md [blocked]
Merge OperationsMERGE performance, parallel merges, optimizationsSee references/merge-operations.md [blocked]

Configuration

TopicDescriptionReference
CheckpointsCheckpoint management and best practicesSee references/checkpoint-best-practices.md [blocked]
Stateful OperationsWatermarks, state stores, RocksDB configurationSee references/stateful-operations.md [blocked]
Trigger & CostTrigger selection, cost optimization, RTMSee references/trigger-and-cost-optimization.md [blocked]

Best Practices

TopicDescriptionReference
Production ChecklistComprehensive best practicesSee references/streaming-best-practices.md [blocked]

Production Checklist

  • Checkpoint location is persistent (UC volumes, not DBFS)
  • Unique checkpoint per stream
  • Fixed-size cluster (no autoscaling for streaming)
  • Monitoring configured (input rate, lag, batch duration)
  • Exactly-once verified (txnVersion/txnAppId)
  • Watermark configured for stateful operations
  • Left joins for stream-static (not inner)

Source and attribution

Source:databricks/databricks-agent-skillsinplugins/databricks/claude/skills/databricks-spark-structured-streamingat commite77e37e

License: No license

Content belongs to its original authors. SourceWeft indexes it from a public repository.

Report or request removal