Skip to content

Associate 学習教材 ② データ取り込み

Databricks 公式ドキュメント(日本語版)の内容を、重要用語・概念を漏らさずまとめた自習用教材です。 参照した公式ページ:


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 ソース。
cloudFilescloudFilesAuto 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)処理済みファイル等の状態を保存する場所。障害時の再開・重複防止に使う。
RocksDBRocksDBチェックポイントのファイルメタデータを永続化するスケーラブルなキーバリューストア。
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 INTOCOPY INTOファイルの場所から Delta テーブルへ読み込む再試行可能・冪等な SQL コマンド。
冪等性Idempotency同じ操作を繰り返しても結果が変わらない性質(既読ファイルはスキップ)。
フォーマットオプションFORMAT_OPTIONSファイル形式リーダー(Spark DataFrameReader)に渡すオプション。
コピーオプションCOPY_OPTIONSCOPY INTO の動作を制御するオプション(force, mergeSchema など)。
検証モードVALIDATE実際に書き込まず、解析・スキーマ・制約を検証するモード。
ストリーミングテーブルStreaming TableCREATE STREAMING TABLE / read_files で作る増分取り込み用テーブル。
read_filesread_files()SQL で Auto Loader を使うためのテーブル値関数。
Delta LakeDelta LakeDatabricks の既定テーブル形式。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/
  • ファイル形式: 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.inferColumnTypestrue にすると、ネストされたデータやその他の列型を推論する(既定は文字列推論)。
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 は文字列推論されるので、selectExprtags:page.nametags: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_OPTIONSinferSchema(CSV 等の型推論)、mergeSchema(複数ソースファイル間でスキーマをマージ推論)。ソースファイルのスキーマが同一なら mergeSchema=false(既定)推奨。
  • COPY_OPTIONSmergeSchema はターゲット Delta テーブル側のスキーマ展開可否。

運用上の注意

  • 削除ベクトル(Deletion Vectors): COPY INTO はワークスペースの削除ベクトル設定を尊重。有効時、Databricks Runtime 14.0 以降で実行するとターゲットテーブルで削除ベクトルが有効化され、DBR 11.3 LTS 以降でないとクエリがブロックされる可能性。
  • 破損ファイル: FORMAT_OPTIONS('ignoreCorruptFiles'='true')(DBR 11.3 LTS 以降)でスキップ。結果に num_skipped_corrupt_filesDESCRIBE HISTORYoperationMetrics.numSkippedCorruptFiles に表示。破損ファイルは追跡されないので、修正後の再実行で再読み込み可能。
  • メタデータのクリーンアップ: DBR 15.2 以降、VACUUMCOPY INTO が作成した参照されないメタデータファイルを削除可能。
  • 同時実行: 同一テーブルへの COPY INTO の同時呼び出しは可能(異なる入力ファイル集合である限りトランザクション競合なく成功)。ただし性能向上目的の同時実行は非推奨(1 コマンドで複数ファイルを扱う方が高速)。ごく多数のファイルを含む巨大ディレクトリは可能なら Auto Loader を推奨。

3-5. Auto Loader と COPY INTO の使い分け

観点Auto Loader(cloudFilesCOPY 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 ソース cloudFilesspark.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.schemaEvolutionModerescuefailOnNewColumns の挙動を区別できる
  • [ ] schemaHints で型を明示できることを知っている
  • [ ] glob の path(プレフィックス)と pathGlobFilter(サフィックス)の使い分けができる
  • [ ] Auto Loader が対応するクラウドストレージとファイル形式を列挙できる
  • [ ] COPY INTO が冪等・再試行可能で既読ファイルをスキップすることを説明できる
  • [ ] COPY_OPTIONS('force'='true') の効果(冪等性無効化)を知っている
  • [ ] COPY INTO のターゲットが既存 Delta テーブルであること、スキーマレステーブルの制約を知っている
  • [ ] FILES(最大 1000, PATTERN と併用不可)と PATTERN(glob)の違いを言える
  • [ ] FORMAT_OPTIONSCOPY_OPTIONS の役割の違いを説明できる
  • [ ] VALIDATE モードで何が検証されるか(解析・スキーマ・NULL/CHECK 制約)を言える
  • [ ] ignoreCorruptFilesnum_skipped_corrupt_files の関係を知っている
  • [ ] Auto Loader と COPY INTO の使い分け基準(規模・継続性・言語)を説明できる
  • [ ] SQL で Auto Loader を使う read_files / CREATE STREAMING TABLE を書ける
  • [ ] COPY INTO の基本 SQL 構文(FROM / FILEFORMAT / FORMAT_OPTIONS / COPY_OPTIONS)を書ける