Azure Data Factory(ADF)でWeb APIを定期的に呼び出してData Lakeへ保存する運用では、「前回の成功実行時刻(=差分取得の起点)をどう管理するか」で設計が大きく変わります。本記事では、全ファイル走査を避けながら確実にウォーターマーク(水位)管理する2つの方法を、ADFの具体設定例つきで解説します。
やりたいこと:ADFパイプラインで「直近の成功実行日時」をAPIに渡して差分取得したい
典型的なシナリオは次のとおりです。
- ADFのCopyアクティビティ(HTTPコネクタ)でWeb APIを呼び出し、返ってきたJSONをAzure Data Lake Storage(またはBlob)へ保存している
- パイプラインは1日数回(例:6回)実行される
- 毎回の実行で、APIのクエリパラメータに前回の成功実行時刻を渡し、差分だけを取りたい(例:
modifiedtimestamp >= 前回成功時刻) - 保存先のフォルダ/ファイルは日時付きで増え続けるため、Get Metadataで毎回一覧を取って「最新ファイル」を探す運用は避けたい
ここで重要なのは、欲しいのが「前回の実行時刻」ではなく前回の“成功”実行時刻である点です。失敗した実行は差分の起点にしてはいけません(起点を進めてしまうと取りこぼしが発生し得ます)。
結論:実現方法は2つ。おすすめは「成功時刻を1か所に保存する」
目的は同じでも、実装パターンは大きく2つに分かれます。
| 方法 | 概要 | 実装難易度 | 権限・前提 | 強み | 注意点 |
|---|---|---|---|---|---|
| 方法① 専用ファイル(レコード)に前回成功日時を保存 | パイプライン開始で読み取り、最後に成功したら上書き更新 | 低 | ストレージへの読み書き権限 | シンプル/高速/失敗時に起点が進まない | 状態管理用のファイル(1つ)を持つ |
| 方法② ADF管理REST APIから実行履歴を取得 | queryPipelineRunsで「Succeeded」の最新を検索 | 中〜高 | Managed Identityに管理APIの権限(Reader等) | 状態ファイルが不要/履歴ベースで統一 | 履歴保持期間・権限・API呼び出し設計が必要 |
運用の安定性と作りやすさを重視するなら、まずは方法①(採用されやすい王道)が鉄板です。監査・統制の都合で「状態ファイルを持ちたくない」「複数パイプラインを履歴で一括管理したい」といった要件がある場合に、方法②が候補になります。
なぜ「保存先の最新ファイルを探す」設計は避けられがちなのか
Get Metadataでコンテナ配下を列挙し、タイムスタンプ付きファイル名から最大値を取る…という発想は一見スマートですが、運用が長期化すると次のような問題が表面化します。
- 列挙コストが増える(フォルダやファイルが増えるほど遅くなる)
- 設計が複雑化(階層が深いほど、絞り込みやソートが煩雑)
- 「成功実行」の判定が難しい(ファイルがあっても中身が不完全、途中失敗などの扱い)
- APIの起点=ウォーターマークの責務が曖昧(「保存された最新」≠「取得に成功した最新」になり得る)
そのため、差分取得の設計では「ウォーターマークは1か所で明示的に管理」し、パイプラインの成功条件と一致させるのが基本方針になります。
方法①:専用ファイルに「最後に成功した日時」を保存して管理する(おすすめ)
この方法の考え方(ウォーターマーク管理)
やることはシンプルです。
- 状態管理用に1つだけファイル(またはレコード)を用意する
- 実行開始時にそのファイルを読み、値をAPIの差分条件に使う
- パイプラインが最後まで成功したら、今回の実行時刻をそのファイルに上書きする
ここで最重要なのは、失敗した実行では更新しないことです。これにより常に「直近の成功実行時刻」だけが残り、取りこぼしリスクを最小化できます。
設計のポイント(実務で効くコツ)
| ポイント | 理由 | 具体例 |
|---|---|---|
| 時刻はUTCで統一 | タイムゾーン差・夏時間で事故りにくい | 2025-12-16T01:23:45Z のようなISO 8601 |
| 今回の実行時刻は変数に一度だけ格納 | パイプライン途中で日付をまたぐ等の不整合を防ぐ | This_Time = @utcnow() を最初に1回だけ |
| 更新は最終成功時のみ | 失敗実行でウォーターマークが進むと取りこぼしになる | 依存関係を「Succeeded」に限定 |
| 状態ファイルはデータ本体と分離 | 誤削除・権限・運用の衝突を減らす | /control/watermark/last_success.txt |
準備:必要なデータセットと変数
最小構成でも以下を用意するとスムーズです。
| 種類 | 名前例 | 用途 | ポイント |
|---|---|---|---|
| ファイル(読み書き) | LastRecord.txt | 前回成功日時を保存 | 1行1列のテキスト(CSVでも可)にするとLookupが楽 |
| ファイル(読み取り) | blank.txt | LastRecord更新用の“ダミー入力” | ADFが「1行ある」と認識できるよう、1文字だけでも入れる |
| 変数 | Last_Time | 前回成功実行日時 | APIに渡す差分起点 |
| 変数 | This_Time | 今回の実行時刻 | フォルダ名/ファイル名/更新値に使い回す |
パイプラインの流れ(完成形イメージ)
アクティビティの並びは、基本的にこの形で安定します。
Set variable (This_Time = utcnow)
↓
Lookup (LastRecord.txt を読む / First row only)
↓
Set variable (Last_Time = Lookup結果)
↓
Copy (HTTP API → Data Lake へ JSON 保存) ※差分条件に Last_Time を渡す
↓
Copy (blank.txt → LastRecord.txt 上書き) ※Additional columns に This_Time を書き込む
(この更新Copyは「Succeeded」のときだけ実行)
最後の「更新用Copy」は、主要な処理がすべて成功したあとにだけ動くように依存関係を設定します。Copyが1つだけで完結しない場合(API呼び出しが複数、変換がある等)は、更新用Copyの依存元を「最終成功の判定点」にまとめるのがコツです。
手順:This_Timeを最初に固定する
まずSet Variableで今回の実行時刻を確定させます。
@utcnow()
ここでの実務上のポイントは、以降は utcnow() を直接呼ばずに This_Time を参照することです。たとえば、出力フォルダ名とファイル名を別々の式で作っているとき、実行が深夜帯にかかると「フォルダは前日、ファイル名は翌日」のような不整合が起こり得ます。This_Timeを固定しておけば、この手の事故をほぼ防げます。
手順:LookupでLastRecordから「前回成功日時」を読む
Lookupアクティビティの設定例です。
- ソース:LastRecord.txt(DelimitedText等)
- First row only:オン
取得値をLast_Timeに入れます。列名はデータセットの定義に依存するため、実際の出力プロパティに合わせて参照してください(例では Prop_1 とします)。
@activity('Lookup saved rundate').output.firstRow.Prop_1
APIが受け付ける形式に合わせて、ここでフォーマットを統一しておくのも有効です。例えば常にISO 8601(Z付き)で渡すなら、This_Timeも同形式に整えます。
@formatDateTime(variables('This_Time'), 'yyyy-MM-ddTHH:mm:ssZ')
手順:CopyでAPIにLast_Timeを渡して差分取得する
HTTPコネクタで呼ぶAPIのURL例(クエリパラメータ型)です。実際のAPI仕様に合わせて、エンコードやパラメータ名を調整してください。
https://api.example.com/orders?modifiedtimestamp=@{variables('Last_Time')}
より現実的には、以下のような「取りこぼし防止のバッファ」を入れる設計もよく使われます。
- API側の更新反映が遅延する
- タイムスタンプの精度が秒単位で同値が出る
- 境界条件(
>=か>か)で重複または欠落が起こり得る
たとえば「前回成功時刻から2分引いた値」を渡し、重複は保存後の処理で吸収する、という戦略です(下流が冪等なら強い)。
@addMinutes(variables('Last_Time'), -2)
ただし、これは「重複を許容できる設計」(キーでの重複排除、Upsert、同一データ上書き等)が前提です。重複が致命的な場合は、API側で安定キー(ID)と時刻の両方を条件にするなど、別の工夫が必要です。
手順:出力パス/ファイル名はThis_Timeで作る(保存の一貫性)
JSONの出力先を日時ごとに整理すると、障害調査や再実行が楽になります。例として、データセットにパラメータ dateparam を用意し、This_Timeを渡します。
フォルダパス例:
@formatDateTime(dataset().dateparam,'yyyy/MM/dd')
ファイル名例(タイムスタンプ+固定名):
@concat(formatDateTime(dataset().dateparam,'yyyyMMdd_HHmmss'), '_orders.json')
“毎回増え続ける”こと自体は問題ではなく、問題は「それを毎回列挙して判定する」ことです。保存は増えてよい、判定は増やさない、がコツです。
手順:成功したらLastRecordにThis_Timeを書き戻す(更新は成功時のみ)
最後に、状態ファイル(LastRecord.txt)を上書き更新します。やり方はシンプルで、Copyアクティビティをもう1つ用意します。
- Source:blank.txt(中身は何でもよいが1行あると扱いやすい)
- Sink:LastRecord.txt(上書き)
- Additional columns:ここに
This_Timeを書く
Additional columns設定例:
- 列名:
LastRunTime(任意) - 値:
@variables('This_Time')
@variables('This_Time')
そしてこの更新Copyは、API取得〜保存が成功したときだけ実行されるよう、依存関係(Success/Succeeded)でつなぎます。失敗時に更新しない=前回成功日時が保持される、というのがこの方法の強さです。
初回実行(LastRecordが空のとき)をどう始めるか
運用開始時は「前回成功日時」が存在しません。ここであいまいに始めると、いきなり大量取得したり、逆に差分が抜けたりします。現場でよく採るスタート方法は次の3つです。
| 開始方法 | 向くケース | やり方 | 注意点 |
|---|---|---|---|
| 手動で基準日時を入れる | すでに別経路でフルロード済み | LastRecord.txtに「フルロード最終時刻」を書く | 基準が不正だと欠落が出る |
| 一度だけフルロードしてから更新 | 初回だけ全件取得が許容 | 初回は固定の過去日時で実行し、成功したらThis_Timeを記録 | API負荷・時間がかかる |
| 過去N日で開始 | 多少の重複が許容、取りこぼし回避重視 | LastRecordに「N日前」を入れて開始 | 重複排除の仕組みが必要 |
迷ったら「取りこぼし回避」を優先し、重複は下流で吸収する設計が安全です。
運用でハマりやすい落とし穴と対策
落とし穴:並列実行(重なり実行)でウォーターマークが競合する
スケジュール間隔が短い、または前回実行が遅延すると、同じパイプラインが同時に動く可能性があります。この状態でLastRecordを更新すると、後から始まった実行が先に更新してしまい、データ欠落や重複の原因になります。
対策:
- パイプラインの同時実行数(Concurrency)を1にする
- トリガーをTumbling Windowにして「時間窓」を厳密に管理する
- どうしても並列が必要なら、ウォーターマークを「1本」ではなく「パーティション単位」(APIの種別や顧客単位など)に分ける
落とし穴:APIの境界条件(>= / >)で同一時刻が重複・欠落する
API側の更新タイムスタンプが秒単位だったり、同一秒に複数更新が入ると、>= だと重複し、> だと欠落する可能性が出ます。
対策:
- APIが「更新時刻+ユニークID」でソート・ページング可能なら、(時刻, ID)の複合ウォーターマークにする
- 難しければ「数分巻き戻し+重複排除」を採用する
- 保存先での重複排除(同一IDの最新だけ残す等)を前提にする
落とし穴:LastRecordのフォーマットがAPIの期待とズレる
APIがローカル時刻前提、ミリ秒必須、タイムゾーン表記必須など、細かい差で「差分が効かない」ケースがあります。
対策:
- 保存する時刻を「APIが受け付ける文字列形式」に固定する(保存前に整形)
- 本番前に、APIのレスポンスが変わる境界(秒・分・日)でテストする
方法②:ADF管理REST APIでパイプライン実行履歴から「最後に成功した実行」を取得する
2つ目の方法は、Azure Data Factoryの管理用REST APIを使い、パイプライン実行履歴から「Succeeded」の最新を取得するやり方です。状態ファイルを持たないため、ガバナンス上の理由で「外部に水位を置きたくない」場合に検討されます。
全体像:queryPipelineRunsで成功実行を検索して先頭のrunStartを使う
ざっくりした流れは次のとおりです。
- 検索期間(earliest / latest)を変数で作る
- Webアクティビティで
queryPipelineRunsを呼ぶ(Status=Succeededでフィルタ) - 返ってきた配列の先頭(新しい順)から
runStartを取り出す - その値をAPIの差分条件に使う
検索期間の作り方(連続失敗を想定して余裕を持たせる)
たとえば「最大でも数日連続で失敗しない」前提なら、少し長めに期間を取ります。
| 変数 | 式例 | 意図 |
|---|---|---|
| earliest | @adddays(utcnow(), -6) | 検索の下限(例:6日前) |
| latest | @adddays(utcnow(), 1) | 検索の上限(例:明日まで含める) |
「何日前まで探せば必ず成功実行が見つかるか」は運用次第です。検索期間が短すぎると、失敗が続いたときに履歴が見つからずパイプラインが止まります。逆に長すぎると応答が重くなる場合があるので、失敗の最大継続時間に合わせて調整します。
Webアクティビティの設定例(Managed Identityで管理APIを呼ぶ)
Webアクティビティで、Data Factory管理APIへPOSTします。URLは次の形式です(値は環境に合わせて置換してください)。
https://management.azure.com/subscriptions/{SUBSCRIPTION_ID}/resourceGroups/{RESOURCE_GROUP}/providers/Microsoft.DataFactory/factories/{DATA_FACTORY_NAME}/queryPipelineRuns?api-version=2018-06-01
認証はManaged Identity(MSI)を使うのが運用しやすいです。
- Authentication:Managed Identity(MSI)
- Resource:
https://management.azure.com/
このとき、ADFのマネージドIDに対して、対象のData Factory(またはリソースグループ)に対する閲覧権限が必要です(一般的にはReader以上が目安)。
Body例:PipelineNameとStatus=Succeededで絞り、RunStartの降順で取得
リクエストボディの例です。成功実行だけに絞り、最新が先頭に来るように並べます。
{
"lastUpdatedAfter": "@{variables('earliest')}",
"lastUpdatedBefore": "@{variables('latest')}",
"filters": [
{
"operand": "PipelineName",
"operator": "Equals",
"values": ["MYPIPELINENAME"]
},
{
"operand": "Status",
"operator": "Equals",
"values": ["Succeeded"]
}
],
"orderBy": [
{
"orderBy": "RunStart",
"order": "DESC"
}
]
}
レスポンスは value 配列で返り、新しいものが先頭に来る想定です。先頭の runStart を参照します。
@activity('Get runs').output.value[0].runStart
取り出した値を Last_Time_From_API のような変数に入れ、以降のAPI呼び出しに使います。
この方法のメリット/デメリット
| 観点 | メリット | デメリット・注意点 |
|---|---|---|
| 状態管理 | 状態ファイルが不要。履歴が真実になる | 履歴の保持期間・削除・権限に依存する |
| 統制 | 複数環境・複数パイプラインで統一的に扱いやすい | 管理API呼び出しの実装がやや複雑 |
| 障害耐性 | 成功履歴があれば復旧が容易 | 検索範囲が短いと「成功が見つからない」失敗モードが出る |
方法②は「仕組みとしては美しい」一方で、権限・保持・呼び出しコストを含めた運用設計が必要です。はじめて差分連携を作る場合は、まず方法①で素早く安定稼働させ、後から要件に応じて方法②へ拡張する、という進め方が現実的です。
どちらを選ぶべきか:実務の判断基準
迷ったときは、以下の基準で決めると失敗しにくいです。
- 最短で安定稼働させたい:方法①(状態ファイル)
- 権限設定を増やしたくない:方法①(ストレージだけで完結)
- 状態ファイルを持ちたくない(統制・監査):方法②(履歴参照)
- 複数パイプラインを履歴ベースで統一管理したい:方法②
- 連続失敗が長期化しやすい/履歴保持に不安:方法①(自前保存のほうが確実)
差分取得を“運用で強くする”ための追加の工夫
最後に、どちらの方式でも差分連携を壊れにくくするための実務的な工夫をまとめます。
API側の仕様を前提に「欠落しない」戦略を決める
差分連携の事故は、ほとんどが「境界(同一時刻、遅延反映、並び順)」で起こります。次の表のどれを採るかを先に決めておくと、運用が安定します。
| 戦略 | 取得条件の例 | 欠落リスク | 重複リスク | 向くケース |
|---|---|---|---|---|
| 厳密(ID併用) | (timestamp, id)でページング | 低 | 低 | APIが強力で、設計を作り込みたい |
| 巻き戻し+重複排除 | 前回成功時刻-数分 〜 現在 | 低 | 中 | APIが単純で、下流で冪等にできる |
| 単純差分(>=のみ) | modifiedtimestamp >= 前回成功 | 中 | 中 | データが少なく境界問題が起こりにくい |
保存先(Data Lake)の命名は“運用の武器”にする
差分取得は「取れるか」だけでなく、「困ったときに追えるか」も同じくらい重要です。おすすめの命名ポリシーは以下です。
- パイプライン名、API名、取得条件(任意)をパスに入れる
- 日付フォルダ+時刻ファイルで、実行単位が一目で分かるようにする
- 後段処理(Databricks / Synapse / Fabric等)で扱いやすい粒度に揃える
/raw/orders/yyyy/MM/dd/20251216_012345_orders.json
/control/watermark/last_success_orders.txt
更新ポイントは「本当に成功と言える地点」に置く
「APIから200が返った」だけで成功にしてしまうと、保存失敗・後処理失敗で欠落が出ます。ウォーターマーク更新は、次のいずれかを満たす地点に置くのが安全です。
- 必要なCopyがすべてSucceeded
- 後続の検証(ファイルサイズ>0、JSONの最小構造チェック等)まで成功
- 必要なら、後段のロード(Silver/Gold)まで成功したことを条件にする
まとめ:ADFの差分連携は「前回成功日時」をどう持つかで安定性が決まる
ADFには「前回成功実行日時」を直接返すシンプルな組み込み変数はないため、差分連携ではウォーターマーク管理が設計の肝になります。最初の一歩としては、専用ファイルに前回成功日時を保存し、成功時のみ上書き更新する方法(方法①)が最も作りやすく、運用も安定しやすい選択です。
一方、運用要件が進んで「状態ファイルを持ちたくない」「履歴ベースで統一したい」となったときには、ADF管理REST APIでSucceededの最新実行を検索する方法(方法②)が強力な選択肢になります。
どちらを選んでも、最終的に差が出るのは「境界条件」と「並列実行」と「更新のタイミング」です。ウォーターマーク更新を成功条件と一致させ、時刻はUTCで統一し、必要なら巻き戻し+重複排除の戦略を組み合わせることで、差分取得は長期運用に耐える仕組みに育ちます。

コメント