ワークフローを使った本番パイプライン

Databricks における Spark SQL を使ったデータ変換

Disha Mukherjee

Lead Data Engineer

Delta Lakeを使う理由

安全に輝くデジタルボールト。整理されたデータテーブルと保護シールドが描かれた、フラットモダンスタイルのイラスト

 

$$

  • ACIDトランザクション → 失敗した書き込みをロールバック
  • スキーマ強制 → 型の不一致をブロック
  • バージョン管理 → 以前の状態をクエリ可能
Databricks における Spark SQL を使ったデータ変換

Deltaへの書き込み

df_valid.write.format("delta") \
    .mode("overwrite") \
    .saveAsTable("transactions_clean")

print(f"Rows written: {df_valid.count():,}")
Rows written: 33,223

$$

$$

  • Unity Catalogに新しいDeltaテーブルが追加されます
Databricks における Spark SQL を使ったデータ変換

ノートブックタスク

task1_ingest: ロードとクリーニング

  • CSVを読み込み → クリーニングを適用 → Deltaテーブルに書き込み
Databricks における Spark SQL を使ったデータ変換

ノートブックタスク

task2_metrics: カテゴリ別売上

  • クリーニング済みテーブルを読み込み → メトリクスを計算 → 新しいテーブルに書き込み
Databricks における Spark SQL を使ったデータ変換

ノートブックタスク

task3_customers: 支出額によるランク付け

  • クリーニング済みテーブルを読み込み → 顧客をランク付け → 新しいテーブルに保存
Databricks における Spark SQL を使ったデータ変換

ジョブの作成

Jobs and Pipelines UIに表示された、task1_ingestからtask2_metrics、task3_customersへと依存関係の矢印が示された完了済み3タスクのジョブDAG

  • 手動実行、スケジュール設定、トリガーの設定が可能
Databricks における Spark SQL を使ったデータ変換

ジョブの実行

ジョブ実行DAG。task1_ingestとtask2_metricsは緑色で成功、task3_customersは赤色で失敗

Databricks における Spark SQL を使ったデータ変換

ジョブの実行

タスク3のエラーを確認する

  • エラー: customer-idCustomer_ID が正しい表記
Databricks における Spark SQL を使ったデータ変換

ジョブの実行

成功したグラフビュー

タイムラインビュー

Databricks における Spark SQL を使ったデータ変換

Lakeflowとは?

$$

$$

比較: 命令的手法(Jobs)| 宣言的手法(Lakeflow)

 

  • Jobs → 各ステップを手動で管理
  • Lakeflow → テーブルの内容を宣言的に定義
  • 順序、リトライ、コンピューティングはDatabricksが管理
Databricks における Spark SQL を使ったデータ変換

@dlt.tableパターン

@dlt.table(name="transactions_bronze")
def transactions_bronze():
    return spark.read.format("csv").schema(schema).load(FILE_PATH)

@dlt.table(name="transactions_silver") def transactions_silver(): return dlt.read("transactions_bronze").na.drop(...).filter(...)
@dlt.table(name="category_revenue_gold") def category_revenue_gold(): return dlt.read("transactions_silver").groupBy("Category").agg(...)
Databricks における Spark SQL を使ったデータ変換

パイプラインの実行

Lakeflowパイプラインのdag。transactions_bronze(10万行)からtransactions_silver(3.3万行)、category_revenue_gold(6行)へと続く、すべて緑色のマテリアライズドビュー

  • Bronze(10万行)→ Silver(3.3万行)→ Gold(6行)
Databricks における Spark SQL を使ったデータ変換

Notebooks、Jobs、Lakeflow の使い分け

$$

レイヤー: Notebooks、Databricks Jobs、Lakeflow Pipelines

 

$$

  • Notebooks → 探索とプロトタイピング
  • Databricks Jobs → 複数ステップのスケジュール済みパイプライン
  • Lakeflow → フルマネージド、宣言的
Databricks における Spark SQL を使ったデータ変換

練習しましょう!

Databricks における Spark SQL を使ったデータ変換

Preparing Video For Download...