🐙

Cloud FunctionからDataformワークフローを連携実行する実装パターン

に公開

はじめに

GCSへのファイルアップロードをトリガーに、Cloud Functionでデータロード(L0層)を実行し、その後Dataformで変換処理(L1層)を自動実行するデータパイプラインを構築しました。

この記事では、Cloud FunctionからDataform APIを使ってワークフローを直接トリガーする実装方法を解説します。

実装の全体像

パイプラインの流れ

  1. GCSにファイルがアップロード → Workflowsがトリガー
  2. Cloud Functionが起動 → L0データロード処理を実行
  3. 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段階で実行されます:

  1. コンパイル: SQLDAファイルを解析し、実行可能なSQLを生成
  2. 実行: 生成された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

デプロイ手順

  1. 上記のコード変更を適用
  2. requirements.txt を更新
  3. 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