mcpbeat Sign in

Databricks Spark Structured Streaming Agent Skill

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.

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)

Other skills for the same job

different authors, same section of the catalogue
Workers Best Practices
by openai
vendor ×1

Reviews and authors Cloudflare Workers code against production best practices. Load when writing new Workers, reviewing Worker code, configuring wrangler.jsonc, or checking for common Workers anti-patterns (streaming, floating promises, global state, secrets, bindings, observability). Biases towards retrieval from Cloudflare docs over pre-trained knowledge.

8k tokens
Azure AI Translation Text Py
by lingxling
×1

Azure AI Text Translation SDK for real-time text translation, transliteration, language detection, and dictionary lookup. Use for translating text content in applications.

2k tokens
Add Model Price
by langfuse
vendor

Use when editing worker/src/constants/default-model-prices.json, packages/shared/src/server/llm/types.ts, pricing tiers, tokenizer IDs, or matchPattern regexes for OpenAI, Anthropic, Bedrock, Vertex, Azure, or Gemini model pricing.

20k tokens scripts
Contributing
by cloudflare
vendor

Use when contributing to the Cloudflare Docs repository — writing or editing documentation pages, choosing content types or components, adding changelog entries, reviewing docs, or learning how to contribute.

12k tokens
Story Setup
by worldwonderer

网文写作工具集基础设施部署。为 Claude Code / OpenCode / Codex / ZCode / OpenClaw / Reasonix 提供内置适配;Web AI / 通用 Agent 可走 skills + AGENTS.md 文件模式。触发方式:/story-setup、$story-setup、「准备写书」「帮我搭一下环境」「配置写作项目」。

305k tokens scripts zh
Authoring Github Workflows
by dotnet

Author and review GitHub Actions workflow YAML safely so syntactically-valid YAML can't ship a workflow that GitHub Actions refuses to run. USE FOR: editing, adding, or reviewing any file under .github/workflows/, writing run-name/name/if/env/run values that contain ${{ }} expressions, diagnosing a run that fails with 'This run likely failed because of a workflow file issue' and no jobs starting, deciding when a workflow scalar must be quoted, validating workflows with actionlint. DO NOT USE FOR: authoring application YAML unrelated to GitHub Actions, Azure Pipelines, GitLab CI, or non-workflow YAML. SCOPE: this skill covers *syntactic/structural* correctness of workflow YAML (quoting, parsing, actionlint); for *semantic and functional* workflow design (what a workflow should do, agentic-workflow behavior), see .github/agents/agentic-workflows.agent.md — the two are complementary. INVOKES: actionlint (downloaded pinned binary) plus git/grep for inspection.

2k tokens
Azure AI Translation Text Py
by microsoft
vendor

| Azure AI Text Translation SDK for real-time text translation, transliteration, language detection, and dictionary lookup. Use for translating text content in applications.

4k tokens
Azure AI Translation Ts
by microsoft
vendor

Build translation applications using Azure Translation SDKs for JavaScript (@azure-rest/ai-translation-text, @azure-rest/ai-translation-document). Use when implementing text translation, transliteration, language detection, or batch document translation.

2k tokens

How to use it

Copy the folder

Take databricks/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.