Databricks公式のハンズオンが新しくなっていたので、ETLパイプラインを実践してみた

Databricksの公開ハンズオンでETLパイプラインを実践。サンプルデータの生成、Auto LoaderによるBronzeへの増分取り込み、Silver・Goldの加工と集計を振り返ります。

目次 12項目

Databricks公式のハンズオンが新しくなったと聞いたので、実践してみました。

今回取り組んだのは、データエンジニアリングワークショップ:データ取り込みから加工・可視化・ジョブ化までです。資料もノートブックも公開されていて、サンプルデータを用意するコードまでそろっています。自分のタイミングで手を動かせる、ありがたい教材です。

特に勉強になったのは、BronzeでのAuto Loaderの書き方でした。データを読み込むだけでなく、チェックポイントの置き方や、取り込み元を後から確認するための列まで含めて学べます。

この記事では、サンプルデータの準備とBronze・Silver・Goldの処理を中心に、触ってみてよかった点をまとめます。

今回のハンズオンで作るもの

題材は、架空のキャッシュレス決済サービスのデータです。決済履歴などのファクトデータ(Fact)と、会員・加盟店・決済手段のマスターデータ(Master)を使って、段階的に分析用のテーブルを作ります。

処理の流れを簡単にすると、次のようになります。

テキスト / 実行結果
admin:サンプルデータを生成
  │
  ▼
共有Volume:Parquetファイルを保存
  │
  ├─ Master → 05_setup ─────┐
  │                         │
  └─ Fact ──→ Bronze        │
                 │          │
                 └────┬─────┘
                      ▼
                   Silver:整形・結合
                      │
                      ▼
                   Gold:分析用に集計

ワークショップ全体では、ここからAI/BIダッシュボードによる可視化や、Lakeflow Jobsによる自動化にも進めます。詳細な操作手順はDatabricks Japanのスライドにまとまっています。

教材の閲覧は無料で、Databricks Free Editionを利用してハンズオンを試すこともできます。Free Editionにはコンピュートの使用量や機能の制限があるため、Free Editionの制限を確認してください。後半で補足する自動クラスタリングについては、別途、利用条件を確認する必要があります。

サンプルデータはadminフォルダーで生成できる

最初に押さえておきたいのが、サンプルデータの準備です。

実践してみると、「あれ?共有カタログがないから、個人では実施できないハンズオンなのかも?」と思いました。ところが、adminフォルダーのノートブックを実行すると、乱数を使ったサンプルデータを生成できる仕組みになっていました。

公開リポジトリのde_workshopフォルダーには、参加者が使うノートブックと、講師向けのadminフォルダーが入っています。個人で試す場合は、講師側で行うデータ準備も自分で実行します。

ノートブック 役割 実行するタイミング
admin/00_setup_and_generate_initial_data.py 初期データを生成し、共有Volumeへ保存 最初の準備
05_setup.py 自分用のスキーマ、チェックポイント用Volume、マスターテーブルを作成 Bronzeの前
10_bronze_autoloader.py ファクトデータをAuto Loaderで取り込み 初回と増分追加後
admin/01_generate_incremental_data.py 増分取り込みを試すための追加データを生成 Bronzeの初回取り込み後
20_silver_cleanse_and_join.py 重複排除、整形、マスター結合 Bronzeの後
30_gold_aggregations.py 分析用の集計テーブルを作成 Silverの後

初期データ作成の手順

ノートブックをDatabricksのワークスペースへ取り込み、フォルダー構成を保ったまま実行します。

まず、admin/00_setup_and_generate_initial_data.pyを開き、「パラメータ設定」のセルにあるcatalogを、利用できる既存のカタログ名に変更します。カタログがまだない場合は、事前に作成する必要があります。

python
# 保存先:利用できる既存のカタログ名に置き換える
catalog = "<既存のカタログ名>"
shared_schema = "de_workshop_shared"  # 共有スキーマ

# レコード件数(実装検証用に縮小)
num_customers = 5000        # 会員数
num_merchants = 500         # 加盟店数
num_payment_methods = int(num_customers * 1.8)  # 決済手段(顧客あたり約1.8個)
num_payments = 100000       # 決済件数
num_point_transactions = 20000
num_charges = 2000
num_login_events = 20000

このノートブックをすべて実行すると、指定したカタログの中に共有スキーマとVolumeが作成され、サンプルデータが生成されます。カタログ自体を作成する処理は含まれていません。

次に、05_setup.pyの「変数定義」セルで、user_nameを自分の識別子に変更し、catalogを先ほど指定したカタログ名とそろえます。user_nameはスキーマ名の一部になるため、taroのような英小文字・数字・アンダースコアを使った値にしておくと扱いやすいです。

python
# 自分の識別子(この例では dew_taro スキーマが作られる)
user_name = "taro"

# admin側と同じ既存のカタログ名に置き換える
catalog = "<既存のカタログ名>"

05_setup.pyを実行し、自分のスキーマにチェックポイント用Volumeと3つのbronze_m_*マスターテーブルが作られれば、準備は完了です。後続のノートブックは%run ./05_setupで共通の変数を読み込みます。

一人で試す場合も、カタログの利用やスキーマ・Volume・テーブルの作成権限は必要です。adminのノートブックにはaccount usersへの共有権限付与も含まれるため、組織のワークスペースで実行する場合は、共有範囲が適切か確認してください。

なお、初期データ生成には上書き処理が含まれています。増分取り込みを試す途中で最初から生成し直すと比較条件が変わるので、初期データの準備と追加データの投入は分けて進めます。

BronzeはAuto Loaderの具体的な書き方が勉強になる

今回いちばん参考になったのが、10_bronze_autoloader.pyです。

ここのコードはしっかり覚えておきたいです。以下は公開ノートブックの取り込み関数です。先に%run ./05_setupを実行し、Fやパスなどの共通変数を読み込んでから使います。

python
def ingest_fact(table: str):
    src  = f"{shared_landing}/{table}/"          # 共有 Volume のソースパス(全員共通)
    ckpt = f"{checkpoint_base}/{table}/"         # 自分専用のチェックポイント
    (spark.readStream
        .format("cloudFiles")                                        # ★Auto Loader(増分・自動検出)
        .option("cloudFiles.format", "parquet")                      # 取り込むファイル形式
        .option("cloudFiles.schemaLocation", ckpt)                   # スキーマ情報の保存先
        .option("cloudFiles.schemaEvolutionMode", "addNewColumns")   # 新規列を自動追加(スキーマ進化)
        .load(src)
        # ★取り込み時刻を東京タイムゾーンで付与(増分取り込みの確認に使用)
        .withColumn("_ingested_at", F.from_utc_timestamp(F.current_timestamp(), "Asia/Tokyo"))
        .withColumn("_source_file", F.col("_metadata.file_path"))    # 由来ファイルパスを記録
     .writeStream
        .clusterBy("customer_id")                                    # ★Bronze も Liquid Clustering(結合キー)
        .option("checkpointLocation", ckpt)                          # ★チェックポイント(既読管理=冪等性)
        .trigger(availableNow=True)                                  # ★到着分を処理して停止(バッチ的)
        .toTable(f"{bp}.bronze_{table}")                             # 自分のスキーマに Bronze テーブル出力
     .awaitTermination())                                            # 取り込み完了まで待機

共有VolumeにあるParquetファイルを、spark.readStream.format("cloudFiles")で読み込み、自分のスキーマのBronzeテーブルへ書き出します。取り込み処理は関数になっていて、決済・ポイント・チャージ・ログインの4種類のデータに同じ形で適用できます。

チェックポイントとスキーマの保存先が分かる

コードを読むときに注目したい設定は、次のとおりです。

設定 この教材での役割
cloudFiles.format 入力ファイルをParquetとして読み込む
cloudFiles.schemaLocation 推論したスキーマの情報を保存する
checkpointLocation ストリームの処理状況を保存し、次回の取り込みに引き継ぐ
trigger(availableNow=True) 実行開始時点までに到着した未処理ファイルを処理し、終了する

教材では、参加者ごとのチェックポイント用Volumeの下に、さらに取り込み対象ごとのディレクトリを作っています。共有の入力データを使いながらも、各自の取り込み状態を分けられる構成です。

同じ処理を再実行するときは、同じチェックポイントを使います。チェックポイントを毎回作り直すと、前回どこまで処理したかを引き継げません。公式の本番運用向けドキュメントでも、チェックポイントの保護とAvailableNowによるバッチ的な実行が説明されています。

「ストリーミングのAPIを使うけれど、到着済みのデータを取り込んだら止める」という書き方は、定期実行のETLにも結び付けて考えやすかったです。

取り込み後に確認するための列も残している

このノートブックでは、取り込み時刻の_ingested_atと、元ファイルのパスを記録する_source_fileも付与しています。

取り込み件数だけでなく、「いつ取り込んだか」「どのファイルから来たか」を追えるようにしているところがよいです。増分取り込みの動きを確認するときにも、この2つの列が役立ちます。

取り込み処理と、その結果を確かめる方法がセットになっているのが、この教材の分かりやすいところだと思いました。

増分取り込みは同じコードをもう一度実行する

教材では、Bronzeを一度実行した後、adminの増分データ生成ノートブックで追加ファイルを配置し、Bronzeの取り込みをもう一度実行します。取り込みコードを書き換える必要はありません。

確認するポイントは、初回のデータが再取り込みされず、追加したファイルの分だけ件数が増えることです。_source_fileを見れば、追加データのbatch=ディレクトリに由来する行も確認できます。追加ファイルがない状態で再実行した場合は、件数が増えないことも確認できます。

ここでの増分取り込みは、ファイル単位の処理状況を使うものです。同じ業務レコードを別の新規ファイルに入れても、自動的に主キーで重複排除されるわけではありません。後続のSilverで行う重複排除とは役割が異なります。

スキーマ変更時の動きは押さえておきたい

教材には、cloudFiles.schemaEvolutionModeをaddNewColumnsにする設定もあります。

この設定は、新しい列が来ても処理が止まらずに続く、という意味ではありません。公式ドキュメントでは、新しい列を検出すると保存済みのスキーマを更新してストリームが停止し、再起動後に処理を再開すると説明されています。運用では、Lakeflow Jobsなどでの再起動に加え、書き込み先のDeltaテーブルのスキーマ変更も考える必要があります。

短いサンプルでも、設定の意味を一つずつ追うと、実際の取り込み処理を書くための参考になります。

SilverとGoldは、加工と集計の基本を確認するパート

SilverとGoldについては、正直なところ、私には大きく目新しいものはありませんでした。PySparkやSQLでデータ加工をしている人なら、見慣れた処理が多いと思います。

Silverでは、決済IDによる重複排除、有効なステータスへの絞り込み、マスターデータとの結合、日付や判定用の列の追加などを行います。Goldでは、その結果から顧客別・加盟店別・決済手段別・日次の集計テーブルを作ります。

ただ、Bronzeで取り込んだデータが、整形されて分析に使う形になるまでを一続きで追えるのはよかったです。取り込みだけのサンプルで終わらないので、Auto LoaderをETL全体のどこで使うかが分かります。

なお、公開コードのSilverとGoldは、読み込んだテーブルを加工・集計して上書きする構成です。Bronzeが増分取り込みだからといって、後続の処理もすべて差分更新になるわけではありません。追加データを集計に反映するには、Bronzeの後にSilver、Goldの順で再実行します。

リキッドクラスタリングの指定も参考になった

目新しいものはないと言っておきながら、一つありました。ハンズオンのコードでは、BronzeのファクトテーブルとSilver・Goldのテーブルに、書き込み時のリキッドクラスタリングを指定しています。結合や集計に使う列をキーにしているところも参考になりました。

教材ではキーを手動で指定していますが、自動クラスタリングのAUTOを使う方法もあります。公式の推奨では、Unity Catalogのマネージドテーブルに自動リキッドクラスタリングと予測的最適化を組み合わせることが勧められています。

ただし、ここからは教材の手動指定に対する発展的な補足です。マネージドDeltaテーブルの自動クラスタリングはDatabricks Runtime 15.4 LTS以降に対応し、キーの自動選択やクラスタリング処理には予測的最適化が必要です。予測的最適化の利用条件にはPremiumプラン以上などが挙げられており、Free Editionでも同じように自動最適化が動くとは限りません。

Pythonの.clusterBy("auto")は、autoという列を指定する書き方です。自動クラスタリングにするには、対応環境で.option("clusterByAuto", "true")を使います。公式ドキュメントによると、このPython APIはDatabricks Runtime 16.4以降で利用できます。

python
# 20_silver_cleanse_and_join.pyの書き込みを置き換える例
# テーブルを作成・置換するときに自動クラスタリングを指定する
(silver_payments.write.mode("overwrite")
    .format("delta")
    .option("overwriteSchema", "true")
    .option("clusterByAuto", "true")
    .saveAsTable(f"{bp}.silver_payments"))

既存テーブルの設定だけを変更する場合は、SQLを使います。

python
# 05_setup.pyで定義したbpを使い、既存テーブルの設定を変更する
spark.sql(f"ALTER TABLE {bp}.silver_payments CLUSTER BY AUTO")

小さなテーブルやクエリの少ないテーブルでは、自動クラスタリングを有効にしてもキーが選ばれない場合があります。設定を付ければ必ず処理が速くなるというものではなく、データ量や利用状況に応じて判断される仕組みです。

Auto Loaderを具体的に学びたいときによい教材だった

今回のハンズオンは、私にとってはAuto Loaderの実装例を読めたことがいちばんの収穫でした。

チェックポイントをどこに置くか、到着済みのファイルを処理して止めるにはどう書くか、取り込み結果をどう確認するか。そうした点を、動かせるノートブックで学べるのがよいです。

SilverとGoldは基本の復習という印象でしたが、データの生成から集計までそろっているので、ETLパイプライン全体を試す題材としても使いやすいと感じました。Auto Loaderの書き方を具体的に知りたい人は、Bronzeのノートブックから読んでみると参考になると思います。

参考資料