Blob Storage に着地した CSV を複数の BigQuery 先テーブルへ取り込む ADF(Azure Data Factory)パイプラインの成否を、外部メタストア無し・Blob 内の 1 つの「ステータス JSON」で管理する——本記事はそのための設計指針と実装を、最新状態の上書きと履歴追記の両パターンで具体的に解説します。
背景とゴール
要件はシンプルです。Blob に CSV が到着し、ADF が BigQuery へロードする複数パイプラインを起動。各パイプラインの完了可否を 1 つの JSON ファイルで管理したい。RDB など外部メタストアは使えないため、Blob 上のファイルだけで実現する必要があります。さらに、運用方針として「最新状態だけを保持(上書き)」と「履歴を残す(追記)」のどちらにも対応できることが求められます。
解決アプローチの全体像(早見表)
| ニーズ | 現実的な実装方法 | 補足ポイント |
|---|---|---|
| 最新状態だけ保持 (過去履歴は不要) | パイプライン終端に Web Activity を配置し、PUT で Blob に JSON をそのまま書き込み(Block Blob 上書き)。 本文(Body)は ADF の式で構築した JSON 文字列をそのまま送信。 認証は Managed Identity(MSI)で完結。 (代替)Copy Activity+ダミー空ファイルでの上書きも可能だが、挙動の明確さ・保守性の観点で Web Activity が推奨。 | ADF はネイティブな「部分更新(追記)」機能を持たないため、上書きか新規作成が基本。 環境分割(DEV/QA/PROD)はパスを分けるだけで衝突回避可能。 Blob バージョン管理を有効化すれば、上書きでも自動で旧版は退避される。 |
| 履歴を残したい/追記したい | 方式A: 毎回タイムスタンプ付き新規ファイルを出力し、定期バッチ(Mapping Data Flow など)で 1 つの履歴ファイルへ集約。 方式B: ADF から Azure Function / Logic Apps / Databricks を呼び出し、JSON を コードで追記(JSONL 形式推奨)して即時に 1 ファイルへ蓄積。 | 方式Aは ADF 内で完結、ただしファイル数増加に応じたライフサイクル管理が必要。 方式Bは外部コード依存だが、リアルタイムに 1 ファイルへ追記できる。 JSON 配列より JSON Lines(JSONL)が追記に強い。 |
アーキテクチャの要点
- イベントドリブン:Blob に CSV が到着 → ADF トリガー(イベントベース or タイムベース)。
- BigQuery ロード:ADF の Google BigQuery コネクタで複数テーブルへ投入。
- 状態管理:パイプライン終端で「最新上書き」もしくは「履歴追記」へ分岐。
- 権限:パイプライン実行 ID(Managed Identity または Service Principal)に Storage Blob Data Contributor を付与。
設計の指針(どちらの方式を選ぶべきか)
- 監査・再現性が不要、運用ダッシュボードだけで十分 → 最新上書きが最もシンプルで低コスト。
- 監査・トラブルシュート・SLA 検証が必要 → 履歴追記を選択。JSONL 形式で 1 ファイル運用、もしくはタイムスタンプ付きファイル+定期集約。
- 書き込み頻度が高く衝突の恐れがある → 方式A(新規ファイル)+日次集約か、方式B(追記)+ETag 条件による同時更新制御を導入。
- リスク軽減:最新上書きでも Blob バージョン管理をオンにしておくと巻き戻しが容易。
実装:最新状態だけ保持(上書きパターン)
方式1-1:Web Activity で Block Blob に PUT(推奨)
ADF の Web Activity は任意の HTTP エンドポイントに対し、Managed Identity(MSI)で Azure AD 認証付きリクエストを送れます。Azure Storage(Blob)は AAD トークンを受け付けるため、PUT で JSON 本文をそのまま Blob に書き込めます。外部コード不要・高い可読性・ヘッダ制御が可能という理由で、この方式を第一候補にしてください。
前提条件
- 対象ストレージアカウントに対して、ADF の Managed Identity(または実行に使う連携 ID)へ Storage Blob Data Contributor を付与。
- (任意)ストレージ側で Blob バージョン管理と Soft Delete を有効化。
パイプラインのパラメータ/変数設計例
| 名前 | 種類 | 例 | 説明 |
|---|---|---|---|
| storageAccount | String(Param) | stacc001 | Blob のアカウント名 |
| container | String(Param) | status | ステータス格納用コンテナ |
| statusPath | String(Param) | dev/pipelines/sales | 環境や業務別のディレクトリ |
| fileName | String(Param) | sales_20250915.csv | 対象 CSV 名(可観測性向上用) |
| targets | Array(Param) | [“dev.sales_hdr”,”qa.sales_hdr”] | BigQuery のロード先テーブル配列 |
| status | String(Var) | COMPLETED / FAILED | 成功・失敗・部分成功など |
Web Activity の設定
- URL(式)
@concat('https://', pipeline().parameters.storageAccount, '.blob.core.windows.net/',
pipeline().parameters.container, '/',
pipeline().parameters.statusPath, '/status.json')
- Method:
PUT - Authentication:
Managed Identity(Resource:https://storage.azure.com/) - Headers(推奨)
{
"x-ms-blob-type": "BlockBlob",
"x-ms-version": "2021-12-02",
"Content-Type": "application/json; charset=utf-8"
}
- Body(式:JSON をパイプラインから組み立て)
@concat(
'{',
'"fileName":"', pipeline().parameters.fileName, '",',
'"processedOn":"', utcnow(), '",',
'"targets":', string(pipeline().parameters.targets), ',',
'"status":"', variables('status'), '",',
'"adfRunId":"', pipeline().RunId, '"',
'}'
)
この設定で、指定パスに status.json が存在しなければ新規作成、存在すれば完全上書きされます。競合制御が必要で「存在しない場合のみ作成」したいなら、ヘッダに If-None-Match: * を追加すると安全に新規作成のみを許可できます(既存があれば 412 応答)。
サンプル JSON スキーマ(上書き運用)
{
"fileName": "sales_20250915.csv",
"processedOn": "2025-09-16T14:30:00Z",
"targets": [
"dev.sales_hdr",
"qa.sales_hdr"
],
"status": "COMPLETED",
"adfRunId": "xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx"
}
方式1-2:Copy Activity+ダミー空ファイル(簡易代替)
Copy Activity の Source に Blob 上の空ファイルを指定し、Sink に status.json を指定する方法もあります。事前に空ファイル(例:empty.json)を用意し、パイプラインの上流で式から JSON を構築してメタデータとして渡すパターンです。ただし、ヘッダ制御や条件付き書き込みが難しく、挙動がわかりづらくなるため、Web Activity を優先することを推奨します。
実装:履歴を残す/追記するパターン
方式A:タイムスタンプ付き新規ファイル → 定期マージ(ADF 内で完結)
各実行で新規ファイル(yyyy/MM/dd/HHmmssZ_status.json など)を書き出し、日次や時間ごとの集計ジョブで 1 つの履歴ファイルへ集約します。最も ADF ネイティブで、スケールや同時実行に強い方式です。
命名規則の例
/status/dev/pipelines/sales/
2025/11/02/20251102T120102Z_status.json
2025/11/02/20251102T121503Z_status.json
2025/11/02/20251102T124259Z_status.json
...
新規ファイルの書き出し(Web Activity 版)
- URL(式)
@concat(
'https://', pipeline().parameters.storageAccount, '.blob.core.windows.net/',
pipeline().parameters.container, '/',
pipeline().parameters.statusPath, '/',
formatDateTime(utcnow(),'yyyy/MM/dd/'),
formatDateTime(utcnow(),'yyyyMMddTHHmmssZ'), '_status.json'
)
Body は上書きパターンと同様に組み立てます。
定期マージ(Mapping Data Flow の一例)
- Source:ワイルドカードで
**/*_status.jsonを読み込み。フォルダ年/月/日を パスから派生列 として抽出し、パーティションに活用。 - Flatten/Select:JSON を正規化(
targetsを explode するなど)。 - Sink:履歴ファイルを JSONL(改行区切り JSON) もしくは Parquet で 1 本に書き出し。ワークロードが重ければ、日次ファイル+月次集約の二段階に分割。
ライフサイクル管理
- ストレージの ライフサイクル管理ポリシーで、90~180 日より古い「原本(タイムスタンプ付き)」を自動削除、あるいは低コスト階層へ移行。
- 履歴の
history.jsonlは直近 N か月を保持。さらに古い期間は月別アーカイブへ。
方式B:Azure Function で JSONL にリアルタイム追記
毎回 1 行の JSON を 追記できると、常に 1 ファイルを参照すれば最新~過去の全履歴が取得できます。Append Blob と JSONL の組み合わせはシンプルで高効率です。以下は Python の最小構成例です。
Azure Function(HTTP トリガー)の例
import azure.functions as func
import json, datetime, uuid
from azure.storage.blob import AppendBlobClient
from azure.identity import DefaultAzureCredential
def main(req: func.HttpRequest) -> func.HttpResponse:
try:
event = req.get_json()
except Exception:
return func.HttpResponse("invalid json", status_code=400)
# 監査項目を補完
event.setdefault("processedOn", datetime.datetime.utcnow().isoformat() + "Z")
event.setdefault("eventId", str(uuid.uuid4()))
# 例: fileName / targets / status / adfRunId などを ADF 側から渡す
account = "stacc001" # 環境に合わせて設定(KeyVault 推奨)
container = "status"
blob_name = "dev/pipelines/sales/history.jsonl"
credential = DefaultAzureCredential()
client = AppendBlobClient(
account_url=f"https://{account}.blob.core.windows.net/",
container_name=container,
blob_name=blob_name,
credential=credential
)
try:
if not client.exists():
client.create() # 初回だけ作成
line = json.dumps(event, ensure_ascii=False)
client.append_block((line + "\n").encode("utf-8"))
return func.HttpResponse("ok", status_code=200)
except Exception as e:
return func.HttpResponse(str(e), status_code=500)
ADF からの呼び出し
- Azure Function Activity(または Web Activity)で HTTP POST。
- Body は上書きパターンと同様の JSON。
Content-Type: application/json。 - 関数の認証は Function Key / AAD いずれでも可。企業ポリシーに従う。
JSONL のメリット
- 追記が O(1)。ファイル全体の読み書きが不要。
- 破損時の影響局所化(行単位)。
- 下流処理(Athena/BigQuery/Dataproc/ADF Data Flow)で扱いやすい。
スキーマ設計の実践
最新のみ(上書き)
{
"fileName": "sales_20250915.csv",
"processedOn": "2025-09-16T14:30:00Z",
"targets": ["dev.sales_hdr", "qa.sales_hdr"],
"status": "COMPLETED",
"adfRunId": "xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx",
"metrics": {
"rowsRead": 125034,
"rowsLoaded": 125034,
"durationSec": 97
}
}
履歴(JSONL 1 行例)
{"eventId":"6a3c...","fileName":"sales_20250915.csv","processedOn":"2025-09-16T14:30:00Z","targets":["dev.sales_hdr","qa.sales_hdr"],"status":"COMPLETED","adfRunId":"...","retries":0,"error":null}
ステータス値の指針
COMPLETED:全テーブルにロード成功。PARTIAL:一部テーブルのみ成功(errorsに失敗対象を列挙)。FAILED:致命的失敗。error.codeとerror.messageを必ず出力。SKIPPED:入力なし、または再実行でスキップ。
エラー耐性と実行ポリシー
- 実行タイミング:パイプライン末尾に「On Success」で上書き/追記すれば、全体成功時のみ
COMPLETEDを残せます。On Completion であれば失敗時にもFAILEDを残せます。 - 再実行:
adfRunIdとfileNameをキーに、履歴から重複検出。JSONL では同一eventIdを避ける。 - 衝突回避:履歴追記の同時実行は Append Blob の直列化に任せるか、Block Blob 上書きなら ETag 条件(
If-Match)を使う。 - 監査:Blob の バージョン管理を有効化すれば、最新上書きでも過去バージョンから差分検証可能。
セキュリティ(権限・秘匿情報)
- 実行 ID に Storage Blob Data Contributor を付与(最小権限)。
- 関数や外部コードが必要な場合は Managed Identity + Key Vault で資格情報を秘匿。
- データ分類に応じてコンテナ・パスを分離(
public/internal/restricted)。 - 監査ログ(Azure Monitor / Storage Diagnostics)を有効化。
環境分離とフォルダ設計
/status/
dev/pipelines/sales/status.json
qa/pipelines/sales/status.json
prod/pipelines/sales/status.json
# 履歴運用の場合
prod/pipelines/sales/history.jsonl
prod/pipelines/sales/2025/11/02/20251102T121503Z_status.json
環境ごとにディレクトリを分けるだけで衝突を防げます。プロジェクトやデータドメイン単位の分割(pipelines/{domain}/{job})にしておくと将来のスケールにも耐えます。
パフォーマンス/コストの要点
- 上書きパターン:常に 1 ファイルのみ読み書きするため、トランザクション数・データ量とも最小。監視・参照も軽い。
- 履歴追記:長期運用でファイル肥大化しやすい。一定サイズ(例:数百 MB)でローテーション、もしくは日次ファイル化が有効。
- 一覧 API コスト:方式A(新規ファイル多数)ではリスト操作が増える。パーティション設計とライフサイクル管理で抑制。
検証チェックリスト
- MSI で
PUTできること(403/401 が出ない)。 Content-Type: application/jsonが付いていること。- JSON が厳密に妥当(全てのキーがダブルクォート)。
- 同時書き込みに対する期待挙動(最後の勝ち / エラーで弾く)がテストできていること。
- 履歴の集約またはローテーションが自動化されていること。
よくある落とし穴と回避策
- BOM 付き UTF-8:不要な BOM は JSON パーサでエラー要因。UTF-8(BOM なし)で書き込む。
- 「部分更新」期待:ADF には PATCH 的な部分更新はありません。上書き or 新規ファイル or 外部コードの 3 択です。
- 「追記」= Append Blob の誤解:Append Blob 自体は追記できますが、ADF 単体で安全に扱う術が乏しいため、Azure Function 経由を推奨。
- 巨大 JSON 配列:履歴を配列に溜めると毎回全体読み書きが必要。JSONLへ切替える。
- マルチ環境衝突:環境別ディレクトリを切り、ファイル名の固定化を避ける(履歴は時刻ベース名)。
運用監視のヒント
- ステータス JSON に
metrics(処理行数・所要秒・エラー件数)を含めると可視化が容易。 - Blob の変更フィード(Change Feed)を使えば、最新ファイルの変更検知から通知を飛ばす運用も可能。
- ダッシュボードは最終的に
status.jsonとhistory.jsonlの 2 ファイルだけを参照する構成にする。
実装スニペット集
(再掲)上書き用 JSON の組み立て(ADF 式)
@concat(
'{',
'"fileName":"', pipeline().parameters.fileName, '",',
'"processedOn":"', utcnow(), '",',
'"targets":', string(pipeline().parameters.targets), ',',
'"status":"', variables('status'), '",',
'"adfRunId":"', pipeline().RunId, '"',
'}'
)
Web Activity のヘッダ例
{
"x-ms-blob-type": "BlockBlob",
"x-ms-version": "2021-12-02",
"Content-Type": "application/json; charset=utf-8",
"If-None-Match": "*" // 既存がある場合はエラーにしたい時に使用
}
Mapping Data Flow(履歴集約の概念)
# Source: **/*_status.json を読み込み
# Derived Columns: dt = toTimestamp(filePath('fullPath')の yyyy/MM/dd 部分)
# Flatten: targets を explode(target 列を展開)
# Sink: history.jsonl(single file または パーティション書き出し)
まとめ
- ADF だけで確実にできるのは「新規作成または上書き」。Web Activity の
PUTがシンプルで推奨。 - 履歴を追記したい場合は、タイムスタンプ付き新規ファイル+定期マージ(純 ADF)か、Azure Function で JSONL に追記(即時 1 ファイル集約)の二択が現実解。
- 最新状態だけでよいなら、Web Activity + MSI で Blob に 1 ファイル上書きが最小コスト・最小構成。
- 安全運用には Blob バージョン管理、ライフサイクル管理、ETag 条件の活用が効きます。
付録:要求に対する実装手順(Checklist)
- 実行 ID に Storage Blob Data Contributor を付与。
- Blob のバージョン管理と Soft Delete を有効化(任意だが推奨)。
- パイプライン末尾に Web Activity を追加、
PUTでstatus.jsonへ上書き。 - 履歴が必要なら
- 方式A:時刻付きファイルを吐き出し、Mapping Data Flow で日次/月次に集約。
- 方式B:Azure Function を用意し JSONL に追記、ADF から POST。
- 監視:
status.jsonのstatusとprocessedOnをダッシュボード化。失敗時は通知。 - 保守:履歴のローテーション/圧縮/アーカイブを自動化。

コメント