mcpbeat

Databricks Spark Structured Streaming

databricks/databricks-agent-copilot-databricks-spark-structured-streaming

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.

This is a copy. The original lives at databricks/databricks-spark-structured-streaming.

43k tokens
context cost
the whole folder, loaded on every use
15
files
instructions only
0
copies elsewhere
how many repositories repackaged it
236
stars on the repo
on the repository, not the skill itself

Install

one command, takes just this skill from the repository
npx skills add https://github.com/databricks/databricks-agent-skills --skill databricks-spark-structured-streaming

The instruction itself

6 sections, as written by the author

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

  • [ ] 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)

How to use it

Copy the folder

Take databricks/databricks-agent-copilot-databricks-spark-structured-streaming from the repository into ~/.claude/skills for personal use, or into .claude/skills inside a project.

Check the name does not clash

The agent identifies a skill by the name field in its header. Two skills with the same name cannot sit side by side — one of them will be ignored.