• Projects
  • Service
  • About
  • branding.bz
  • Podcast
  • Tips
  • FAQ
  • Recruit
  • Download
  • Contact
  • branding.bz(ブランド構築SaaS)
  • DESIGN NOW(デザインメディア)
  • X
  • LinkedIn
  • Spotify
  • Facebook

213-0011 神奈川県川崎市高津区久本3-6-7-303

© 2026 ID INC. All rights reserved

claude-skills/スキル
SKILLOfficialdatabase

databricks-spark-structured-streaming

プラグイン
databricks
ソース
GitHub で見る ↗
説明

Spark Structured Streamingの本番環境での活用に関する包括的なガイドです。 次のような場合に使用: ストリーミング・パイプライン(データを継続的に処理する仕組み)の構築、Kafka(分散メッセージング基盤)からのデータ取り込み、リアルタイムモード(RTM)の実装、トリガー設定(処理時間やデータ利用可能時の自動実行条件)の構成、ウォーターマーク(遅延データ処理の時間枠)を使った状態管理の操作、チェックポイント(処理の復旧地点)の最適化、ストリーム同士またはストリームと静的データ(テーブル)の結合処理、複数の出力先への書き込み、またはストリーミングのコスト・パフォーマンス調整が必要な場合

原文を表示

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.

ユースケース
  • ストリーミング・パイプラインの構築
  • Kafkaからのデータ取り込み
  • リアルタイムモードの実装
  • トリガー設定の構成
  • ウォーターマークを使った状態管理
  • チェックポイントの最適化
本文(日本語訳)

Spark Structured Streaming

Spark Structured Streamingを使った本番環境対応のデータ処理パイプライン(リアルタイムデータの一連の処理フロー)です。このスキルは、詳細なパターンとベストプラクティス(最適な実装方法)への案内を提供します。

クイックスタート

from pyspark.sql.functions import col, from_json

# Kafkaから Delta へのシンプルなストリーミング処理
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")

主要パターン

パターン 説明 詳細
Kafkaストリーミング Kafkaから Delta へ、Kafkaから Kafka へ、リアルタイムモード references/kafka-streaming.md を参照
リアルタイムモード(RTM) 1秒未満の低遅延処理 — クラスタ構成、計算リソース配分、対応操作(DBR 18以降でストリーム間結合を含む)、状態付き変換、監視可能性、エラー分類、配信保証 references/real-time-mode.md を参照
Lakebaseシンク ストリーミングレコード(データの一行一行)を Lakebase Postgres に書き込み、トランザクション(一連の処理操作)保証によるデータの挿入・更新を行う。ネイティブな format("postgresql") シンク(DBR 18.3以降)と、代替手段としての foreach シンク references/lakebase-sink-python.md を参照
ストリーム結合 ストリーム同士の結合、ストリームと静的データの結合 references/stream-stream-joins.md、references/stream-static-joins.md を参照
複数の出力先への書き込み 複数のテーブルへの書き込み、並列マージ処理 references/multi-sink-writes.md を参照
マージ操作 マージのパフォーマンス、並列マージ、最適化 references/merge-operations.md を参照

設定

トピック 説明 詳細
チェックポイント チェックポイント(処理状態の保存地点)の管理とベストプラクティス references/checkpoint-best-practices.md を参照
状態保持操作 ウォーターマーク(時間経過追跡)、状態ストア、RocksDB設定 references/stateful-operations.md を参照
トリガーとコスト最適化 トリガー選択、コスト最適化、RTM references/trigger-and-cost-optimization.md を参照

ベストプラクティス

トピック 説明 詳細
本番環境チェックリスト 包括的なベストプラクティス references/streaming-best-practices.md を参照

本番環境チェックリスト

  • [ ] チェックポイント保存場所が永続的である(UC ボリューム、DBFS 不可)
  • [ ] ストリームごとに独立したチェックポイント
  • [ ] 固定サイズのクラスタ(ストリーミングに自動スケーリング不可)
  • [ ] 監視設定済み(入力レート、遅延、バッチ処理時間)
  • [ ] 1回限りの配信を確認済み(txnVersion/txnAppId)
  • [ ] 状態保持操作にウォーターマーク設定済み
  • [ ] ストリームと静的データの結合では左結合を使用(内部結合不可)
原文(English)を表示

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)

原文・著作権は Anthropic および各プラグイン作者に帰属します。日本語訳は Claude API による自動翻訳です。