Cloud FunctionからDataformワークフローを連携実行する実装パターン
はじめに
GCSへのファイルアップロードをトリガーに、Cloud Functionでデータロード(L0層)を実行し、その後Dataformで変換処理(L1層)を自動実行するデータパイプラインを構築しました。
この記事では、Cloud FunctionからDataform APIを使ってワークフローを直接トリガーする実装方法を解説します。
実装の全体像
パイプラインの流れ
- GCSにファイルがアップロード → Workflowsがトリガー
- Cloud Functionが起動 → L0データロード処理を実行
- Dataform APIを呼び出し → L1変換処理を自動実行
このパターンにより、手動でのDataform実行が不要になり、完全に自動化されたデータパイプラインを実現できます。
実装内容
1. Dataformトリガー関数の追加
まず、main.py の上部に以下の関数を追加します。
# main.py の上部に追加
from google.cloud import dataform_v1beta1
def trigger_dataform_workflow(
project_id: str,
region: str,
repository_id: str,
workspace_id: str,
tag: str
):
"""
Dataformのワークフロー実行を、特定のタグでトリガーします。
この関数は、DataformのAPIを直接呼び出し、コンパイルと実行の2ステップを処理します。
実行の完了は待たずに、呼び出しを開始した時点で終了します(Fire and Forget)。
Args:
project_id (str): Google CloudプロジェクトID。
region (str): Dataformリポジトリが存在するリージョン。
repository_id (str): Dataformのリポジトリ名。
workspace_id (str): Dataformのワークスペース名。
tag (str): 実行したいDataformのアクションに付けられたタグ。
"""
try:
logger.info(f"Starting Dataform invocation with tag: {tag}")
client = dataform_v1beta1.DataformClient()
repository_path = f"projects/{project_id}/locations/{region}/repositories/{repository_id}"
# 1. Dataformのコードをコンパイルするリクエストを作成
compilation_result = client.create_compilation_result(
parent=repository_path,
compilation_result={
"workspace": f"{repository_path}/workspaces/{workspace_id}"
},
)
logger.info(f"Created Dataform compilation result: {compilation_result.name}")
# 2. コンパイル結果を使って、ワークフロー実行をリクエスト
invocation = client.create_workflow_invocation(
parent=repository_path,
workflow_invocation={
"compilation_result": compilation_result.name,
"invocation_config": {"included_tags": [tag]},
},
)
logger.info(f"Created Dataform workflow invocation: {invocation.name}")
return invocation.name
except Exception as e:
logger.exception("Failed to trigger Dataform workflow.")
raise
ポイント解説
Fire and Forget パターン
この実装では、Dataformの実行完了を待たずに処理を返します(非同期実行)。理由は以下の通りです:
- Cloud Functionのタイムアウトを気にせずに済む
- L1変換が長時間かかる場合でも問題ない
- 実行状態の監視は別途Cloud LoggingやDataform UIで確認可能
2ステップ実行
Dataform APIは以下の2段階で実行されます:
- コンパイル: SQLDAファイルを解析し、実行可能なSQLを生成
- 実行: 生成されたSQLをBigQueryで実行
2. メインエントリポイントの実装
次に、既存の trigger_gcs_to_workflows 関数を以下で置き換えます。
# main.py の既存の trigger_gcs_to_workflows 関数を、以下の内容で完全に置き換え
def trigger_gcs_to_workflows(request: Request):
"""
WorkflowsからHTTPで呼び出されるメインのエントリポイント。
L0ロードを実行し、その後Dataformのワークフローをトリガーします。
"""
logger.info("==== Cloud Function Main Entrypoint Called ====")
try:
# Workflowsから渡されたJSONペイロードを解析
payload = request.get_json(silent=True) or {}
bucket = payload.get("bucket")
name = payload.get("name")
if not bucket or not name:
logger.warning("Bad request, missing bucket or name in payload")
return make_response(("Bad request: need JSON {bucket,name}", 400))
logger.info(f"Payload parsed: gs://{bucket}/{name}")
# --- ステップ1: L0ロードの実行(既存のロジック) ---
_prepare_local_files_from_gcs(bucket, name)
run()
logger.info("L0 Load process completed successfully.")
# --- ステップ2: DataformのL1変換をトリガー ---
try:
# GCSのパスからカテゴリ名を抽出
# 例: "client_data/campaigns/file.xlsx" -> "campaigns"
category = name.split('/')[1]
dataform_tag = f"l1_{category}"
# Dataformの実行に必要な設定
region = "asia-northeast1"
repository_id = "data_pipeline"
workspace_id = "prd"
# Dataform実行関数を呼び出し
invocation_name = trigger_dataform_workflow(
project_id=config.project_id,
region=region,
repository_id=repository_id,
workspace_id=workspace_id,
tag=dataform_tag
)
response_message = f"L0 Load OK. Triggered Dataform with invocation: {invocation_name}"
logger.info(response_message)
return make_response((response_message, 200))
except (IndexError, TypeError) as e:
# パスからカテゴリ名が取得できなかった場合
logger.warning(f"Could not determine Dataform category from path '{name}'. Dataform not triggered. Error: {e}")
return make_response(("L0 Load OK. Dataform not triggered due to invalid path.", 200))
except Exception as e:
logger.exception("Unhandled exception in main entrypoint.")
return make_response((f'{{"error":"{str(e)}"}}', 500))
ポイント解説
動的タグ生成
GCSのファイルパスからカテゴリを抽出し、対応するDataformタグを動的に生成します:
# "client_data/campaigns/report.xlsx" の場合
category = name.split('/')[1] # "campaigns"
dataform_tag = f"l1_{category}" # "l1_campaigns"
これにより、カテゴリごとに異なるDataform変換処理を自動選択できます。
エラーハンドリング
パスが想定と異なる形式の場合でも、L0ロードは成功として扱い、Dataformのトリガーのみスキップします。これにより、部分的な失敗で全体が停止することを防ぎます。
3. 依存パッケージの追加
requirements.txt に以下を追加します:
google-cloud-dataform
デプロイ手順
- 上記のコード変更を適用
-
requirements.txtを更新 - Cloud Functionをデプロイ
gcloud functions deploy your-function-name \
--runtime python39 \
--trigger-http \
--entry-point trigger_gcs_to_workflows
実行フロー図
GCS Upload
↓
Workflows (Trigger)
↓
Cloud Function
├─ L0 Load (BigQueryへ生データ投入)
└─ Dataform API Call
├─ Compilation (SQLDAをSQLに変換)
└─ Invocation (BigQueryで実行)
この実装のメリット
1. 完全自動化
手動でのDataform実行が不要になり、データが届いた瞬間から変換まで自動で完了します。
2. カテゴリ別の柔軟な処理
ファイルパスに基づいて適切なDataform変換を自動選択できます。
3. Fire and Forget
長時間実行にも対応でき、Cloud Functionのタイムアウトを気にする必要がありません。
4. エラーの部分的許容
パス解析に失敗してもL0ロードは完了するため、データ欠損を防げます。
まとめ
Cloud FunctionからDataform APIを直接呼び出すことで、GCSアップロードからデータ変換までの完全自動パイプラインを実現できました。
この実装パターンは、以下のようなユースケースに応用できます:
- マルチテナントでクライアント別に処理を分岐
- データソース種別ごとに異なる変換ロジックを適用
- 時間帯やファイルサイズに応じた処理の最適化
ぜひ参考にしてみてください!
Discussion