Databricks Spark Structured Streaming

作者 databrickse77e37e8a4da無授權條款345 個星標收錄於 2026年10月8日更新於 2026年10月8日儲存庫今天更新

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 產生的概覽

在 Databricks 上建置正式環境 Spark Structured Streaming 管線的指南,涵蓋 Kafka、聯結、狀態、檢查點與成本。

功能
此技能提供在 Databricks 上用於正式環境工作負載的 Spark Structured Streaming 模式導覽指南。它指向多份參考文件,涵蓋 Kafka 擷取、即時模式、Lakebase 接收端、串流對串流與串流對靜態聯結、多接收端寫入、合併作業、檢查點、具狀態作業、觸發條件與成本最佳化。它也包含一段快速入門程式碼片段,以及一份正式環境檢查清單,涵蓋檢查點持久化、叢集規模、監控、精確一次語意與水位線。
適用情境
在建置或調校串流管線時使用,例如從 Kafka 擷取到 Delta、實作即時模式、設定觸發條件、以水位線處理具狀態作業,或最佳化串流成本與效能。
執行需求
需要 Databricks CLI(>= v1.0.0),以及具備 Spark Structured Streaming 的 Databricks 環境。不隨附指令碼,僅有指示說明、參考 Markdown 檔案與圖片素材。

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)

來源與署名

來源:databricks/databricks-agent-skills位於plugins/databricks/claude/skills/databricks-spark-structured-streaming提交e77e37e

授權條款: 無授權條款

內容歸原作者所有。SourceWeft 從公開儲存庫中收錄這些內容。

檢舉或申請下架