外部SaaSやAPIからHTTP GETだけでデータを取得し、Azure上のData Lakeに蓄積しつつ、そのままSQLでクエリできる基盤を作りたい──そんな要件は、今や多くの企業で「標準パターン」になりつつあります。本記事では、Azure Data Factory・Azure Data Lake Storage・Synapse Serverless SQLを組み合わせて、堅牢かつコスト最適な「Azureネイティブ データパイプライン」の設計と実装ポイントを、アーキテクチャから細部のパラメータまで一気に解説します。
Azureネイティブなデータパイプラインの全体像
まず、本記事で目指すアーキテクチャを整理します。要件を満たすためのキーワードは次の通りです。
- データ取得は 外部HTTP API(GETのみ)
- 処理は Azureネイティブサービスのみ で完結
- 生データは Data Lake(ADLS Gen2)へそのまま保存
- 分析用には 圧縮バイナリ形式(主にParquet) に変換
- データサイエンティストは SQLで直接クエリ(サーバーレスSQL)
- エラーハンドリング・再試行・監視・スケール・コストを考慮した設計
これを満たす「最小かつ実践的な」アーキテクチャは下図のようになります。
[ 外部 HTTP API (GET) ]
|
| Azure Data Factory (HTTP/REST コネクタ)
v
[ ADLS Gen2 /raw/... ] <--- 失敗時 /deadletter/...
|
| ADF Mapping Data Flow(型変換・正規化・圧縮)
v
[ ADLS Gen2 /curated/... (Parquet または Delta) ]
|
| Synapse Serverless SQL (外部テーブル / OPENROWSET)
v
[ データサイエンティスト / BI / Notebook からのSQLクエリ ]
この構成の良いところは、
- 常にデータは Data Lake上に静的ファイルとして存在 し、どのレイヤからも参照できる
- クエリエンジン(Synapse Serverless SQL)は コンピュートのみ で、データは動かさずにクエリできる
- サービス構成がシンプルで、段階的な拡張(Spark、Delta、メダリオンアーキテクチャ)も容易
それでは、このアーキテクチャを支える各サービスと設計ポイントを詳しく見ていきます。
採用するAzureサービスと役割
本パイプラインで利用する主なAzureサービスと役割を一覧にすると、次のようになります。
| サービス | 役割 | 主な採用理由 |
|---|---|---|
| Azure Data Factory | HTTP取り込み・変換・オーケストレーション | HTTP/RESTコネクタ、Copy activity、Mapping Data Flow、トリガー・監視などが一体で提供される |
| Azure Data Lake Storage Gen2 | 生データ/整形データの保管 | 階層型名前空間(HNS)、ACL、ライフサイクル管理により、大量データの長期保存とパーティション管理が容易 |
| Azure Synapse Analytics (Serverless SQL Pool) | Data Lake上のファイルに対するSQLクエリ | データを移動せずにParquetなどを直接クエリ可能。サーバーレスでスケールし、運用負荷が低い |
| Azure Key Vault | APIキー・シークレット管理 | 外部HTTP APIの認証情報を安全かつ一元的に管理し、ローテーションも容易 |
| Azure Monitor / Log Analytics | 監視・ログ収集・アラート | ADFやストレージ、Synapseなどの診断ログを集約し、SLA監視とトラブルシュートを行える |
この構成は「HTTP → Data Lake → 圧縮バイナリ → SQLクエリ」の要件に対して、過不足なく必要なコンポーネントだけを揃えたバランスの良い組み合わせです。
データフロー詳細:HTTP GETからSQLクエリまで
HTTP GETでのデータ取り込み(ADF)
外部ツールのAPIがHTTP GETのみ提供している場合でも、Azure Data FactoryのHTTP/RESTコネクタを利用すれば、次のような機能を標準で享受できます。
- ヘッダーやクエリパラメータの設定(APIキー、認証トークンなど)
- ページング(page / offset / nextLink など)のサポート
- レスポンスをそのままADLS Gen2へ書き込み
- 再試行ポリシー、タイムアウト、コンカレンシー制御
構成のポイントは以下の通りです。
- HTTP Linked Service:ベースURL・認証方式(Key Vault参照)を定義
- HTTP Dataset:エンドポイントパス・クエリパラメータのテンプレート化
- Copy Activity:ソースをHTTP、シンクをADLSに指定
- ページング:page番号やnextLinkを使ったループ構造(ForEach)を設定
これにより、APIからの生JSONやCSVレスポンスを、ほぼノーコードでData Lakeに保存できます。
生データの保存設計(Rawレイヤ)
生データは、後から何度でも再処理できるよう「そのまま」保存します。推奨されるディレクトリ構造の一例は以下の通りです。
/raw/
source=<システム名>/
entity=<テーブル名・エンドポイント名>/
ingest_date=YYYY/MM/DD/
ingest_ts=YYYYMMDDHHMMSS/
page=0001.json
page=0002.json
ポイントは、
- source:外部システム名やSaaS名(例:salesforce, zendesk)
- entity:エンドポイント単位(例:tickets, users)
- ingest_date:パイプライン実行日(再取り込みや差分取り込みの基準)
- ingest_ts:1回の取り込みバッチを識別するタイムスタンプ
この構造にすることで、
- どの日時にどのデータを取得したかが明確
- 誤った変換をしても、rawを元に再実行できる
- 特定期間のデータだけを再取り込み・再処理しやすい
Mapping Data Flowでの変換とParquet出力
生データをそのままクエリすることもできますが、実務上は以下の課題があります。
- JSONのネストが深く、SQLから扱いづらい
- 型が曖昧(文字列になりがち)で、数値や日付の比較が難しい
- テナントや期間ごとに構造が微妙に変わることがある(スキーマドリフト)
そこで、ADFのMapping Data Flowを用いて、次のような処理を行います。
- JSONのフラット化(必要なフィールドを抽出し、列として展開)
- 型変換(string → int / decimal / datetime など)
- 不要なフィールドの削除
- ビジネスキーやイベント日時の正規化
- 出力形式をParquet(Snappy圧縮)に設定
- パーティションキー(例:event_date, dt)でディレクトリ分割
出力ディレクトリの例:
/curated/
source=<システム名>/
entity=<テーブル名>/
dt=YYYY/MM/DD/
part-00000.snappy.parquet
part-00001.snappy.parquet
Mapping Data Flowでは、スキーマドリフト機能を有効にしておくことで、APIのレスポンスに新しいフィールドが追加された場合でも、致命的なエラーを避けつつ、必要に応じて後から列を取り込むことができます。
SQLクエリ提供(Synapse Serverless SQL)
整形済みParquetファイルをSQLから利用するには、Synapse Serverless SQL PoolからADLS Gen2を参照する設定を行います。
- Synapse WorkspaceとADLS Gen2を紐付け(Managed Identityに対してストレージACLを付与)
- Serverless SQLで外部データソースを作成
- Parquet用の外部ファイルフォーマットを作成
- 外部テーブルを作成し、クエリしやすいビューを定義
SQLのイメージは次のようになります。
-- 外部データソース
CREATE EXTERNAL DATA SOURCE ds_datalake
WITH (
LOCATION = 'https://<account>.dfs.core.windows.net',
CREDENTIAL = [Managed Identity]
);
-- ファイルフォーマット(Parquet)
CREATE EXTERNAL FILE FORMAT ff_parquet
WITH (
FORMAT_TYPE = PARQUET
);
-- 外部テーブル
CREATE EXTERNAL TABLE ext_tickets
(
ticket_id BIGINT,
status NVARCHAR(50),
created_at DATETIME2,
updated_at DATETIME2,
priority NVARCHAR(50),
dt DATE
)
WITH (
LOCATION = '/curated/source=zendesk/entity=tickets/',
DATA_SOURCE = ds_datalake,
FILE_FORMAT = ff_parquet
);
-- データサイエンティスト向けビュー
CREATE VIEW vw_tickets AS
SELECT *
FROM ext_tickets
WHERE dt >= DATEADD(DAY, -30, CAST(GETDATE() AS DATE));
これにより、データサイエンティストは SELECT * FROM vw_tickets WHERE status = 'open'; のように、通常のテーブルと同じ感覚でData Lake上のデータをクエリできるようになります。
圧縮バイナリフォーマットの選定:ParquetとDelta
「圧縮バイナリ形式」としては、主に次の2つを押さえておけば十分です。
| フォーマット | 特長 | 向いているケース |
|---|---|---|
| Parquet(Snappy) | 列指向・圧縮率が高い・Serverless SQLでネイティブ対応 | 追記中心、バッチ処理、クエリ・イン・プレースでの分析 |
| Delta Lake | ParquetにACID・スキーマ進化・タイムトラベルを付加 | アップサートが多い、レイクハウスとしてデータ品質管理を厳密に行いたい |
本記事の要件(HTTP GETの結果を貯めて分析する)であれば、まずは Parquet一択 と考えて問題ありません。将来、
- 特定のキーに対するUPSERT/MERGEが頻繁に発生する
- レコード単位の更新・削除が必要
- データ品質管理(ゴールドレイヤ)を厳密に行いたい
といった要件が出てきた段階で、Delta Lakeへの切り替え・併用を検討するとよいでしょう。
HTTP GETの再試行・冪等性・エラーハンドリング設計
外部HTTP APIを叩くパイプラインでは、ネットワークやAPI側の問題により、どうしても失敗が発生します。障害に強いパイプラインにするための具体的な設計ポイントを解説します。
再試行ポリシーの設計
Azure Data Factoryのアクティビティには、標準で Retry設定 が用意されています。代表的な設定例は次の通りです。
- Retry:3〜5回
- Retry interval:30〜120秒(徐々に伸ばすエクスポネンシャルバックオフを意識)
- Timeout:APIのSLAに合わせて(例:60〜120秒)
さらに、HTTPステータスコードによって挙動を変えると安定性が高まります。
- 429 / 5xx:再試行対象(Rate Limit・一時的な障害)
- 4xx(400, 401, 403 など):設定ミスや認証ミスの可能性が高いので早めにFailさせる
冪等性を担保するディレクトリ/ファイル命名
パイプラインが再実行されたときに、同じデータを二重に蓄積しないことが重要です。そのためには、
- 「いつ・何ページ目を取ったか」に基づく決定的なファイル名
- 取り込み済みかどうかを記録するインジェスションログ
を組み合わせます。ディレクトリ構造の例:
/raw/source=foo/entity=orders/ingest_date=2025/12/07/
page=0001.json
page=0002.json
このとき、「2025-12-07のpage=0001」が既に存在するかどうかをチェックし、存在する場合はスキップするようなロジックをパイプライン側に組み込みます。
インジェスションログ(メタデータ)の実装
重複処理の防止やトラブルシュートのため、各取り込みバッチの情報を簡単なメタデータとして保存しておくと便利です。実装手段は以下のいずれかです。
- ADLS上の小さなCSV/Parquetファイル(後でServerless SQLから参照可能)
- Azure Table Storage(シンプルなキー・バリュー構造)
インジェスションログに持たせる代表的なカラムは次の通りです。
| 項目名 | 説明 | 例 |
|---|---|---|
| pipeline_name | パイプライン名 | pl_ingest_foo_orders |
| source | 外部システム名 | foo |
| entity | エンドポイント名 | orders |
| ingest_date | 取り込み日 | 2025-12-07 |
| page | ページ番号 | 1 |
| status | 処理結果 | success / failed |
| file_path | 実際に書き込んだパス | /raw/source=foo/…/page=0001.json |
| etag | API側のETagなど(差分取得用) | “a1b2c3…” |
ADFのLookup+If Conditionアクティビティを用いて、「同じsource/entity/ingest_date/pageのレコードがsuccessで存在する場合は処理をスキップ」といった制御を行うことができます。
ETagやIf-Modified-Sinceによる差分取得
APIがETagやIf-Modified-Sinceヘッダーに対応している場合、差分取得も実現できます。
- 前回取得時のETagをインジェスションログに保存
- 次回のHTTP GETで
If-None-Match: <前回ETag>を指定 - レスポンスが304(Not Modified)の場合は、データが変わっていないと判断してスキップ
これにより、無駄なデータ転送や処理を減らし、コスト削減と安定性向上が期待できます。
エラー時のデッドレター戦略
完全に失敗したリクエストについても、「何が起きたか」を残しておかないと調査が困難になります。そこで、デッドレターパターンを採用します。
- HTTPステータスが4xx/5xxのレスポンスボディを
/deadletter/...以下に保存 - 保存するファイルには、エラーコード・URL・ヘッダー・ボディ・タイムスタンプを含める
- Log Analyticsにもエラーイベントを送信し、アラートルール(メール、Teams通知など)を設定
これにより、後から手動でデッドレターを再処理したり、API提供元に問い合わせる際の証拠として活用できます。
スケーラビリティとコスト最適化のベストプラクティス
取り込みパイプラインの並列化
HTTP APIのスループットを最大限活かしつつ、API提供元に迷惑をかけないようにするには、適切な並列度の設計が重要です。
- 日付やIDレンジ単位での分割(例:1日ごと、ユーザーIDの範囲ごと)
- ページングをForEachで並列に処理しつつ、最大コンカレンシーを制限
- API側のRate Limitに合わせてスロットリング(Waitアクティビティなど)を挿入
ADFのパイプラインレベルで Max concurrent runs を設定し、環境全体での同時実行数を制御することも有効です。
ファイルサイズと「small files問題」への対応
Serverless SQLは「スキャンしたデータ量」によって課金されるため、Parquetを使うこと自体がコスト削減に大きく貢献します。しかし、ファイルが小さすぎると、
- ファイル数が増えすぎてメタデータ処理がボトルネックになる
- クエリ時のオープン・クローズ回数が増え、レイテンシが悪化する
一般的な目安として、1ファイルあたり 100〜500MB程度 に収まるよう調整するとよいとされています。Mapping Data Flowや他の変換処理では、
- パーティション数を制御する(パーティションキーの選び方を工夫)
- 必要に応じて「Coalesce」的な処理でファイルをまとめる
といった工夫を行います。
パーティション設計とディレクトリ構造
パーティション設計は、クエリパターンとスキャンコストを大きく左右します。よく使われるのが以下の2軸です。
- event_date:ビジネスイベントが発生した日(例:注文日、チケット作成日)
- ingest_date:データを取り込んだ日
クエリで「過去30日」などの絞り込みを行うことが多い場合、dt=YYYY/MM/DD のように日付でディレクトリを切っておくと、Serverless SQLで WHERE dt >= '2025-11-01' のように条件を付けるだけで、該当パーティションだけをスキャンすることができます。
Serverless SQLコスト最適化テクニック
Serverless SQLの課金は「スキャンしたデータ量(バイト数)」に比例するため、次のようなクエリ設計が重要です。
- Parquetを利用する(CSVなどのテキスト形式は極力避ける)
- SELECT * を多用せず、必要な列だけを指定する
- WHERE句で列のフィルタに加えて、パーティションキー(dtなど)へのフィルタを必ず含める
- よく使うクエリはビューに切り出し、適切なフィルタをデフォルトにしておく
例えば、次の2つのクエリでは、スキャン対象のデータ量が大きく異なります。
-- NG: 全パーティションをスキャン
SELECT *
FROM ext_tickets
WHERE status = 'open';
-- OK: dtで期間を絞った上で必要な列だけ取得
SELECT ticket_id, status, created_at
FROM ext_tickets
WHERE dt >= DATEADD(DAY, -7, CAST(GETDATE() AS DATE))
AND status = 'open';
このようなガイドラインをチームで共有しておくと、予測しやすいコストで運用しやすくなります。
ストレージ階層とライフサイクル管理
大量データを長期間保存する前提であれば、ストレージ階層とライフサイクル管理も必須です。
- /raw:一定期間(例:90日)はHot、以降はCool、さらに古いものはArchiveへ自動移行
- /curated:直近6〜12ヶ月はHot、それ以前はCoolに移行
Azure Storageのライフサイクル管理ポリシーを設定しておけば、条件に応じて自動的に階層が変更され、コストを継続的に最適化できます。
SQL提供レイヤの実装詳細
外部データソースと外部テーブルの定義
前述したように、Synapse Serverless SQLからADLSを参照するには、外部データソースと外部テーブルを定義します。主要なステップは次の通りです。
- Synapse WorkspaceのManaged Identityに対して、ADLSコンテナへの読み取り権限(Storage Blob Data Readerなど)を付与
- Serverless SQLで外部データソースを作成
- 外部ファイルフォーマットを作成
- 外部テーブルを作成し、論理スキーマを定義
サンプルSQL:
-- 1. 外部データソース
CREATE EXTERNAL DATA SOURCE ds_datalake
WITH (
LOCATION = 'https://<account>.dfs.core.windows.net/<container>',
CREDENTIAL = [Managed Identity]
);
-- 2. 外部ファイルフォーマット
CREATE EXTERNAL FILE FORMAT ff_parquet
WITH (
FORMAT_TYPE = PARQUET
);
-- 3. 外部テーブル作成
CREATE EXTERNAL TABLE ext_orders
(
order_id BIGINT,
customer_id BIGINT,
amount DECIMAL(18,2),
currency NVARCHAR(10),
order_date DATE,
dt DATE
)
WITH (
LOCATION = '/curated/source=foo/entity=orders/',
DATA_SOURCE = ds_datalake,
FILE_FORMAT = ff_parquet
);
ビューでの抽象化と権限管理
外部テーブルは「物理レイヤ」に近いため、データサイエンティストやBIツールには、ビューを通して提供するのがおすすめです。
- 列名のエイリアスを付与し、ビジネス用語に揃える
- dtの既定値(例:過去1年)を仕込んでおく
- 非公開列(個人情報など)はビューから除外し、ビューに対してのみ権限を付与
CREATE VIEW vw_orders AS
SELECT
order_id AS id,
customer_id AS customer_id,
amount,
currency,
order_date
FROM ext_orders
WHERE dt >= DATEADD(YEAR, -1, CAST(GETDATE() AS DATE));
ユーザーやグループには、ベースとなる外部テーブルではなくビューへのSELECT権限を付与することで、権限管理とスキーマ進化を柔軟に行えます。
監視・可観測性・運用設計
技術的な構成がどれだけ良くても、「動いているかどうか」「どのくらい遅れているか」「どこで失敗しているか」が見えなければ運用は破綻します。ここでは、監視・可観測性の設計ポイントを整理します。
ADFのモニタリングとLog Analytics連携
- ADFの「Monitor」画面で、各パイプラインの成功/失敗を一覧確認
- 診断ログをLog Analyticsワークスペースに送信し、実行ログをKustoクエリで分析
- 失敗時・遅延時のアラートルールをAzure Monitorで定義(メールやTeams通知)
例えば、過去15分で失敗したパイプラインが1件でもあればアラートを飛ばす、といった設定が可能です。
メトリクスの例
- パイプラインごとの成功/失敗回数
- 平均・最大実行時間
- インジェストしたレコード数・データ量(Mapping Data Flow側でカウント)
- Serverless SQLのスキャンデータ量(クエリログから集計)
これらをダッシュボード化しておくことで、異常の早期検知や容量計画にも役立ちます。
発展パターン:Sparkやより高度なレイクハウスへの拡張
本記事では「あくまで最小構成」に絞って解説しましたが、将来的には次のような拡張も検討できます。
- 複雑な結合や機械学習前処理が増えてきたら、Synapse SparkやDatabricksによる変換レイヤを追加
- Bronze / Silver / Gold などのメダリオンアーキテクチャに発展させ、Delta Lakeでゴールドレイヤの信頼性を高める
- イベントドリブンなNear-Real-Time処理が必要になったら、Event HubsやStream Analyticsを組み合わせて拡張
重要なのは、最初から完璧なレイクハウスを目指すのではなく、HTTP → Data Lake → Parquet → Serverless SQLという「シンプルで拡張可能なコア」を作ることです。そこから、ビジネス要件に応じて少しずつ機能を足していく方が、運用負荷とリスクを抑えつつ価値を積み上げていけます。
実装チェックリスト(詳細版)
最後に、本記事で紹介した内容をもとに、設計・実装時に確認すべきチェックポイントを整理します。
- HTTP取り込み
- ADF HTTP/RESTコネクタで、ヘッダー・クエリパラメータ・認証を正しく設定したか
- Retry回数・間隔・TimeoutをAPIのSLAに合わせて設定したか
- 429/5xxと4xxで挙動を変える(再試行/即Fail)設計になっているか
- ページングの方式(page / offset / nextLink)をAPI仕様に合わせて実装したか
- Data Lake構造
- /rawと/curatedのディレクトリ構造・命名規則が整理されているか
- source, entity, ingest_date, dtなどのキーが一貫しているか
- 1ファイルあたりのサイズが概ね100〜500MBに収まるよう設計されているか
- 冪等性・差分取得
- 決定的なファイル名・パスを採用しているか
- インジェスションログをどこに格納し、どのように参照するか決めているか
- APIがETag / If-Modified-Sinceに対応している場合、それを活用する実装になっているか
- 変換・Parquet出力
- Mapping Data Flowでスキーマドリフトを適切に許容しているか
- 型変換(数値・日付)はすべて明示的に行っているか
- パーティションキー(dtなど)がクエリパターンに合っているか
- SQL提供レイヤ
- SynapseのManaged IdentityにADLSのアクセス権限が付与されているか
- 外部データソース・外部テーブル・ビューが正しく作成されているか
- データサイエンティストやBIツールにはビューを介してアクセスさせているか
- クエリガイドライン(SELECT *禁止など)がチームで共有されているか
- 監視・運用
- ADFの診断ログがLog Analyticsに送信されているか
- 失敗・遅延・SLA逸脱のアラートルールが設定されているか
- ライフサイクル管理ポリシーで、古い/raw・/curatedファイルの階層が自動で切り替わるか
- Key VaultでAPIキーを管理し、ローテーション手順が定義されているか
まとめ:AzureネイティブでHTTPデータ基盤を構築する
本記事では、外部HTTP GET APIからのデータをAzure上だけで取り込み、Data Lakeに生データを保存しつつ、圧縮バイナリ形式でSQLクエリ可能にするアーキテクチャを解説しました。
- 取り込み・変換・オーケストレーションには Azure Data Factory
- 生データ/整形データの永続化には Azure Data Lake Storage Gen2
- クエリ・イン・プレースには Synapse Serverless SQL Pool
- 機密情報管理には Azure Key Vault
- 監視・可観測性には Azure Monitor / Log Analytics
という、比較的シンプルな構成です。
フォーマットはまず Parquet を基本とし、将来のアップサートや高品質なレイクハウス要件が出てきたときに Delta Lake を検討する、という2段階アプローチが現実的です。
最後にもう一度強調すると、重要なのは「Azureネイティブでシンプルなコアパターン」を作ることです。HTTP → Data Lake → Parquet → Serverless SQLという流れさえきちんと押さえておけば、サービスの追加やデータ量の増加、分析ニーズの高度化にも柔軟に対応できます。ぜひ本記事の内容を、自社のデータパイプライン設計・見直しに活かしてみてください。

コメント