Azure Databricks の Auto Loader(cloudFiles)はとても便利ですが、「includeExistingFiles と modifiedAfter を設定したのに、なぜか全部のファイルが読み込まれてしまう」という声は少なくありません。本記事では、ADLS Gen2 を前提に、なぜその現象が起こるのか、そしてどう直せばよいのかを、コード例とともに丁寧に解説します。
Auto Loader で意図せず「すべてのファイル」が読まれてしまう問題
想定しているシナリオは次のようなものです。
- データソース:Azure Data Lake Storage Gen2(ADLS Gen2)のサブフォルダー
- 取り込み方法:Azure Databricks Auto Loader(
spark.readStream.format("cloudFiles")) - 要件:
- ジョブ開始時点より前から存在する「既存ファイル」は読みたくない(新着ファイルだけ処理したい)
- 特定の日時以降に更新されたファイルだけを対象にしたい
そのため、次のような設定をしているとします。
cloudFiles.includeExistingFiles = "false"cloudFiles.modifiedAfter = "<日時>"
ところが、こうした設定を入れているにもかかわらず、ディレクトリ配下のファイルが「すべて」処理されてしまう、というのが今回のテーマです。
主な原因と解決アクションの一覧
まずは原因の候補をざっと俯瞰しておきましょう。よくある落とし穴と、その対処を表にまとめます。
| 主な原因 | 解決アクション |
|---|---|
Auto Loader のフォーマット設定が二重指定になっており、後から指定した format() で上書きされている | format("cloudFiles") は一度だけ書き、ファイル形式は option("cloudFiles.format","text") のようにオプションで指定する |
cloudFiles.modifiedAfter の日時書式が ISO-8601(UTC)になっていない | 2025-09-09T00:00:00.000Z のような UTC の ISO-8601 形式で指定する |
| Event Grid 通知モードとの組み合わせで期待どおりのフィルタになっていない | まずは cloudFiles.useNotifications = "false" にしてリストモードで動作を固定し、シンプルな構成で検証する |
| 過去の実行で作成されたチェックポイントが残っており、新しい設定が反映されていない | 設定を変えた場合は checkpointLocation を別パスに変えるか、古いチェックポイントを削除してから再実行する |
この4つを順に潰していくことで、「なぜ全部読まれるのか?」をかなりの確度で切り分けられます。以降ではそれぞれを詳しく見ていきます。
原因1:format() の二重指定で Auto Loader が無効化される
もっとも多いミスが「フォーマットの二重指定」です。次のようなコードになっていないでしょうか。
df_stream = (spark.readStream
.format("cloudFiles") # Auto Loader を指定
.option("cloudFiles.format", "text")
.option("cloudFiles.includeExistingFiles", "false")
.option("cloudFiles.modifiedAfter", "2025-09-09T00:00:00.000Z")
.format("text") # <= ここで再度 format() を呼んでしまう
.load(LANDED_PATH)
)
上記のように format("cloudFiles") の後に format("text") を呼ぶと、Spark の仕様として 最後に呼ばれた format() が有効 になります。
- 結果として Auto Loader の設定はすべて無効化される
- 単なる
spark.readStream.format("text")と同じ扱いになる - そのため、
cloudFiles.includeExistingFilesやcloudFiles.modifiedAfterは完全に無視される
この状態だと、「設定を入れているのに効かない」という見え方になりがちです。
正しいフォーマット指定の方法
Auto Loader を使うときは、ストリーム側の format() は必ず一度だけにし、ファイル形式は cloudFiles.format オプションで指定します。
df_stream = (spark.readStream
.format("cloudFiles") # <= ここだけ
.option("cloudFiles.format", "text") # ファイル形式はオプションで指定
.option("cloudFiles.includeExistingFiles", "false")
.option("cloudFiles.modifiedAfter", "2025-09-09T00:00:00.000Z")
.option("cloudFiles.useNotifications", "false")
.load(LANDED_PATH)
)
このように書くことで、Auto Loader のオプションが正しく認識され、指定したフィルタが効くようになります。
原因2:modifiedAfter の日時書式とタイムゾーン
cloudFiles.modifiedAfter は「この時刻より後に更新されたファイルだけを読む」という条件を指定するオプションです。ただし、書式を間違えると 条件が無視されて全件読み込み になってしまうことがあります。
modifiedAfter の必須フォーマット
Databricks の Auto Loader では、modifiedAfter は基本的に次の形式で指定します。
- ISO-8601 形式
- UTC タイムゾーン
- ミリ秒まで含めるのが無難
具体例:
.option("cloudFiles.modifiedAfter", "2025-09-09T00:00:00.000Z")
ここで末尾の Z は「UTC」を意味します。日本時間(JST)で 2025/09/09 9:00:00 を指定したい場合は、UTC に変換して 2025/09/09 0:00:00Z として指定するイメージです。
よくあるミス例
"2025/09/09 00:00:00"のようにスラッシュ区切り+ローカル時刻で書いてしまう"2025-09-09 00:00:00"のようにTやZを省略してしまう"2025-09-09T00:00:00"と UTC やオフセットを明示しない(Runtime によっては解釈が変わる可能性)
こうした書き方でもエラーにならずに実行できてしまうことがあり、その場合「modifiedAfter が効いていないのに気づきにくい」という問題が発生します。
Databricks 上で UTC の文字列を生成する例
手入力が不安な場合、Python などで UTC の文字列を生成してからオプションに渡す方法も有効です。
from datetime import datetime, timezone
dt = datetime(2025, 9, 9, 0, 0, 0, tzinfo=timezone.utc)
modified_after_str = dt.strftime("%Y-%m-%dT%H:%M:%S.%f")[:-3] + "Z"
print(modified_after_str)
# 2025-09-09T00:00:00.000Z
こうして得た文字列を option("cloudFiles.modifiedAfter", modified_after_str) に渡せば、書式ミスの心配がかなり減ります。
原因3:通知モード(Event Grid)とリストモードの違い
Auto Loader は、ファイル検出の方法として大きく2種類をサポートしています。
- 通知モード: Event Grid + ストレージイベントで新着ファイルを検出
- リストモード: ストレージ API でディレクトリをスキャンして検出
cloudFiles.useNotifications = "true" にしていると、Event Grid 通知を前提とした振る舞いになります。このとき、環境や構成によっては「modifiedAfter で想定どおりに絞られない」「通知タイミングと改行時刻がズレる」などの要因が絡み、挙動がわかりにくくなることがあります。
まずはリストモードでシンプルに検証する
問題切り分けの第一歩としては、次のように useNotifications を明示的に "false" にし、リストモードに固定して検証するのがおすすめです。
.option("cloudFiles.useNotifications", "false")
リストモードは構成要素が少ないため、
includeExistingFilesが効いているかmodifiedAfterで意図した時刻より前のファイルがちゃんと除外されているか
といった点をストレージ側のタイムスタンプと突き合わせて確認しやすくなります。動作が期待どおりであることを確認できたら、必要に応じて通知モードへ移行する、という順番が安全です。
原因4:チェックポイントの残骸により設定変更が反映されない
Structured Streaming(Auto Loader を含む)は、「どのファイルをどこまで処理したか」をチェックポイントディレクトリに保存しています。checkpointLocation を同じパスにしたまま設定だけ変えて再実行すると、次のような現象が起こります。
- 過去の実行の進捗情報が優先される
includeExistingFilesやmodifiedAfterの値を変えても、ジョブ再開として扱われるため、期待どおりに効かない
設定変更時はチェックポイントを必ずリセットする
Auto Loader 周りの設定を変えたときは、次のいずれかの対処を行うようにしましょう。
checkpointLocationを別ディレクトリにする(新しいストリームとして実行)- 既存の
checkpointLocationを削除してから再実行する
特に、初回の検証で「includeExistingFiles を設定し忘れた」「modifiedAfter の書式を間違えた」というケースでは、同じチェックポイントを使い回している限り、何度設定を変えても挙動が変わらないように見えます。
そのため、
- 設定を変える → チェックポイントを削除 → 再実行
という手順をテンプレート化しておくと、安全に検証できます。
修正版サンプルコードと実行の流れ
以上を踏まえた、Auto Loader の基本的なサンプルコードを改めて整理します。
df_stream = (spark.readStream
.format("cloudFiles") # Auto Loader を一度だけ指定
.option("cloudFiles.format", "text") # ファイル形式(text/csv/json/parquet 等)
.option("cloudFiles.includeExistingFiles", "false")
.option("cloudFiles.modifiedAfter", "2025-09-09T00:00:00.000Z")
.option("cloudFiles.useNotifications", "false") # まずはリストモードで検証
.load(LANDED_PATH) # ADLS Gen2 上のパス
)
(df_stream.writeStream
.format("delta") # 出力形式(例:Delta Lake)
.trigger(availableNow=True) # 一度きりのバッチ処理に近い動作
.foreachBatch(processdata) # バッチ単位の処理ロジック
.option("checkpointLocation", CHECKPOINT_PATH)
.start())
ポイントを整理すると次のようになります。
- フォーマットは
format("cloudFiles")のみ:二重指定をしない - ファイル形式は
cloudFiles.formatで指定:text・csv・json・parquet など includeExistingFiles="false"で「ジョブ開始後の新着ファイルのみ」modifiedAfterでさらに対象を時刻で絞り込み- 設定を変えたらチェックポイントをクリアして再実行
includeExistingFiles と modifiedAfter の役割の違い
両者は似たような用途で使われますが、役割は少し異なります。
| オプション | 役割 | 主な用途 |
|---|---|---|
cloudFiles.includeExistingFiles | ストリーム開始「時点」で既に存在しているファイル(既存ファイル)を読むかどうかを制御する | 初回実行時に「過去分を読みたくない」場合に "false" にする |
cloudFiles.modifiedAfter | ファイルの更新日時が指定時刻より新しいかどうかで絞り込み | バックフィルや、ある境界日付以降のデータだけを処理したいとき |
特に重要なのは、includeExistingFiles="false" は「ストリーム開始」のタイミングを基準に判定されるのに対し、modifiedAfter は「ファイルの更新時刻」を基準にしているという点です。
- 「初回実行開始後に到着したファイルだけ」であれば
includeExistingFiles="false"だけでも実現可能 - 「2025-09-09 以降のファイルだけ」というように、カレンダー上の境界で切りたい場合は
modifiedAfterを併用するのが安全
実務でのおすすめパターン
実運用でよく使われる構成パターンを、3つに分けて紹介します。
パターンA:一度だけフルロード → 以降は新着のみ
- 最初のジョブ:
includeExistingFiles="true"modifiedAfterは指定しない- 全ての既存ファイルを Delta テーブルにロード
- 2回目以降のジョブ:
- チェックポイントを引き継いで
includeExistingFiles="false"に変更 - 以降は新着ファイルのみ増分で取り込み
- チェックポイントを引き継いで
このパターンは「まず過去分をすべて取り込む」ことが明確な場合に有効です。
パターンB:特定日付以降のみを取り込み(バックフィル兼用)
includeExistingFiles="true"modifiedAfterにビジネス上の境界日(例:システム切り替え日)を指定
この構成だと、初回実行で modifiedAfter 以降のファイルのみがまとめて取り込まれ、その後は新着分のみが追加されます。「それ以前のファイルは古すぎて使わない」というケースに向きます。
パターンC:検証環境での動作確認用
- 小さなテスト用ディレクトリを用意
- 毎回チェックポイントをクリアしながら、
includeExistingFilesとmodifiedAfterの組み合わせを検証 - 本番と同じ設定をテストフォルダーで試し、期待どおりかどうかを確認してから本番パスへ切り替える
特に Auto Loader の挙動は「時系列」と「チェックポイント」の影響を強く受けるため、検証環境で手順を確立しておくと本番トラブルを大きく減らせます。
ファイルの実際の更新時刻を必ず確認する
設定を何度変えても思ったとおりに絞り込めない場合、そもそも「ストレージ側の更新時刻」が想定とズレている可能性があります。ADLS Gen2 上のファイルのタイムスタンプを、Databricks 上から確認してみましょう。
display(dbutils.fs.ls(LANDED_PATH))
この出力には modificationTime(更新時刻)が含まれます。ここで見えている時間は基本的に UTC であるため、
- 日本時間でいつのファイルか
- 指定した
modifiedAfterの時刻とどちらが新しいか
を丁寧に突き合わせることで、「本当に条件に合うファイルだけが対象になっているか」を確認できます。
Databricks Runtime バージョンも確認しておく
Auto Loader も Databricks Runtime のバージョンにより細かな挙動やサポートされるオプションが変わることがあります。特定のオプションが思ったように動かないと感じたときは、
- 使用している Databricks Runtime のバージョン
- そのバージョンでサポートされている Auto Loader 機能
を確認しておくと安心です。特に古い Runtime を長期間使い続けている環境では、
- ドキュメントやサンプルコードは新しい Runtime を前提に書かれている
- そのため一部のオプションの挙動が違うように見える
といったギャップが発生しがちです。
典型的なトラブルシューティングの流れ
最後に、「全部読まれてしまう」問題を調査するときのチェックリストをまとめます。
| ステップ | 確認内容 |
|---|---|
| 1 | format("cloudFiles") が二重に指定されていないか(format("text") 等で上書きしていないか) |
| 2 | cloudFiles.modifiedAfter が UTC の ISO-8601 形式か(末尾に Z がついているか) |
| 3 | cloudFiles.useNotifications を一旦 "false" にして、シンプルなリストモードで動作を確認する |
| 4 | checkpointLocation を新規パスに変えるか、既存チェックポイントを削除してから再実行する |
| 5 | dbutils.fs.ls() でファイルの更新時刻を確認し、指定した modifiedAfter と本当に整合しているかをチェックする |
この順番で確認すれば、多くのケースで原因を早期に特定できます。
まとめ:既存ファイルをスキップしつつ新着だけ安全に取り込むには
本記事で解説したポイントを改めて整理します。
- Auto Loader を使うときは
format("cloudFiles")を一度だけ記述 し、ファイル形式はoption("cloudFiles.format", "text")のようにオプションで指定する cloudFiles.includeExistingFiles="false"だけでも、「最初の処理開始後に到着したファイルのみ」を対象にできる- 取り込み対象の時刻を厳密に絞りたい場合は
cloudFiles.modifiedAfterを併用 し、UTC の ISO-8601 形式(例:2025-09-09T00:00:00.000Z)で記述する - Event Grid などの通知モードを使う前に、まずは リストモード(
useNotifications="false")で動作確認 するとトラブルシューティングが容易 - 設定を変更したら チェックポイントをリセット しないと、過去の実行情報が優先されて「設定が効いていない」ように見える
- 最終的には、ストレージ上の 実際のファイル更新時刻 と 指定した条件 の整合性を丹念に確認することが重要
これらのポイントを押さえておけば、Auto Loader で「なぜか全部読まれてしまう」「既存ファイルがスキップされない」といった悩みから解放され、ADLS Gen2 上の新着ファイルだけを安全かつ効率的に取り込めるようになります。自動化されたデータパイプラインを安定運用するうえで、ぜひ押さえておきたいベストプラクティスです。

コメント