テーマ切替
Associate 学習教材 ② データ取り込み
Databricks 公式ドキュメント(日本語版)の内容を、重要用語・概念を漏らさずまとめた自習用教材です。 参照した公式ページ:
- Auto Loader とは(https://docs.databricks.com/ja/ingestion/cloud-object-storage/auto-loader/index.html)
- 一般的なデータ読み込みパターン / Auto Loader patterns(https://learn.microsoft.com/ja-jp/azure/databricks/ingestion/auto-loader/patterns)
- COPY INTO を使用してデータを読み込む(https://learn.microsoft.com/ja-jp/azure/databricks/ingestion/copy-into/)
- COPY INTO を使用した一般的なデータ読み込みパターン(https://learn.microsoft.com/ja-jp/azure/databricks/ingestion/copy-into/examples)
- (補足)COPY INTO SQL 言語リファレンス(https://learn.microsoft.com/ja-jp/azure/databricks/sql/language-manual/delta-copy-into)
- (補足)Lakeflow Connect / インジェスト概要(https://learn.microsoft.com/ja-jp/azure/databricks/ingestion/)
1. このセクションの概要
Databricks における「データ取り込み(Ingestion)」とは、クラウドオブジェクトストレージやメッセージバスなどの外部ソースから、データを Databricks(主に Delta Lake テーブル)へ読み込む処理です。Associate 試験ではとくに クラウドオブジェクトストレージからの増分取り込み(incremental ingestion) を実現する 2 つの中核機能が重要です。
- Auto Loader(自動ローダー) … 新しいファイルがクラウドストレージに到着すると、追加設定なしで段階的(incrementally)かつ効率的に処理する Structured Streaming ソース(
cloudFiles)。大規模・継続的な取り込みに適する。 - COPY INTO … ファイルの場所から Delta テーブルへデータを読み込む SQL コマンド。再試行可能(retryable)で冪等(idempotent)。SQL ユーザーや小〜中規模のバッチ取り込みに適する。
インジェスト製品の 3 レイヤー(カスタマイズ性 → 自動化の順)
| レイヤー | 説明 |
|---|---|
| 構造化ストリーミング(Structured Streaming) | Apache Spark のストリーミングエンジン。エンドツーエンドのフォールトトレランスと **1 回だけの処理(exactly-once)**を保証。最もカスタマイズ性が高い。 |
| Lakeflow パイプライン(Lakeflow Pipelines) | 構造化ストリーミングを拡張した宣言型フレームワーク。オーケストレーション、監視、データ品質、エラー処理を管理。自動化が多くオーバーヘッドが少ない。 |
| マネージドコネクタ(Managed Connectors) | Lakeflow パイプライン上に構築。Salesforce や SQL Server などのソースに対し、認証・CDC・自動スキーマ進化・自動再試行などを提供。最も自動化されている。 |
Databricks は「最も管理されたレイヤーから始め、要件を満たさない場合に下位(よりカスタマイズ可能な)レイヤーに落とす」ことを推奨。
インジェストのスケジュール
| 利用シーン | パイプラインモード |
|---|---|
| バッチインジェスト | トリガー(triggered): スケジュールに従って、または手動でトリガーされたときに新しいデータを処理 |
| ストリーミングデータ取り込み | 継続的(continuous): ソースに到着した新しいデータを処理 |
重要な傾向: Databricks は、SQL ユーザーに対して
COPY INTOよりもストリーミングテーブル(Streaming Tables) /CREATE STREAMING TABLEを推奨するようになっている(よりスケーラブルで堅牢なため)。ただし試験ではCOPY INTOの仕様も問われる。
2. 重要用語集
| 用語(日本語) | English | 説明 |
|---|---|---|
| 自動ローダー | Auto Loader | 新規ファイルの到着を自動検知し、段階的・効率的に処理する Structured Streaming ソース。 |
| cloudFiles | cloudFiles | Auto Loader を利用するための Structured Streaming ソース/フォーマット名。format("cloudFiles") で指定。 |
| 増分取り込み | Incremental Ingestion | 既に処理したファイルは飛ばし、新規到着ファイルのみを取り込む方式。 |
| スキーマ推論 | Schema Inference | 到着データからカラム構造・型を自動検出する機能。 |
| スキーマ進化 | Schema Evolution | 時間とともにスキーマ(新列・型変更)へ自動追従する機能。 |
| スキーマの場所 | Schema Location (cloudFiles.schemaLocation) | 推論したスキーマの変遷を保存するディレクトリ。スキーマ進化に必須。 |
| スキーマヒント | Schema Hints (cloudFiles.schemaHints) | 特定カラムの型を明示的に指定して推論を補正する仕組み。 |
| レスキューデータ列 | Rescued Data Column (_rescued_data) | スキーマと一致しない・解析できないデータを取りこぼさず退避する特別列。 |
| チェックポイント | Checkpoint (checkpointLocation) | 処理済みファイル等の状態を保存する場所。障害時の再開・重複防止に使う。 |
| RocksDB | RocksDB | チェックポイントのファイルメタデータを永続化するスケーラブルなキーバリューストア。 |
| 1 回のみ処理 | Exactly-once | データが 1 回だけ処理されることを保証するセマンティクス。 |
| ディレクトリ一覧モード | Directory Listing Mode | ディレクトリをスキャンして新規ファイルを検出する既定モード。 |
| ファイル通知モード | File Notification Mode | クラウドのイベント通知を使って新規ファイルを検出するモード。大規模で低コスト。 |
| バックフィル | Backfill | 既存(過去)ファイルの取り込み。非同期に実行可能。 |
| glob パターン | Glob Pattern | * ? [abc] {ab,cd} などでパス・ファイルをフィルタする表記。 |
| 厳密グロッバー | Strict Globber (cloudFiles.useStrictGlobber) | 他の Spark ファイルソースと同じグロブ動作にする設定。 |
| バリアント型 | VARIANT | 半構造化データを 1 列で柔軟に保持する型。スキーマ/型変更に強い。 |
| COPY INTO | COPY INTO | ファイルの場所から Delta テーブルへ読み込む再試行可能・冪等な SQL コマンド。 |
| 冪等性 | Idempotency | 同じ操作を繰り返しても結果が変わらない性質(既読ファイルはスキップ)。 |
| フォーマットオプション | FORMAT_OPTIONS | ファイル形式リーダー(Spark DataFrameReader)に渡すオプション。 |
| コピーオプション | COPY_OPTIONS | COPY INTO の動作を制御するオプション(force, mergeSchema など)。 |
| 検証モード | VALIDATE | 実際に書き込まず、解析・スキーマ・制約を検証するモード。 |
| ストリーミングテーブル | Streaming Table | CREATE STREAMING TABLE / read_files で作る増分取り込み用テーブル。 |
| read_files | read_files() | SQL で Auto Loader を使うためのテーブル値関数。 |
| Delta Lake | Delta Lake | Databricks の既定テーブル形式。ACID トランザクションを提供。 |
| 削除ベクトル | Deletion Vectors | 行の論理削除を記録する Delta の機能。COPY INTO は本設定を尊重。 |
3. 詳細解説
3-1. Auto Loader とは / 仕組み(増分検知・チェックポイント・exactly-once)
定義: Auto Loader は「新しいデータファイルがクラウドストレージに到着すると、追加設定なしで段階的(incrementally)かつ効率的に処理する」Databricks の機能。実体は Structured Streaming ソースで、cloudFiles というインターフェース(フォーマット)として提供される。入力ディレクトリのパスを指定すると、新規到着ファイルを自動検出・処理する。
サポート対象
- クラウドストレージ:
- Amazon S3(
s3://) - Azure ADLS Gen2(
abfss://) - Google Cloud Storage(
gs://) - Azure Blob Storage(
wasbs://) - Unity Catalog ボリューム(
/Volumes/)
- Amazon S3(
- ファイル形式: JSON, CSV, XML, PARQUET, AVRO, ORC, TEXT, BINARYFILE(圧縮ファイルにも対応)。
仕組みの核となる要素
チェックポイント(Checkpoint)
- 処理済みファイルのメタデータを、スケーラブルなキーバリューストアである RocksDB に永続化する。
- これにより、どのファイルを処理済みかを記録し、重複処理を防止する。
- 障害が発生しても、中断した箇所から再開できる(フォールトトレランス)。
- 書き込み側では
checkpointLocationオプションで場所を指定する。
exactly-once(1 回のみ)セマンティクス
- データが 1 回だけ処理されることを保証する。障害時も Delta Lake への書き込みでこの保証が維持される。
- 自動で実現され、ユーザーによる状態管理は不要。
増分処理(Incremental Ingestion)
- 新規到着ファイルのみを段階的に取り込む。
- バックフィル(Backfill)(既存ファイルの取り込み)は非同期に実行できるため、コンピューティングリソースの無駄を避けられる。
従来方式(spark.readStream.format(fileFormat).load(dir))に対する優位性
- スケーラビリティ: 数十億のファイルを効率的に検出できる。
- パフォーマンス: ファイル検出コストが、ディレクトリ全体の再スキャンに依存しにくい。
- コスト: ネイティブクラウド API を用い、ストレージ/リスティングのコストを削減。
注意点
- 順序保証なし: Auto Loader は「ファイルの検出または処理の順序を保証しない」。順序が重要な場合は、データ内のタイムスタンプ比較やソフト削除(soft-delete)戦略で対応する。
- AUTO CDC / 順不同到着: Lakeflow Pipelines と組み合わせる場合、順不同で到着するファイルに対応するため、
pipelines.cdc.tombstoneGCThresholdInSecondsで削除レコード(tombstone)の保持期間を設定する。
3-2. Auto Loader のスキーマ推論・スキーマ進化
スキーマ推論(Schema Inference): 到着データから自動的にカラム構造を検出する。既定では、列は文字列型(string)として推論される(明示的に型推論を有効にしない限り)。
スキーマ進化(Schema Evolution): 時間の経過とともにスキーマを推論・進化させる能力。新規カラムや型変更に追従する。スキーマ進化には推論スキーマの履歴を保存する cloudFiles.schemaLocation(スキーマの場所)が必要。
重要オプション
| オプション | 説明 |
|---|---|
cloudFiles.format | ソースファイル形式(json, csv, parquet, binaryFile など)。必須。 |
cloudFiles.schemaLocation | 推論したスキーマの変遷を保存するディレクトリ。スキーマ推論/進化に必要。 |
cloudFiles.schemaEvolutionMode | スキーマ進化の挙動。rescue(新フィールドや型不一致を _rescued_data に収集)、failOnNewColumns(新列出現時にストリームを失敗させる)など。 |
cloudFiles.schemaHints | 特定カラムの型を明示(例: headers map<string,string>, statusCode SHORT)。 |
cloudFiles.inferColumnTypes | true にすると、ネストされたデータやその他の列型を推論する(既定は文字列推論)。 |
rescuedDataColumn / _rescued_data | スキーマと一致しないデータ・解析エラーを退避する列。データ損失防止。 |
mergeSchema(書き込み側 writeStream のオプション) | 進化したスキーマをターゲット Delta テーブルへマージ。 |
既定の挙動(簡単な ETL パターン)
- 既定では スキーマは文字列型として推論される。
- 解析エラー(すべて文字列のままなら基本発生しない)は
_rescued_dataに移動する。 - 新しい列が出現するとストリームは一旦失敗し、スキーマを進化させる。そのため、ソーススキーマ変更時に自動でストリームを再起動する運用(Databricks ジョブでの実行)が推奨される。
設計上の使い分け
- スキーマが不明でとにかく損失なく取り込みたい → スキーマ推論を有効化 +
mergeSchema+_rescued_data。 - スキーマは既知だが予期しないデータを捕捉したい →
.schema(expected_schema)を与えつつrescuedDataColumn(またはschemaEvolutionMode=rescue)。 - 新列が出たら止めたい →
cloudFiles.schemaEvolutionMode = failOnNewColumns。 - ネストされた JSON →
cloudFiles.inferColumnTypes=true、または最上位を文字列のまま取り込み、半構造化アクセス API(col:field.subfield::int構文)で後段変換。 - 半構造化で最も堅牢 → VARIANT 型として 1 列で取り込む(大文字小文字・NULL・型変更に強い)。
3-3. Auto Loader のファイル検知モード(ディレクトリ一覧 / ファイル通知)と実践パターン
ファイル検知モード
- ディレクトリ一覧モード(Directory Listing Mode)
- 既定モード。入力ディレクトリをスキャンして新規ファイルを検出する。
- セットアップが簡単だが、ファイル数が増えるとリスティングコスト・時間が増加する傾向。
- ファイル通知モード(File Notification Mode)
- 推奨モード(大規模時)。クラウドが提供するファイルイベント通知(通知サービス+キュー)を利用。
- ディレクトリのリストアップを回避するため、クラウド費用を削減し、大量ファイルでもスケールする。
glob パターンによるフィルタ
パスに glob を書くことでディレクトリ/ファイルをフィルタできる。
| パターン | 意味 |
|---|---|
? | 任意の 1 文字 |
* | 0 個以上の文字 |
[abc] | 文字セット {a,b,c} の 1 文字 |
[a-z] | 文字範囲 {a…z} の 1 文字 |
[^a] | セット/範囲 {a} 以外の 1 文字(^ は左角かっこの直後) |
{ab,cd} | 文字列セット {ab, cd} のいずれか |
{ab,c{de,fh}} | {ab, cde, cfh} のいずれか |
- プレフィックスフィルタは
path(.load(...))で指定する。 - サフィックスフィルタ(例:
*.pngだけ)はpathではできず、pathGlobFilterオプションで指定する。 - Auto Loader の既定グロブ動作は他の Spark ファイルソースと異なる。標準の Spark 動作に揃えたい場合は
.option("cloudFiles.useStrictGlobber", "true")を付ける。
実践パターン(patterns ページの要点)
- VARIANT として取り込む: すべてのデータをターゲットテーブルの 1 つの
VARIANT列に読み込む。スキーマ・型変更に柔軟で、大文字小文字と NULL を保持するため、多くのシナリオで堅牢。 - 簡単な ETL(データ損失なし): スキーマ推論を有効化し
mergeSchemaを付けて Delta へ書き込む。ソーススキーマ変更時にジョブで自動再起動する運用が推奨。 - 適切に構造化されたデータの損失防止:
rescuedDataColumnを使い、想定外データを_rescued_dataに集める。 - 柔軟な半構造化パイプライン:
schemaHintsでベンダー提供の一部フィールド型を固定しつつ、スキーマ進化に任せる。 - ネスト JSON の変換: 最上位 JSON は文字列推論されるので、
selectExprでtags:page.name、tags:page.id::intのように半構造化アクセスして変換。 - ネスト JSON の推論:
cloudFiles.inferColumnTypes=true。 - ヘッダーなし CSV:
.schema(<schema>)を与えrescuedDataColumnで損失防止。 - ヘッダーあり CSV:
header=true+ スキーマ適用 +rescuedDataColumn。 - 画像・バイナリを ML 用に取り込む:
cloudFiles.format = binaryFileで Delta に取り込み、分散推論へ。 - Lakeflow パイプライン構文: Python は
@dp.tableデコレータ、SQL はCREATE OR REFRESH STREAMING TABLE ... AS SELECT * FROM STREAM read_files(...)。- Lakeflow パイプラインは Auto Loader 使用時にスキーマとチェックポイントのディレクトリを自動構成・管理する。手動構成したディレクトリは完全更新(full refresh)の影響を受けないため、自動構成ディレクトリの使用が推奨。
3-4. COPY INTO の使い方・冪等性・オプション
定義: COPY INTO は、ファイルの場所から Delta テーブルへデータを読み込む SQL コマンド。**再試行可能(retryable)かつ冪等(idempotent)**で、既に読み込み済みのファイルは後続の実行でスキップされる(ファイルが読み込み後に変更されていてもスキップ)。
主な機能
- S3 / ADLS / ABFS / GCS / Unity Catalog ボリュームなど、クラウドストレージのファイル・フォルダフィルタを簡単に構成。
- 複数フォーマット対応: CSV, JSON, XML, Avro, ORC, Parquet, TEXT, BINARYFILE。
- 既定で冪等な 1 回限りの処理。
- ターゲットテーブルスキーマの推論・マッピング・マージ・進化。
構文(言語リファレンス)
sql
COPY INTO target_table [ BY POSITION | ( col_name [, ...] ) ]
FROM { source_clause | ( SELECT expression_list FROM source_clause ) }
FILEFORMAT = data_source
[ VALIDATE [ ALL | num_rows ROWS ] ]
[ FILES = ( file_name [, ...] ) | PATTERN = glob_pattern ]
[ FORMAT_OPTIONS ( { data_source_reader_option = value } [, ...] ) ]
[ COPY_OPTIONS ( { copy_option = value } [, ...] ) ]
source_clause
source [ WITH ( [ CREDENTIAL { credential_name | (temporary_credential_options) } ]
[ ENCRYPTION (encryption_options) ] ) ]句・パラメータの意味
| 句 / オプション | 説明 |
|---|---|
target_table | 既存の Delta テーブルを指定(時間指定やオプション指定は不可)。delta. `/path` `` 形式で外部ロケーションへの書き込みも可。 |
BY POSITION / (col, ...) | ソース列とターゲット列を序数位置で照合(ヘッダーなし CSV のみ、headers=false 必須)。IDENTITY/GENERATED 列は照合対象外。列数不一致はエラー。 |
SELECT expression_list | コピー前にソースから列・式を選択・変換(キャスト、名前変更、定数追加、ウィンドウ式など)。 |
FILEFORMAT = data_source | ソースファイル形式(CSV, JSON, AVRO, ORC, PARQUET, TEXT, BINARYFILE)。 |
VALIDATE | データを検証するがテーブルには書き込まない(Databricks Runtime 10.4 LTS 以降)。解析可否・スキーマ一致/展開要否・NULL 許容と CHECK 制約を検証。VALIDATE 15 ROWS で行数指定、50 未満だとプレビュー返却。 |
FILES = (...) | 読み込むファイル名リスト(最大 1000 ファイル)。PATTERN と併用不可。 |
PATTERN = glob | 読み込むファイルを glob で指定。FILES と併用不可。 |
FORMAT_OPTIONS (...) | 形式リーダー(Spark DataFrameReader)に渡すオプション(header, delimiter, multiLine, inferSchema, mergeSchema, ignoreCorruptFiles など)。 |
COPY_OPTIONS (...) | COPY INTO の動作制御。force(既定 false、true で冪等性を無効化し既読でも再読み込み)、mergeSchema(既定 false、true で受信データに応じてターゲットスキーマを展開)。 |
WITH (CREDENTIAL ... ENCRYPTION ...) | ソースへのアクセス資格情報・暗号化を指定。ADLS/Blob は AZURE_SAS_TOKEN、S3 は AWS_ACCESS_KEY/AWS_SECRET_KEY/AWS_SESSION_TOKEN、暗号化は S3 の TYPE='AWS_SSE_C' + MASTER_KEY 等。 |
冪等性(Idempotency)のポイント
- 同じ
COPY INTOを繰り返しスケジュール実行しても、新規データのみが Delta テーブルに読み込まれる。 - 既読ファイルは(内容が後から変わっていても)スキップされる。
COPY_OPTIONS('force'='true')を指定すると冪等性が無効化され、既読ファイルも再読み込みされる。
スキーマ関連
- Databricks Runtime 11.3 LTS 以降、スキーマ展開に対応する形式ではターゲットテーブルのスキーマ定義を省略可能(空のプレースホルダ Delta テーブルを作り、
FORMAT_OPTIONS('mergeSchema'='true')+COPY_OPTIONS('mergeSchema'='true')で推論させる)。 - 空の Delta テーブルは
COPY INTO以外では使えない(INSERT INTO/MERGE INTOはスキーマレステーブルへの書き込み不可)。COPY INTOでデータが入るとクエリ可能になる。 FORMAT_OPTIONSのinferSchema(CSV 等の型推論)、mergeSchema(複数ソースファイル間でスキーマをマージ推論)。ソースファイルのスキーマが同一ならmergeSchema=false(既定)推奨。COPY_OPTIONSのmergeSchemaはターゲット Delta テーブル側のスキーマ展開可否。
運用上の注意
- 削除ベクトル(Deletion Vectors):
COPY INTOはワークスペースの削除ベクトル設定を尊重。有効時、Databricks Runtime 14.0 以降で実行するとターゲットテーブルで削除ベクトルが有効化され、DBR 11.3 LTS 以降でないとクエリがブロックされる可能性。 - 破損ファイル:
FORMAT_OPTIONS('ignoreCorruptFiles'='true')(DBR 11.3 LTS 以降)でスキップ。結果にnum_skipped_corrupt_files、DESCRIBE HISTORYのoperationMetrics.numSkippedCorruptFilesに表示。破損ファイルは追跡されないので、修正後の再実行で再読み込み可能。 - メタデータのクリーンアップ: DBR 15.2 以降、
VACUUMでCOPY INTOが作成した参照されないメタデータファイルを削除可能。 - 同時実行: 同一テーブルへの
COPY INTOの同時呼び出しは可能(異なる入力ファイル集合である限りトランザクション競合なく成功)。ただし性能向上目的の同時実行は非推奨(1 コマンドで複数ファイルを扱う方が高速)。ごく多数のファイルを含む巨大ディレクトリは可能なら Auto Loader を推奨。
3-5. Auto Loader と COPY INTO の使い分け
| 観点 | Auto Loader(cloudFiles) | COPY INTO(SQL) |
|---|---|---|
| 実体 | Structured Streaming ソース | SQL コマンド |
| 主な言語 | Python / Scala(SQL は read_files/ストリーミングテーブル経由) | SQL(Python/Scala/R からも spark.sql で実行可) |
| ファイル規模 | 数百万〜数十億ファイルでもスケール | 少〜中規模。1 回で最大 1000 ファイル(FILES)指定。巨大ディレクトリは非推奨 |
| 状態管理 | チェックポイント(RocksDB)で処理済みを追跡 | Delta のトランザクションログで読み込み済みファイルを追跡 |
| 保証 | exactly-once(自動) | 冪等・再試行可能(既読スキップ) |
| ファイル検知 | ディレクトリ一覧 / ファイル通知モード | 実行時にソースをリスト |
| 継続処理 | ストリーミング(継続的)に強い | バッチ的な定期実行に向く |
| 推奨シーン | 継続的・大規模・低レイテンシ、ファイルが継続到着 | 定期バッチ、SQL 中心、ファイル数が数千程度 |
選択の指針
- 取り込むファイル数が数千を超える/継続的に到着するなら Auto Loader。
- SQL ユーザーで、周期的なバッチ取り込みや手軽さ重視なら COPY INTO。
- ただし公式は、SQL ユーザーには**ストリーミングテーブル(
CREATE STREAMING TABLE+read_files)**を、よりスケーラブル・堅牢としてCOPY INTOより推奨する方向。read_filesは内部的に Auto Loader を利用する。
4. 構文・コード例
4-1. Auto Loader: 基本の readStream(Python)
python
df = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "json") \
.option("cloudFiles.schemaLocation", "<path-to-schema-location>") \
.load("/Volumes/catalog_name/schema_name/volume_name/source_data")
df.writeStream \
.option("mergeSchema", "true") \
.option("checkpointLocation", "<path-to-checkpoint>") \
.start("<path_to_target>")4-2. Auto Loader: glob フィルタ / サフィックスフィルタ
python
# プレフィックス(path 内の * )
df = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "<format>") \
.schema(schema) \
.load("/Volumes/catalog_name/schema_name/volume_name/*/files")
# サフィックス(pathGlobFilter で .png のみ)
df = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "binaryFile") \
.option("pathGlobfilter", "*.png") \
.load("/Volumes/catalog_name/schema_name/volume_name/path")4-3. Auto Loader: スキーマ進化モード / スキーマヒント
python
# 想定スキーマ + 新フィールド・型不一致を _rescued_data に収集
spark.readStream.format("cloudFiles") \
.schema(expected_schema) \
.option("cloudFiles.format", "json") \
.option("cloudFiles.schemaEvolutionMode", "rescue") \
.load("...source...") \
.writeStream.option("checkpointLocation", "...").start("...target...")
# 新列が来たら停止させる
.option("cloudFiles.schemaEvolutionMode", "failOnNewColumns")
# スキーマヒント(型を明示)
.option("cloudFiles.schemaHints", "headers map<string,string>, statusCode SHORT")4-4. Auto Loader: ネスト JSON の型推論 / 半構造化アクセス
python
# ネスト構造も推論
spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "json") \
.option("cloudFiles.schemaLocation", "<path-to-checkpoint>") \
.option("cloudFiles.inferColumnTypes", "true") \
.load("...nested_json...")
# 半構造化アクセス(selectExpr)
.selectExpr("*",
"tags:page.name", # {"tags":{"page":{"name":...}}}
"tags:page.id::int", # int にキャスト
"tags:eventType")4-5. Auto Loader: SQL / Lakeflow パイプライン構文(read_files)
sql
CREATE OR REFRESH STREAMING TABLE booking_updates
AS SELECT * FROM STREAM read_files(
"/Volumes/my_catalog/my_schema/my_volume/wanderbricks/booking_updates",
format => "json",
multiLine => true,
inferColumnTypes => true
);python
@dp.table
def reviews():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("multiLine", "true")
.load("/Volumes/my_catalog/my_schema/my_volume/wanderbricks/reviews")
)4-6. COPY INTO: スキーマを定義して JSON をロード
sql
CREATE TABLE <catalog>.<schema>.booking_updates_upload (
booking_id BIGINT,
user_id BIGINT,
status STRING,
total_amount DOUBLE
);
COPY INTO <catalog>.<schema>.booking_updates_upload
FROM '/Volumes/<catalog>/<schema>/<volume>/wanderbricks/booking_updates'
FILEFORMAT = JSON
FORMAT_OPTIONS ('multiLine' = 'true');4-7. COPY INTO: スキーマレステーブルに推論ロード(mergeSchema)
sql
CREATE TABLE IF NOT EXISTS <catalog>.<schema>.booking_updates_schemaless;
COPY INTO <catalog>.<schema>.booking_updates_schemaless
FROM '/Volumes/<catalog>/<schema>/<volume>/wanderbricks/booking_updates'
FILEFORMAT = JSON
FORMAT_OPTIONS ('mergeSchema' = 'true', 'multiLine' = 'true')
COPY_OPTIONS ('mergeSchema' = 'true');4-8. COPY INTO: 冪等性の確認(同じ FILES を 2 回実行)
sql
COPY INTO my_json_data
FROM 'abfss://container@storageAccount.dfs.core.windows.net/base/path'
FILEFORMAT = JSON
FILES = ('f1.json', 'f2.json', 'f3.json', 'f4.json', 'f5.json');
-- 2 回目は既に読み込み済みのためデータはコピーされない
COPY INTO my_json_data
FROM 'abfss://container@storageAccount.dfs.core.windows.net/base/path'
FILEFORMAT = JSON
FILES = ('f1.json', 'f2.json', 'f3.json', 'f4.json', 'f5.json');4-9. COPY INTO: PATTERN と SELECT で変換しながら CSV をロード
sql
COPY INTO target_table
FROM (SELECT key, index, textData, 'constant_value'
FROM 'abfss://container@storageAccount.dfs.core.windows.net/base/path')
FILEFORMAT = CSV
PATTERN = 'folder1/file_[a-g].csv'
FORMAT_OPTIONS('header' = 'true');
-- ヘッダーなし CSV をキャスト・列名変更してロード
COPY INTO target_table
FROM (SELECT _c0::bigint key, _c1::int index, _c2 textData
FROM 'abfss://container@storageAccount.dfs.core.windows.net/base/path')
FILEFORMAT = CSV
PATTERN = 'folder1/file_[a-g].csv';4-10. COPY INTO: Avro を SQL 式付きでロード / 破損ファイル無視 / VALIDATE
sql
-- Avro を式変換しながらロード
COPY INTO my_delta_table
FROM (SELECT to_date(dt) dt, event as measurement, quantity::double
FROM 'abfss://container@storageAccount.dfs.core.windows.net/base/path')
FILEFORMAT = AVRO;
-- 破損ファイルをスキップ(結果に num_skipped_corrupt_files)
COPY INTO my_table
FROM '/path/to/files'
FILEFORMAT = <format>
FORMAT_OPTIONS ('ignoreCorruptFiles' = 'true');
-- 書き込まずに検証のみ
COPY INTO my_table
FROM '/path/to/files'
FILEFORMAT = <format>
VALIDATE ALL
FORMAT_OPTIONS ('mergeSchema' = 'true');5. 試験で問われるポイント
- Auto Loader = 増分ファイル取り込みであり、実体は Structured Streaming ソース
cloudFiles。spark.readStream.format("cloudFiles")で使う。 - Auto Loader は **チェックポイント(RocksDB に永続化)**で処理済みファイルを追跡し、exactly-once を自動で保証する。ユーザーは状態管理不要。
- Auto Loader のファイル検知は 2 モード: ディレクトリ一覧(既定)とファイル通知(大規模・低コストで推奨)。違いを説明できること。
- スキーマ推論は既定で文字列型。型推論には
cloudFiles.inferColumnTypes=true。スキーマ進化にはcloudFiles.schemaLocationが必要。 _rescued_data(rescuedDataColumn) はスキーマ不一致・解析失敗データの退避先で、データ損失防止が目的。cloudFiles.schemaEvolutionModeの値(rescue/failOnNewColumnsなど)と挙動。- COPY INTO は冪等(idempotent)・再試行可能で、既読ファイルはスキップ。
COPY_OPTIONS('force'='true')で冪等性を無効化して再読み込み。 - COPY INTO のターゲットは既存 Delta テーブル(
INSERT INTO/MERGE INTOはスキーマレステーブルに書けない)。 FILEFORMAT(CSV/JSON/AVRO/ORC/PARQUET/TEXT/BINARYFILE)、FILES(最大 1000)とPATTERN(glob)は併用不可。FORMAT_OPTIONS(リーダー設定: header, multiLine, inferSchema, mergeSchema, ignoreCorruptFiles)とCOPY_OPTIONS(force, mergeSchema)の違い。- VALIDATE モードは書き込まずに解析・スキーマ・制約を検証する。
- Auto Loader と COPY INTO の使い分け(大規模・継続 → Auto Loader、SQL 定期バッチ → COPY INTO、ただし公式はストリーミングテーブルを推奨)。
- 対応クラウド(S3 / ADLS / GCS / Blob / Unity Catalog Volumes)とパススキーム(
s3://,abfss://,gs://,wasbs://,/Volumes/)。 - Auto Loader はファイルの処理順序を保証しない。
6. 理解度チェックリスト
- [ ] Auto Loader が Structured Streaming ソースであり
cloudFilesフォーマットで使うことを説明できる - [ ] Auto Loader がチェックポイント(RocksDB)と exactly-once でどう重複を防ぐか説明できる
- [ ] ディレクトリ一覧モードとファイル通知モードの違い・使いどころを言える
- [ ]
cloudFiles.schemaLocationの役割と、スキーマ推論/進化の関係を説明できる - [ ] 既定のスキーマ推論が文字列型であり、
inferColumnTypesで型推論できることを知っている - [ ]
_rescued_data(rescuedDataColumn)の目的(データ損失防止)を説明できる - [ ]
cloudFiles.schemaEvolutionModeのrescueとfailOnNewColumnsの挙動を区別できる - [ ]
schemaHintsで型を明示できることを知っている - [ ] glob の
path(プレフィックス)とpathGlobFilter(サフィックス)の使い分けができる - [ ] Auto Loader が対応するクラウドストレージとファイル形式を列挙できる
- [ ]
COPY INTOが冪等・再試行可能で既読ファイルをスキップすることを説明できる - [ ]
COPY_OPTIONS('force'='true')の効果(冪等性無効化)を知っている - [ ]
COPY INTOのターゲットが既存 Delta テーブルであること、スキーマレステーブルの制約を知っている - [ ]
FILES(最大 1000, PATTERN と併用不可)とPATTERN(glob)の違いを言える - [ ]
FORMAT_OPTIONSとCOPY_OPTIONSの役割の違いを説明できる - [ ]
VALIDATEモードで何が検証されるか(解析・スキーマ・NULL/CHECK 制約)を言える - [ ]
ignoreCorruptFilesとnum_skipped_corrupt_filesの関係を知っている - [ ] Auto Loader と COPY INTO の使い分け基準(規模・継続性・言語)を説明できる
- [ ] SQL で Auto Loader を使う
read_files/CREATE STREAMING TABLEを書ける - [ ]
COPY INTOの基本 SQL 構文(FROM / FILEFORMAT / FORMAT_OPTIONS / COPY_OPTIONS)を書ける