Spark Structured Streaming

SkillDev tools

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.

Available today. Use it from your connected AI after setup.

Connect ahel once, and every AI you use reads what you have installed.

Then ask your AI: use the Spark Structured Streaming skill

What this skill tells your AI

The instructions your AI receives, as published by kilo-org/kilo-marketplace in skills/databricks-spark-structured-streaming/SKILL.md and read by ahel’s review.

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

PatternDescriptionReference
Kafka StreamingKafka to Delta, Kafka to Kafka, Real-Time ModeSee references/kafka-streaming.md
Stream JoinsStream-stream joins, stream-static joinsSee references/stream-stream-joins.md, references/stream-static-joins.md
Multi-Sink WritesWrite to multiple tables, parallel mergesSee references/multi-sink-writes.md
Merge OperationsMERGE performance, parallel merges, optimizationsSee references/merge-operations.md

Configuration

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

Best Practices

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

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)

Signals

GitHub stars
175
Forks
159
Last commit
Aug 2026
Advanced
Catalog kind
skill
Gateway key
databricks-spark-structured-streaming
Source
github.com/kilo-org/kilo-marketplace