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.
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 を参照 |
Production-ready streaming pipelines with Spark Structured Streaming. This skill provides navigation to detailed patterns and best practices.
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")
| 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 |
| 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 |
| Topic | Description | Reference |
|---|---|---|
| Production Checklist | Comprehensive best practices | See references/streaming-best-practices.md |
原文・著作権は Anthropic および各プラグイン作者に帰属します。日本語訳は Claude API による自動翻訳です。