Databricks Certified Data Engineer Associate 教科書
第2章 データ取り込みとロード(Data Ingestion and Loading, 21%)
🎯 この節の学習目標
trigger(availableNow=True) によるバッチ的運用を説明できるAuto Loader は、クラウドオブジェクトストレージに到着する新しいファイルを増分的に検出してロードする Structured Streaming のソースです。コード上はフォーマット名 cloudFiles として登場します。COPY INTO と同じく「一度処理したファイルは二度処理しない」性質を持ちますが、その記録方法が異なります。Auto Loader はチェックポイント(checkpoint)に処理済みファイルの情報を永続化し、障害時もチェックポイントから正確に一度(exactly-once)の保証つきで再開できます。
💡 具体例:S3 の JSON を Auto Loader で取り込む(PySpark)
(spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", "s3://my-bucket/_schemas/events/")
.load("s3://my-bucket/landing/events/")
.writeStream
.option("checkpointLocation", "s3://my-bucket/_checkpoints/events/")
.trigger(availableNow=True)
.toTable("main.bronze.events")
)
ポイントは3つです。(1) format("cloudFiles") が Auto Loader の指定、(2) cloudFiles.schemaLocation に推論したスキーマの保存先を指定、(3) checkpointLocation が既処理ファイルの記録場所です。書き込み先は toTable() で UC 管理テーブルを指定します。
Auto Loader の強みは、スキーマの扱いが高機能なことです。3つの概念を区別して押さえます。
スキーマを明示しなくても、Auto Loader が初回にファイルをサンプリングしてスキーマを推論し、cloudFiles.schemaLocation に保存します。以降の実行はこの保存済みスキーマを起点に動くため、実行のたびに推論結果が揺れることはありません。
運用中に新しい列が現れたときの挙動は cloudFiles.schemaEvolutionMode で制御します。
| モード | 新しい列が現れたときの挙動 |
|---|---|
| addNewColumns(スキーマ推論時の既定) | ストリームを一度停止してスキーマに新列を追加。再起動後は新列込みで処理が続く |
| rescue | スキーマは変えず、新列のデータを rescued data column に退避して処理を続行 |
| failOnNewColumns | ストリームを失敗させる。人が確認してスキーマを更新するまで再開しない |
| none | 新列を無視する(データは捨てられる) |
📝 試験のポイント
addNewColumns の「新列を検知するとストリームがいったん停止し、再起動後に新列が反映される」という挙動は問われやすい細部です。Lakeflow Jobs 上で実行していればジョブのリトライ設定により自動的に再起動され、運用は途切れません。また rescued data column(既定の列名 _rescued_data)は「スキーマに合わなかったデータを捨てずに JSON として退避する列」で、型不一致や新列のデータ救済に使われます。
書き込み先の Delta テーブルにはスキーマ強制が働きます。テーブルのスキーマに合わないデータは書き込み時に拒否されるのが原則です(第3章で詳述)。Auto Loader は「読み取り側でスキーマを推論・進化させ、合わないデータは rescued data column に退避する」ことで、スキーマ強制と共存しながらデータを失わずに取り込みます。スキーマを固定したい場合は .schema(...) で明示指定すれば、推論を行わないスキーマ強制的な運用もできます。
「新しいファイルが来たことをどう知るか」には2つのモードがあります。
| 観点 | directory listing(既定) | file notification |
|---|---|---|
| 検出方法 | ディレクトリを一覧(リスト)して新ファイルを見つける | クラウドのイベント通知サービス(キュー)経由でファイル到着の通知を受け取る |
| セットアップ | 追加設定不要で簡単。ストレージへの読み取り権限だけで動く | 通知・キューのリソース設定が必要(自動セットアップ機能あり)。クラウド側の追加権限が要る |
| 得意な規模 | 小〜中規模のディレクトリ | 大規模(ファイル数が非常に多い、または高頻度で到着する)ディレクトリでも検出コストが増えにくい |
| 指定方法 | 既定(指定不要) | 推奨:cloudFiles.useManagedFileEvents = true(旧方式:cloudFiles.useNotifications = true) |
directory listing はリストのコストがファイル数に比例して増えるため、ディレクトリが巨大になるほど不利になります。file notification は通知ベースなので検出コストが到着ファイル数にしか依存せず、大規模・高頻度のワークロードに向きます。
file notification モードには新旧2つの方式がある点に注意してください。現在の推奨は、Unity Catalog の外部ロケーションでファイルイベント(file events)を有効化したうえで cloudFiles.useManagedFileEvents = true を指定する managed file events 方式です。通知・キューのリソースを Databricks 側が管理するため、ストリームごとにクラウドリソースを個別作成する必要がありません。cloudFiles.useNotifications = true でストリームごとに通知リソースを構成する旧方式(レガシー)も引き続き利用可能ですが、新規構築では managed file events が推奨です。
Auto Loader は常時稼働のストリーミングだけでなく、バッチ的な定期実行にも使えます。trigger(availableNow=True) を指定すると、「起動した時点で未処理のファイルをすべて処理し、終わったら自動停止する」動きになります。
図:Auto Loader による増分取り込みの全体像
2-2 の比較表を判断基準の形でまとめ直します。
📝 試験のポイント
迷ったら「大規模・継続的・スキーマ進化 → Auto Loader」と覚えてください。逆方向の「小規模・SQL・単純バッチ → COPY INTO」とセットで、どちらの言い回しで問われても答えられるようにしておきましょう。
✅ この節のまとめ
schemaLocation に保存。新列への対応は schemaEvolutionMode(既定 addNewColumns:停止→再起動で新列反映)。合わないデータは rescued data column(_rescued_data)に退避される。trigger(availableNow=True) で「未処理分を処理して停止する」バッチ的運用ができ、Lakeflow Jobs のスケジュールと組み合わせるのが典型。問1. Auto Loader が「一度処理したファイルを二度処理しない」ことを実現している仕組みはどれか。
正解:B
Auto Loader は checkpointLocation に処理済みファイルの情報を永続化し、障害からの再開時も含めて exactly-once を保証します。Aのようにソースファイルを削除・変更することはありません(Dも同様に誤り)。Cのような全行照合は行っておらず、ファイル単位の記録で増分性を実現しています。
問2. スキーマ進化モードが既定の addNewColumns のとき、ソースデータに新しい列が現れた。Auto Loader ストリームの挙動として正しいものはどれか。
正解:B
addNewColumns では、新列の検知時にストリームが UnknownFieldException で停止し、保存済みスキーマに新列が追記されます。再起動後は更新済みスキーマで処理が続くため、Lakeflow Jobs のリトライと組み合わせれば運用は自動で回復します。Aは none モード、データを退避しつつ続行するのは rescue モードの挙動です。Cは誤りで、再起動すれば復旧します。Dのような列のマージは行われません。
問3. 1時間あたり数万ファイルが到着する大規模なランディングディレクトリを Auto Loader で取り込むと、ファイル検出のコストと遅延が問題になってきた。最も適切な対策はどれか。
正解:A
directory listing はディレクトリ内のファイル数に比例してリスティングのコストが増えます。file notification モードならクラウドのイベント通知で到着を知るため、大規模・高頻度でも検出コストが増えにくくなります(なお現在の推奨は、外部ロケーションでファイルイベントを有効化して cloudFiles.useManagedFileEvents = true を使う managed file events 方式で、useNotifications は旧方式です)。Bは逆方向で、COPY INTO はこの規模には向きません。Cはチェックポイントの削除により全ファイルの再処理(重複)を招く危険な操作です。Dは実行スタイルの変更にすぎず、検出方式のコスト構造は変わりません。
問4. Auto Loader のパイプラインを「1日1回、その時点までの未処理ファイルを処理したら停止する」形で動かし、クラスタの常時稼働コストを避けたい。正しい設定はどれか。
trigger(processingTime="24 hours") を指定して常時稼働させるtrigger(availableNow=True) を指定し、Lakeflow Jobs で日次スケジュール実行する正解:B
availableNow=True は「未処理分をすべて処理して自動停止する」トリガーで、Lakeflow Jobs のスケジュールと組み合わせると常時稼働なしの増分バッチになります。Aは24時間間隔で動くもののストリーム自体は稼働し続けるため、コスト回避の要件に合いません。Cは誤りで、Auto Loader はバッチ的運用が公式にサポートされた使い方です。Dは誤りで、チェックポイントは増分性の要であり無効化できません。