Azure Databricks×ADLS Gen2で実現するイベント駆動ETL設計とネットワーク構成

Adobe やマーケティング基盤から JSON ファイルが ADLS Gen2 に到着したタイミングで、Azure Databricks の ETL を自動起動したい――多くのデータレイクプロジェクトで出てくるテーマです。本記事では、イベント駆動とポーリングの設計、ネットワーク(VNet インジェクト/NPIP)やセキュリティまで含めて、実践的なアーキテクチャと実装パターンを整理します。

目次

シナリオ概要:Adobe → ADLS Gen2 → Azure Databricks

本記事で想定するシナリオは次の通りです。

  • Adobe 側で Cron ジョブが動き、定期的に JSON ファイルを出力する
  • 出力先は Azure Data Lake Storage Gen2(ADLS Gen2)のコンテナー(Bronze レイヤー)
  • ファイル到着を検知して Azure Databricks(Unity Catalog なし)のジョブを自動起動したい
  • Databricks ワークスペースは VNet インジェクトかつ NPIP(Public IP なし)で閉域化

ポイントは大きく次の 3 つです。

  1. ADLS Gen2 にファイルが「いつ到着したか」をどう検知するか
  2. 閉域ネットワーク(VNet インジェクト+Private Link)でどうやって安全に連携するか
  3. イベント駆動とポーリング、それぞれのコストと運用のトレードオフをどう見るか

以下では、トリガー設計からネットワーク、実装サンプルまで一気に整理していきます。

ADLS Gen2 のファイル到着を検知する基本パターン

まず押さえておきたい前提は、ADLS Gen2 単体では「ファイル到着を Databricks に直接通知する機能はない」という点です。したがって、次のどちらかの方式で「新規ファイル到着」を検知し、Databricks ジョブを起動する必要があります。

  • イベント駆動(Event Grid ベース) … 推奨
  • ポーリング(Databricks のみで完結)

イベント駆動アーキテクチャ(推奨)

推奨パターンは、Blob 作成イベントを Azure Event Grid で受け取り、サーバーレスコンポーネントから Databricks ジョブを起動する方式です。

典型的なフローは以下のようになります。

  1. Adobe から ADLS Gen2(Bronze コンテナー)の特定パスへ JSON ファイルを書き込み
  2. ADLS Gen2 が BlobCreated イベントを Azure Event Grid に発行
  3. Event Grid のサブスクリプション先として Azure Functions または Logic Apps を設定
  4. Functions / Logic Apps が Event データ(パスなど)を受け取り、Databricks Workflows(ジョブ)REST API を呼び出して「Run now」
  5. Databricks ジョブが起動し、abfss 経由で Bronze の JSON ファイルを読み取り、Silver/Gold に加工

構成要素と役割は次のように整理できます。

コンポーネント役割ポイント
ADLS Gen2ファイル格納(Bronze)Blob 作成イベントを Event Grid に発行
Azure Event GridイベントルーターBlobCreated イベントを Functions / Logic Apps / Queue へ配信
Azure Functions / Logic Apps軽量オーケストレーションEvent を Databricks 形式に変換し、REST API を実行
Azure Databricks WorkflowsETL 実行基盤外部からの「Run now」トリガーでジョブを起動

この方式の利点は多数あります。

  • 疎結合:ADLS と Databricks が「イベント」でつながるため、構成変更の影響が小さい
  • 低コスト:ポーリングしないので、新規ファイルが来ない時間帯はほぼゼロコスト
  • 低遅延:Blob 作成直後にイベントが発生し、ほぼリアルタイムで ETL を起動可能
  • 高スケール:多数のファイル到着にも Event Grid & Functions が自動スケールで追随

ポーリング(Databricks のみで完結)

もう一つのやり方は「Databricks だけで完結させる」方式です。構成はシンプルで、次の 2 パターンがあります。

  • Databricks Auto Loader(cloudFiles) で対象ディレクトリを監視し、新着のみ段階的に取り込む
  • dbutils.fs.ls や Spark APIs で定期的にディレクトリ一覧を取得し、「既処理ファイル一覧(メタデータ表)」と突き合わせて新規のみ処理する

イメージとしては「Databricks Workflows ジョブを 5~10 分間隔でスケジューリングし、その中で『新着ファイルを探して処理』する」形です。

この方式のデメリットは、次のようなポイントです。

  • 新規ファイルが無い場合でも、クラスター起動/ファイル一覧取得のコストが発生する
  • 検知遅延はジョブ実行間隔に依存(例:10 分おきに実行なら、最大 10 分の遅延)
  • ファイル数が増えてくると、一覧取得そのものが無視できないコストになる

一方で、メリットとしては以下があります。

  • Event Grid や Functions を新たに構築・運用する必要がない
  • PoC や小規模環境では「とにかく Databricks だけで完結させたい」ときに手早く組める

イベント駆動 vs ポーリングの比較

観点イベント駆動(推奨)ポーリング(Databricks のみ)
アーキテクチャEvent Grid+Functions などで疎結合Databricks ワークフローのみで完結
検知方式BlobCreated イベントをトリガー一定間隔でディレクトリを一覧
コストアイドル時ほぼゼロ、実際のイベント数に比例新規なしでもクラスター起動&一覧コスト
遅延秒〜数十秒オーダージョブ間隔依存(分オーダー)
スケーラビリティEvent Grid / Functions が自動スケールファイル数増加で一覧コストが逓増
設計・運用の複雑さコンポーネントが増えるが、ロジックはシンプル構成は簡単だが、冪等性や再実行設計は自前
向いているケース本番運用、即時性が必要、到着頻度が低〜中PoC / 小規模、即時性があまり要らない場合

Databricks Auto Loader を使った新規ファイル検知

ポーリング方式であっても、ADLS Gen2 との連携に Databricks Auto Loader(cloudFiles) を使うと、手組みよりはるかに安全かつ高機能になります。

Auto Loader のポイントは以下です。

  • 対象ディレクトリに到着するファイルをストリーミング的に常時監視できる
  • 内部で「チェックポイント」を管理し、既処理ファイルは自動スキップしてくれる
  • Azure では 「ディレクトリ一覧モード」と「イベント通知モード(Event Grid 連携)」を選択可能

最初は構築が簡単な「ディレクトリ一覧モード」で開始し、ファイル数や即時性の要件がシビアになってきたら「イベント通知モード」に切り替える、というステップアップも現実的です。

Auto Loader の基本サンプル(ディレクトリ一覧モード)

Bronze に到着する JSON を Delta(Silver 相当)に書き出す最小サンプルは次のようになります。


df = (
    spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "json")
        .load("abfss://<container>@<account>.dfs.core.windows.net/bronze/incoming/")
)

(
    df.writeStream
        .format("delta")
        .option(
            "checkpointLocation",
            "abfss://<container>@<account>.dfs.core.windows.net/_chk/bronze_incoming"
        )
        .outputMode("append")
        .table("silver.adobe_json")
)

この構成の特徴は次の通りです。

  • チェックポイント(checkpointLocation) に処理済みオフセットが保存され、ジョブ再起動時も冪等性が保たれる
  • Silver テーブル(Delta)に到着順に追記されるため、Downstream のバッチ/BI からは通常のテーブルとして参照可能
  • 将来的にイベント通知モードへ切り替える場合も、Auto Loader の設定追加のみで済みやすい

Auto Loader をイベント通知モードに拡張する考え方

大量ファイルや即時性(数秒〜十数秒単位)が重要になってきたら、Auto Loader を Event Grid ベースのイベント通知モードにすることで、一覧コストを大幅に下げられます。

イメージとしては次のような構成です。

  • ADLS Gen2 → Event Grid → Auto Loader(Notification モード)
  • Auto Loader が Event Grid サブスクリプションを利用し、「どのファイルが新規か」をイベント経由で受け取る
  • ディレクトリ一覧よりもスケーラブルかつ高速に新規ファイルを検知

イベント駆動でありながら、Databricks 側のコードはディレクトリ一覧モードとほぼ同じで済む点も実運用では大きなメリットです。

Databricks Workflows を外部シグナルで起動する方法

イベント駆動アーキテクチャを採用する場合、Databricks Workflows(ジョブ)を外部の Event Handler から起動できる設計が鍵になります。代表的なパターンは以下です。

  • Azure Functions / Logic Apps から Databricks REST API(/api/2.1/jobs/run-now) を呼び出す
  • Event Grid → Storage Queue → Databricks ジョブ(キューを監視)という間接トリガー

REST API を使ったシンプルな起動例

Azure Functions(Python)や Logic Apps から呼び出す際のイメージとして、curl 風の例を示します。


curl -X POST \
  -H "Authorization: Bearer <DATABRICKS_PERSONAL_ACCESS_TOKEN>" \
  -H "Content-Type: application/json" \
  https://<databricks-instance>/api/2.1/jobs/run-now \
  -d '{
        "job_id": 123,
        "notebook_params": {
          "input_path": "abfss://<container>@<account>.dfs.core.windows.net/bronze/incoming/<file>.json"
        }
      }'

実際には Personal Access Token ではなく、サービス プリンシパル+OAuth や Azure Managed Identity を活用するパターンが推奨されます。重要なのは、「外部 Event Handler から Databricks ジョブを起動する標準的な入り口は REST API」であるという点です。

VNet インジェクト&NPIP 環境でのネットワーク設計

今回の前提は、Databricks ワークスペースが VNet インジェクト&NPIP(Public IP なし)であることです。この条件で Adobe → ADLS Gen2 → Databricks のパスを成立させるには、次の設計が重要です。

ADLS Gen2 と Private Link(プライベート エンドポイント)

まず、ストレージ側では Private Link(プライベート エンドポイント) を用意し、VNet から閉域でアクセスできるようにします。

  • dfs.core.windows.net 用のプライベート エンドポイント
  • blob.core.windows.net 用のプライベート エンドポイント

ADLS Gen2 は DFS(Hadoop 互換 API)と Blob の両方を利用するため、双方に PE を作成しておくと安心です。

また、プライベート DNS ゾーンを作成し、VNet とリンクさせます。

  • privatelink.dfs.core.windows.net
  • privatelink.blob.core.windows.net

ストレージ アカウント側のネットワーク設定は、原則として次のようにします。

  • ネットワーク規則:「選択したネットワークのみ許可」
  • パブリック ネットワーク アクセス:無効

Databricks 側のネットワーク設計

VNet インジェクトされた Databricks クラスターは、同じ VNet(またはピアリングされた VNet)上の プライベート エンドポイントを経由して ADLS Gen2 にアクセスします。

対外通信(コントロールプレーンや必要な外部 API 呼び出し)については、次のような設計が一般的です。

  • VNet からインターネットへの出口を NAT Gateway に集約
  • リソース自体には Public IP を付与せず、アウトバウンドのみ NAT 経由に統一

ADLS への認証は、Microsoft Entra ID(旧 Azure AD)ベースの RBAC を第一候補にします。

  • Databricks クラスターに紐づく サービス プリンシパル を作成
  • 対象コンテナーに Storage Blob Data Reader / Contributor などの RBAC ロールを付与
  • クレデンシャル情報は Azure Key Vault に保存し、Databricks の Key Vault バックド シークレットスコープから参照
要素推奨設定補足
ストレージ公開方法Private Endpoint+Private DNSPublic アクセスは無効、選択したネットワークのみ許可
Databricks クラスターVNet インジェクト+NPIPアウトバウンドは NAT Gateway へ集約
認証方式Entra ID RBAC+サービス プリンシパル共有キー/SAS の利用は最小化
シークレット管理Azure Key Vault+Databricks シークレットスコープアプリケーションコードに秘密情報を書かない

イベント基盤(Event Grid・Functions・Logic Apps)の閉域化

Event Grid や Functions を利用する場合も、NPIP 方針に合わせて閉域化することが重要です。

  • Event Grid:Private Link エンドポイントを利用し、必要に応じてプライベート DNS を設定
  • Azure Functions / Logic Apps:VNet 統合や Private Endpoint を活用して、内部通信は VNet 内で完結

さらに、運用監視の観点からは次のような設定が有効です。

  • Event Grid の Dead Letter Queue(DLQ) を Storage や Service Bus に設定し、配信失敗時のイベントを後追い可能にする
  • ストレージ/Event Grid/Functions の診断ログを Log Analytics に集約し、アラートルールを設定

Adobe から直接 Bronze に書き込むパターン

要件によっては、Adobe 側の送信元 IP をホワイトリストし、外部システムから ADLS Gen2(Bronze)へ直接ファイル投入したいというケースもあります。ここでは代表的な 3 パターンを整理します。

パターン A:パブリック エンドポイント+IP 制限

最もシンプルなのは、ADLS Gen2 のパブリック エンドポイントを開けて、ストレージ ファイアウォールで Adobe の送信元 IP のみ許可する方式です。

  • 利点:構成が簡単で、Adobe 側から見て「普通の HTTPS/REST で書き込むだけ」で済む
  • 欠点:環境方針が「全リソース私設+Public 無し(NPIP)」の場合、ガバナンス的に例外となる

セキュリティレビューが厳格な環境では、このパターンは採用が難しい場合も多いです。

パターン B:ADLS Gen2 の SFTP 機能を利用

ADLS Gen2 が持つ SFTP 機能を利用し、Adobe から SFTP 経由でファイルをアップロードする方式です。

  • 認証:公開鍵認証を用い、最小権限のロール/ローカルユーザーで運用
  • ネットワーク:
    • Private Endpoint+プライベート DNS で閉域
    • または必要最小限の IP のみ許可して公開

SFTP を有効化する際は、既存の Private Link 設定やファイアウォールとの整合性を必ず確認してください。

パターン C:受け口プロキシ(API Management / Functions)を経由(推奨)

NPIP が強く求められる環境では、受け口となる API やサーバーレスを 1 箇所だけパブリックにし、そこから内部の ADLS Gen2 に保存する方式がバランスが良く、現場でも好まれる構成です。

  • 外部公開部分:API Management または Functions(HTTP トリガー)
    • Adobe の送信元 IP のみ許可(IP 制限)
    • 必要に応じて Basic 認証や OAuth で追加防御
  • 内部:受け口から VNet 内の Private Endpoint 経由で ADLS Gen2 に書き込み

この方式では、パブリックに露出するリソースを「受け口」のみに限定でき、データレイクや Databricks は完全にプライベートなまま運用できます。

パターンメリットデメリット向いているケース
A: Public+IP 制限構成が最もシンプルNPIP 方針と衝突しやすい小規模環境、厳格な閉域要件がない場合
B: ADLS SFTP既存の SFTP 運用フローに乗せやすいSFTP 特有のネットワーク設定、キー管理が必要SFTP での連携が標準になっている組織
C: 受け口プロキシ外部公開範囲を最小化しつつ柔軟API 実装・運用のひと手間は必要NPIP/閉域方針が強い本番環境

軽量ポーリング実装の疑似コード例

Databricks のみでポーリングを行う場合、「どのファイルを処理済みとみなすか」を管理するメタデータ表が重要です。以下は、簡略化した疑似コードのイメージです。


# 1. 既処理ファイル情報を Delta テーブルから取得
processed_df = spark.table("meta.processed_files")  # columns: path, processed_at など
processed_paths = set(row["path"] for row in processed_df.select("path").collect())

# 2. 現在ディレクトリに存在するファイルの一覧を取得
latest_files = dbutils.fs.ls("abfss://<container>@<account>.dfs.core.windows.net/bronze/incoming/")
latest_paths = [f.path for f in latest_files]

# 3. 新規ファイルだけを抽出
new_paths = [p for p in latest_paths if p not in processed_paths]

# 4. 新規ファイルを順次処理
for path in new_paths:
    df = spark.read.json(path)
    # ここで必要な変換処理(スキーマ変換、正規化、アップサートなど)を実施
    # 例: df.write.format("delta").mode("append").saveAsTable("silver.adobe_json")

# 5. 処理済みファイルを meta.processed_files に登録(冪等性の要)
new_processed = [(p,) for p in new_paths]
spark.createDataFrame(new_processed, ["path"]) \
     .withColumn("processed_at", current_timestamp()) \
     .write \
     .format("delta") \
     .mode("append") \
     .saveAsTable("meta.processed_files")

本番運用では、次のような点を追加で検討してください。

  • リトライ時に部分的な失敗をどう扱うか(トランザクション境界)
  • ファイルの再投入や上書きをどう扱うか(バージョン管理/ハッシュによる重複検知)
  • スキーマ変更時の扱い(スキーマ進化や互換性チェック)

ADLS Gen2 への認証と接続サンプル(Databricks)

Databricks から ADLS Gen2 に接続する際は、サービス プリンシパル+OAuth を用いるのがベストプラクティスです。代表的な設定イメージを示します。


configs = {
  "fs.azure.account.auth.type.<account>.dfs.core.windows.net": "OAuth",
  "fs.azure.account.oauth.provider.type.<account>.dfs.core.windows.net":
      "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider",
  "fs.azure.account.oauth2.client.id.<account>.dfs.core.windows.net":
      dbutils.secrets.get(scope="kv-scope", key="sp-client-id"),
  "fs.azure.account.oauth2.client.secret.<account>.dfs.core.windows.net":
      dbutils.secrets.get(scope="kv-scope", key="sp-client-secret"),
  "fs.azure.account.oauth2.client.endpoint.<account>.dfs.core.windows.net":
      "https://login.microsoftonline.com/<tenant-id>/oauth2/token"
}

for k, v in configs.items():
    spark.conf.set(k, v)

この設定を行ったうえで、abfss:// スキームを使ってファイルにアクセスします。


df = spark.read.json(
    "abfss://<container>@<account>.dfs.core.windows.net/bronze/incoming/sample.json"
)

シークレット(クライアント ID / シークレット)は Key Vault バックド シークレットスコープで管理し、ノートブックやジョブ設定には平文を絶対に埋め込まないようにします。

設計・運用のチェックリスト

Adobe → ADLS Gen2 → Databricks のアーキテクチャを設計する際に、抜け漏れがないか確認しやすいよう、チェック項目を整理しておきます。

カテゴリチェック項目概要
ネットワークPrivate Endpoint(DFS/Blob)ADLS Gen2 に対して PE を両方作成しているか
ネットワークプライベート DNSVNet とリンクし、名前解決が PE 側に向くように設定しているか
ネットワークPublic アクセス無効ストレージのパブリック ネットワーク アクセスを無効化しているか
イベント基盤Event Grid 閉域化必要に応じて PE+プライベート DNS を構成しているか
イベント基盤Dead Letter Queueイベント配信失敗時の DLQ を必ず設定しているか
DatabricksVNet インジェクト & NPIPワークスペースの構成方針と整合しているか
認証Entra ID RBACサービス プリンシパルに最小限の RBAC ロールを付与しているか
認証シークレット管理Key Vault+シークレットスコープで秘密情報を集中管理しているか
監視Log Analyticsストレージ/Event Grid/Functions/Databricks のログを集約しているか
監視アラート失敗率増加、遅延増大などに対するアラート閾値を定義しているか
データレイクレイヤー分離Bronze / Silver / Gold のレイヤリングを明示しているか
データレイク命名規則コンテナー/ディレクトリ/テーブルの命名ルールを定めているか
冪等性重複投入防止チェックポイントやメタデータ表で重複処理を避けているか
冪等性再実行設計ジョブ再実行で破綻しない設計になっているか

まとめ:小さく始めてイベント駆動へ

Adobe から ADLS Gen2 への JSON ファイル到着を Azure Databricks で検知・起動するには、次のポイントを押さえておくと設計がスムーズになります。

  • ADLS Gen2 は「到着通知」を直接は行わないため、イベント駆動またはポーリングでトリガーを設計する必要がある
  • 本番運用の観点では、Event Grid+Functions などのイベント駆動アーキテクチャが原則ベストプラクティスとなる
  • ネットワークは Private Link+プライベート DNS を前提に閉域化し、認証は Entra ID+RBAC を軸に設計する
  • Adobe からの直接投入は、方針が「全リソース私設+Public 無し」であれば、受け口プロキシ(API Management / Functions)や ADLS SFTP を優先する
  • 最初は Auto Loader のディレクトリ一覧モード+シンプルなワークフローで小さく始め、負荷や即時性の要求が高まったタイミングで イベント通知モードや Event Grid 連携へ拡張する

これらのパターンを組み合わせることで、「NPIP かつ VNet インジェクトされた Databricks」でも、セキュアかつ運用しやすいイベント駆動 ETL 基盤を実現できます。自社のネットワーク方針や Adobe 側の制約に合わせて、最適な組み合わせを検討してみてください。

この記事を書いた人

実務の現場で詰まりがちなポイントを地図にするITブログ「IT trip」を運営。Windows/Office(Teams・Excel)からSQL、サーバ運用、ガジェットまで、再現性のある手順と“なぜそうなるか”を丁寧に解説します。読んだらすぐ試せること、そして迷った人の次の一歩が見えることを大切にしています。

コメント

コメントする

目次