Data Agent Kit と Antigravity IDE を使用した不正行為検出パイプライン

1. はじめに

大量の決済を処理する Cymbal Financial の データ サイエンティストであるとします。決済の遅延が多発しており、コンプライアンス チームは組織的な不正行為を疑っています。未加工のクリアリングハウス トランザクション ログを取り込み、データをクリーンアップし、ML モデルをトレーニングし、バッチ推論を実行して、リスクの高いトランザクションを手動監査用の Cloud Spanner レビューキューにシンクするパイプラインを構築する必要があります。

通常、これには、繰り返し設定コード(Spark ノートブック、dbt 構成、トレーニング スクリプト、Airflow DAG)の作成に数日を要し、コンソール インターフェースとエディタの間でコンテキストを常に切り替える必要があります。

この Codelab では、Antigravity IDE 内の Google Cloud Data Agent Kit(DAK)を使用して、エージェントとペアプログラミングを行います。会話形式の自然言語を使用して、エージェントは Spark ノートブックの生成、dbt プロジェクトのコンパイル、推論ループの構築、Managed Service for Apache Airflow を使用したワークフローのオーケストレーションを支援します。

演習内容

  • Managed Service for Apache Spark(Spark サーバーレス)を使用して、Cloud Storage から クリアリングハウス ログを取り込み、BigQuery テーブルに保存します。
  • dbt を使用してトランザクションを重複除去して正規化し、クリーンなデータレイヤ(Raw、Staging、Enriched)を確立します。
  • Spark Serverless で分散ランダム フォレスト分類モデルをトレーニングする(RandomForestClassifier)。
  • 新しいトランザクションでバッチ推論を実行し、高リスクのアラートを Cloud Spanner に直接書き込みます。
  • Managed Service for Apache Airflow と IDE 内のインタラクティブ DAG モニタリングを使用して、パイプライン全体をオーケストレート、視覚的に構成、デプロイします。

必要なもの

  • ウェブブラウザ(Chrome など)
  • 課金が有効になっている Google Cloud プロジェクト(ハンズオンラボには新しい専用プロジェクトを使用することをおすすめします)。
  • SQL、Python、PySpark に関する基本的な知識。
  • Google AI Pro サブスクリプション付きの Antigravity IDE(推奨)

この Codelab で作成するリソースの費用は 5 ドル未満です。ラボの最後にあるクリーンアップの手順に沿って、プロビジョニングされたリソースを削除してください。

2. 環境のセットアップ

ラボを開始するには、ブートストラップ スクリプトを実行します。このスクリプトは、必要な GCP API を自動的に有効にし、取り込み用の Cloud Storage バケットを作成し、トランザクションとディレクトリのモック データセットを生成し、参照ディレクトリを BigQuery に読み込み、Cloud Spanner と Managed Service for Apache Airflow(旧称 Cloud Composer)のバックグラウンド プロビジョニングを開始します。

プロジェクトの選択または作成

Google Cloud コンソールで、既存のプロジェクトを選択するか、新しいプロジェクトを作成します。

お支払いを確認する

Google Cloud プロジェクトの課金が有効になっていることを確認します。詳しくは、こちらのガイドをご覧ください。

設定スクリプトを実行する

Google Cloud Shell(または Google Cloud CLI で構成されたローカルシェル)を使用して、環境設定を起動します。

  1. Google Cloud Console を開きます。
  2. 右上のツールバーにある [Cloud Shell をアクティブにする] をクリックします。

Cloud Shell を開く

  1. Cloud Shell ターミナルで、アクティブなプロジェクトを構成します。
gcloud config set project <<YOUR_PROJECT_ID>>
export PROJECT_ID=$(gcloud config get-value project)
  1. Codelab リポジトリのクローンを作成し、scripts フォルダに移動します。
cd ~/
git clone --filter=blob:none --no-checkout https://github.com/GoogleCloudPlatform/devrel-demos.git
cd ~/devrel-demos
git sparse-checkout init --cone
git sparse-checkout set codelabs/agentic-data-labs/data-science
git checkout main
cd codelabs/agentic-data-labs/data-science/scripts
  1. ブートストラップ設定スクリプトを実行して、すべてのリソースを us-central1 にデプロイします。
chmod +x setup.sh setup_spanner.sh setup_composer.sh
export REGION=us-central1
./setup.sh
  1. スクリプトが終了すると、BigQuery データセットと Cloud Storage バケットの準備ができたことを示す概要出力が表示されます。バックグラウンドで、Cloud Spanner(約 2 分)と Managed Airflow(約 20 分)のプロビジョニングが続行されます。進行状況は、次のコマンドを実行していつでもモニタリングできます。
tail -f /tmp/spanner_setup.log
tail -f /tmp/composer_setup.log

Antigravity IDE を開く

  1. Google Antigravity のダウンロード ページから Antigravity IDE をダウンロードしてインストールします。
  2. Antigravity IDE を起動します。
  3. ローカルマシンに新しい空のフォルダ(agentic-data-labs など)を作成し、[フォルダーを開く] を選択して IDE で開きます。これが、Codelab のローカル ワークスペースとして機能します。

Antigravity IDE プロジェクト フォルダを構成する

Data Agent Kit 拡張機能をインストールする

Google Cloud Data Agent Kit 拡張機能は、エディタ内で Google Cloud データ サービスと深く統合されており、コンテキストを切り替えることなく BigQuery、Cloud SQL、Cloud Storage などを操作できます。

  1. Antigravity IDE で、画面の左端にあるアクティビティ バーの [拡張機能] アイコン(4 つの正方形のようなアイコン)をクリックします。
  2. [拡張機能] ペインの上部にある検索バーに「Google Cloud Data Agent Kit」と入力します。
  3. googlecloudtools が公開した Google Cloud Data Agent Kit という名前の拡張機能を見つけます。
  4. [Install] ボタンをクリックします。
  5. 「パブリッシャー「googlecloudtools」とその拡張機能を信頼しますか?」というメッセージが表示されることがあります。[Trust Publishers & Install] をクリックして続行します。

Data Agent Kit 拡張機能をインストールする

インストールが完了すると、Antigravity IDE の左端にあるアクティビティ バーに、新しい Google Cloud Data Agent Kit アイコンが表示されます。

  1. 「Google Cloud Data Agent Kit へようこそ」というタイトルのオンボーディング ページが自動的に開きます。Cloud アカウントにログインしていない場合は、画面の指示に沿ってアクセスを許可します。
  2. [構成の概要] セクションで、プロジェクト フィールドを見つけます。プルダウンをクリックして、Google Cloud プロジェクトを選択します。リージョンを us-central1 に設定します。[Configure MCP Servers] を選択します。

Data Agent Kit 拡張機能の初期構成

  1. [MCP サーバーを構成] を選択します。[MCP 構成] ペインで、次のリモート MCP サーバーが有効になっていることを確認します。
    • BigQuery
    • Spanner
    • Notebooks

[使ってみる] をクリックします。

MCP サーバーを構成する

構成オプションを確認する

設定が完了すると、[Google Cloud Data Agent Kit を使ってみる] ページが表示されます。

  1. [設定と構成] で、[開始] をクリックします。
  2. [Data Agent Kit Configuration] パネルが開きます。タブを確認します。
    • プロジェクトとリージョン: 選択したプロジェクト ID を確認し、設定スクリプトがすべての必要な API(Compute Engine、Cloud Storage、BigQuery、Spanner など)を有効にしていることを確認します。
    • BigQuery: BigQuery クエリのデフォルトの場所を構成します。リージョン us-central1 を使用します。
    • MCP サーバーを構成する: AI エージェントがデータを安全に操作できるようにする有効な MCP サーバー(BigQuery、Notebooks、Spanner など)を表示します。
    • スキル: エージェントに複雑なデータタスク用の特別な機能を提供する事前構築済みのスキルを確認します。

Data Agent Kit 設定パネル

セクションのまとめ: ブートストラップ スクリプトを実行して GCS アセットと BigQuery アセットを作成しました。Spanner と Airflow はバックグラウンドでビルドされます。次に、Antigravity IDE でプロジェクトを開き、Google Cloud Data Agent Kit 拡張機能を有効にしました。これで、最初のノートブックを作成する準備が整いました。

3. Spark Serverless を使用して未加工ログを取り込む

このセクションでは、未加工の JSON トランザクション ログをデータレイクに取り込みます。Managed Service for Apache Spark(Spark Serverless)は、BigQuery のネイティブ ストレージに直接接続します。標準の BigQuery コネクタを使用して、表形式のデータを管理し、直接クエリと分析を有効にします。

事前構成済みの Spark サーバーレス ランタイムを確認する

Spark コードを実行する前に、設定スクリプトによって事前構成された Serverless ランタイム テンプレートを調べます。このテンプレートは、ターゲット実行環境のバックエンドを定義し、必要なコネクタの依存関係をバンドルします。

  1. IDE のアクティビティ バーで、[Google Cloud Data Agent Kit] パネルを開きます。
  2. [Apache Spark] プルダウン メニューを開き、[Serverless] を開きます。
  3. fraud-pipeline-runtime を右クリックして [Profile] を選択し、エディタで構成ビューを開きます。
  4. [プロファイル] タブで、下にスクロールして [プロパティ] を展開し、環境にアタッチされているカスタム依存関係を調べます。
    • spark.jars: gs://spark-lib/spanner/spark-3.5-spanner-1.4.0.jar を含みます。これは、Spark Spanner コネクタを使用して、ラボの後半で Spark ジョブが推論結果を Cloud Spanner に直接書き込めるようにします。(注: Dataproc Serverless には、Google Cloud の Spark BigQuery コネクタがデフォルトで含まれているため、BigQuery テーブルの読み取りと書き込みに jar の追加構成は必要ありません)。

Spark サーバーレス ランタイム プロパティを確認する

  1. 左側の [Interactive Sessions] タブに注目してください。まだコードを実行していないため、現在は空になっています。次のステップでノートブックを実行すると、ライブ サーバーレス コンピューティング セッションが動的にプロビジョニングされ、ここに表示されます。

Data Agent Kit を使用してデータを取り込む

Spark セッションを手動で構成したり、PySpark ローディング スクリプトをゼロから記述したりする代わりに、Data Agent Kit を使用してエージェントとペアプログラミングを行います。

  1. 右上にあるツールバーの [Toggle Agent] アイコンをクリックして、[Agent Chat] ペインを開きます。
  2. 次のプロンプトをチャットに貼り付けます(${PROJECT_ID} は実際の Google Cloud プロジェクト ID に置き換えてください)。
Create a PySpark notebook (01_ingestion.ipynb) to ingest JSON transaction logs
from gs://${PROJECT_ID}-fin-clearing-raw/ into a BigQuery table
`${PROJECT_ID}.transactions_dataset_evals.raw_transactions`
using the Spark BigQuery connector (`format("bigquery")`) with overwrite mode.
  1. エージェントがバックグラウンド検証コマンドの実行権限を求めてきた場合(「このコマンドの実行を許可しますか?」など)、提案されたコマンドを確認し、[はい、今回のみ許可します](または [はい、常に許可します])を選択します。
  2. エージェントがファイルの生成を完了したら、チャット ペインの下部にある青い [すべて承認] ボタン(またはチェックマーク アイコン)をクリックして、notebooks/01_ingestion.ipynb をワークスペースに保存します。

取り込みノートブックを生成するエージェント

ノートブックを確認して実行する

  1. IDE で新しく生成された notebooks/01_ingestion.ipynb を開きます。
  2. BigQuery コネクタの書き込みロジックの PySpark コードを確認します。
  3. IDE のノートブック ツールバーで [すべてを実行] をクリックします。
  4. リモート Spark ノートブックを初めて実行する場合は、IDE でローカル依存関係のインストールを求めるメッセージが表示されることがあります。プロンプトが表示されたら、[リモート Spark カーネルの依存関係をインストール] をクリックしてインストール ダイアログを確認し、[すべて実行] をもう一度クリックします。
  5. [カーネルを選択] プルダウン メニューで、[Remote Spark Kernels] -> [fraud-pipeline-runtime on Serverless Spark] を選択します。(ヒント: 事前構成済みのランタイム テンプレートが表示されない場合は、カーネル選択ツールのプルダウンの右上にある更新アイコンをクリックして、使用可能なリモート カーネルを再読み込みします)。
  6. エディタの左下にあるステータスバーを確認します。Connecting to kernel: fraud-pipeline-runtime on Serverless Spark... が表示されます。これは Spark サーバーレス ランタイム カーネル バックエンドの初回起動であるため、プロビジョニングと起動に数分かかります。
  7. カーネルの接続が完了すると、ノートブックはすべてのセルを順番に自動的に実行し、未加工のトランザクション ログを BigQuery データセットに処理します。

確認

実行が完了したら、Data Agent Kit カタログを確認して、テーブルの作成を確認します。

カタログ エクスプローラで Raw テーブルを確認する

  1. IDE のアクティビティ バーで、[Google Cloud Data Agent Kit] パネルを開きます。
  2. [カタログ] セクションを開きます。
  3. プロジェクト ID を開きます。
  4. [BigQuery] を開きます。
  5. transactions_dataset_evals データセットを開きます。
  6. raw_transactions テーブルをクリックして、メイン エディタで詳細ビューを開きます。
  7. 左側のナビゲーションで、[データ]、[スキーマ]、[詳細] の各タブを確認して、取り込まれたレコードとメタデータを調べます。

セクションのまとめ: エージェント チャットで自然言語を使用して、完全な Spark Serverless ワークロードを生成しました。次に、それを実行して、非構造化 JSON ログを BigQuery(未加工)テーブルに処理しました。

4. dbt を使用して重複除去と正規化を行う

ML モデルをトレーニングする前に、重複するストリーミング ログを削除し、無効なレコード(空のトランザクション ID など)を分離し、ディメンション データ(支払い者と受取人)を結合して、データ品質を強化します。このプロセスには、べき等で信頼性の高い SQL 変換が必要であるため、dbt(データビルド ツール)が最適です。

dbt パイプラインをスキャフォールディングする

エージェントを使用して、BigQuery データセット上に dbt プロジェクトを生成します。

  1. [エージェント チャット] ペインに戻ります。
  2. 次の手順で dbt プロジェクトを生成します。
Scaffold a us-central1 dbt project in dbt_project/ that maps raw_transactions
through to an enriched_transactions model in dataset transactions_dataset_evals.

Deduplicate by transaction_id in staging. Quarantine null IDs to an invalid_transactions model.
Join the valid staging records with dim_payers and dim_payees for the enriched_transactions model,
preserving the historical `is_fraud` label column, and finally add a transaction uniqueness test.

Create an implementation plan first.
  1. エージェントは、メイン エディタ ペインに実装計画アーティファクトを表示します。提案されたファイル構造と SQL ロジックを確認します。
  2. [続行](次に [すべて承認])をクリックして、エージェントがワークスペースでファイルを生成できるようにします。

[Proceed] ボタン付きの実装計画

  1. 生成が完了すると、新しいコンポーネントの概要を示すチュートリアルが表示されます。プロンプトが表示されたら、すべての変更を承認します。

Chat ペインで生成されたすべてのファイルを受け入れる

ビルドとテスト

エージェントは、生成された SQL の構文が有効であることを確認するために dbt compile を自動的に実行しましたが、これらのビューとテーブルを BigQuery にマテリアライズし、ローカル検証用のデータ品質テストを実行します。(注: このラボの後半で、この dbt ステップをエンドツーエンドの Airflow DAG の一部として自動化します)。

  1. 左端のアクティビティ バーで、[エクスプローラ] アイコンをクリックします(または Cmd/Ctrl+Shift+E キーを押します)。
  2. dbt_project -> models を展開して、生成された SQL モデルを調べます。enriched_transactions.sql をクリックして、エディタで変換と不正行為の特徴のロジックを開いて確認します。
  3. ファイル エクスプローラで、dbt_project フォルダを右クリックし、[Open in Integrated Terminal] を選択します。これにより、必要な dbt_project 作業ディレクトリに直接設定されたターミナル ペインが自動的に開きます。
  4. dbt がまだインストールされていない場合は、dbt_project/ の外部(ホームまたはワークスペースのルート)に仮想環境を作成し、BigQuery アダプタをインストールします。
python3 -m venv ~/.venv/dbt
source ~/.venv/dbt/bin/activate
pip install dbt-bigquery
  1. dbt モデルとそれに関連付けられたデータ品質テストを実行します。
dbt build
  1. ターミナル出力を確認します。dbt は SQL をコンパイルし、BigQuery でステージング テーブルとエンリッチ テーブルを具体化して、データテストを実行します。

統合ターミナルで dbt プロジェクトをビルドしてテストする

  1. ビルドが完了したら、ターミナル ペインを閉じて、残りの手順のために画面スペースを解放します。

セクションのまとめ: エージェントを使用して dbt プロジェクトを生成し、データ品質テストを実行して、未加工のレコードをステージング テーブルとエンリッチ テーブルの BigQuery テーブルに変換しました。

5. ランダム フォレストを使用して分散型不正検出モデルをトレーニングする

BigQuery でマテリアライズされた拡張トランザクションを使用して、不正なイベントを分類する ML モデルを構築します。ランダム フォレストは、表形式の分類データに適したアンサンブル学習手法です。Spark Serverless で RandomForestClassifier を実行すると、インフラストラクチャを管理することなく、モデル トレーニングがワーカーノードに分散されます。

このステップでは、エージェントを使用して Spark ML トレーニング パイプラインを生成します。

ML トレーニング ノートブックを生成する

  1. [Agent Chat] ペインを開きます。
  2. 次のプロンプトを指定して、モデルのトレーニング シーケンスを設計します(${PROJECT_ID} はアクティブなプロジェクト ID に置き換えてください)。
Create a PySpark notebook (02_training.ipynb) to train a distributed Random Forest
(RandomForestClassifier) model on the BigQuery table
`transactions_dataset_evals`.`enriched_transactions`, predicting the `is_fraud` label.
Train only on historically labeled records where `is_fraud` is not null.

One-hot encode categorical strings, scale amounts, cache the dataset in memory,
evaluate AUC, and save the evaluated model to gs://${PROJECT_ID}-models/fraud_model.
  1. エージェントのプランまたは生成されたコードを確認し、[続行] / [すべて承認] をクリックして、notebooks/02_training.ipynb をワークスペースに保存します。

トレーニング ノートブックを生成するエージェント

ノートブックを確認して実行する

  1. エディタで notebooks/02_training.ipynb を開きます。
  2. 特徴量エンコード、ベクトル アセンブリ、ランダム フォレスト分類ロジックの PySpark ML パイプライン ステージを確認します。
  3. IDE のノートブック ツールバーで [すべてを実行] をクリックします。
  4. [Select Kernel] プルダウン ピッカーが開いたら、[fraud-pipeline-runtime on Serverless Spark] を選択します。

トレーニング ノートブックの Serverless Spark カーネルを選択する

確認

実行が完了したら、モデルが正しくトレーニングされ、エクスポートされたことを確認します。

  1. ノートブックの下部にある評価セルの出力を確認して、報告された ROC 曲線の下の面積(AUC)スコアを確認します。
  2. モデル アーティファクトが GCS に正常に保存されたことを確認するには、Data Agent Kit のサイドバーで [ストレージ] エクスプローラ ペインを開きます。
  3. -models で終わるバケット(アクティブなプロジェクト ID に関連付けられている)を見つけて展開し、ドリルダウンして fraud_model ディレクトリとそのパイプライン ステージが存在することを確認します。

GCS に保存されたモデルを確認する

セクションのまとめ: エージェントを使用して PySpark ML トレーニング パイプラインを作成し、拡充された BigQuery テーブルでランダム フォレスト モデルをトレーニングして、モデルを Cloud Storage にエクスポートしました。

6. バッチ推論と Cloud Spanner 書き込み

Cloud Storage に保存されているトレーニング済みの予測モデルを使用して、BigQuery を通過する新しいトランザクションに対してバッチ推論を実行します。高リスクの取引は、コンプライアンス チームが審査できるように、運用システムに転送する必要があります。Cloud Spanner は、この審査キューにスケーラブルなトランザクション データベースを提供します。

バッチ推論ノートブックを生成する

エージェントを使用して、BigQuery、Cloud Storage、Cloud Spanner を接続する推論ノートブックを作成します。

  1. [Agent Chat] ペインを開きます。
  2. 次のプロンプトを入力します。
Create an inference notebook (03_inference.ipynb) that loads the RandomForestClassifier
model to score unlabeled records (where `is_fraud` is null) from the BigQuery table
`transactions_dataset_evals`.`enriched_transactions`.

Filter for high-risk transactions with a 50%+ fraud probability score (probability >= 0.50)
and write them to the Cloud Spanner table SparkEvalFraudReviewQueue in the cymbal-fraud instance
under fraud-db.
  1. 生成されたノートブックを承認して、notebooks/03_inference.ipynb をワークスペースに保存します。

推論ノートブックを生成するエージェント

ノートブックを確認して実行する

  1. 新しく生成された notebooks/03_inference.ipynb をエディタで開きます。
  2. PySpark 推論シーケンスを確認します。
    • 依存関係: Serverless ランタイム テンプレートは、Spark 実行に必要な cloud-spanner JAR 依存関係を提供します。
    • データ形式: スクリプトは、Spanner テーブル スキーマに合わせて書き込む前に、複雑な Spark ML ベクトル列(未加工の特徴や確率など)を削除します。
    • Spanner コネクタ: .format("cloud-spanner") を使用してフラグ付きの行を書き込み、審査キューに直接追加します。
  3. IDE のノートブック ツールバーで [すべてを実行] をクリックします。
  4. カーネルの選択を求められたら、[fraud-pipeline-runtime on Serverless Spark] を選択します。

確認

推論ノートブックの処理が完了したら、IDE 内で運用 Spanner データベースに直接クエリを実行できます。

  1. IDE のアクティビティ バーで、[Google Cloud Data Agent Kit] パネルを開きます。
  2. [カタログ] セクションを開きます。
  3. プロジェクト ID を開き、[Spanner] を開きます。
  4. cymbal-fraud -> fraud-db -> Tables -> SparkEvalFraudReviewQueue に移動します。
  5. テーブルを右クリックして [テーブルのクエリ] を選択し、クエリを実行します。
SELECT *
FROM `SparkEvalFraudReviewQueue`
LIMIT 100;
  1. 下の [クエリ結果] ペインに、手動レビューの対象としてフラグが設定されたリスクの高いトランザクションを表す新しい行が挿入されていることがわかります。

Cloud Spanner の行を確認する

セクションのまとめ: エージェントを使用してバッチ推論ノートブックを作成し、トレーニング済みモデルを使用してラベルのない BigQuery レコードをスコアリングし、リスクの高いトランザクションを Cloud Spanner に直接書き込みました。

7. Managed Airflow でスキャフォールディングとオーケストレーションを行う

現在のパイプラインは、取り込みノートブック、dbt 変換プロジェクト、バッチ推論ノートブックという個別のステップで構成されています。これをプロダクション レディにするには、スケジュールされた依存関係グラフに統合します。

Managed Service for Apache Airflow(旧称 Cloud Composer)は、このワークフローのマネージド オーケストレーション エンジンを提供します。Data Agent Kit には、宣言型 YAML パイプライン定義を Airflow DAG に直接変換するOrchestration Pipelines機能が含まれています。

パイプラインの定義

エージェントを使用して、オーケストレーション パイプライン構成を生成します。

  1. [エージェント チャット] で、次のプロンプトを入力します(${PROJECT_ID} の置き換えを忘れないでください)。
Initialize and define an orchestration pipeline (fraud_analysis_pipeline) triggering
the ingestion notebook, dbt project, and inference notebook in sequential order.
For Dataproc Serverless engine configs in us-central1, use resourceProfile.inline
(defining properties with spark.jars: "gs://spark-lib/spanner/spark-3.5-spanner-1.4.0.jar"
for inference) rather than resourceProfile.path or overrides.

Set the schedule interval to run daily at midnight, and use
gs://${PROJECT_ID}-airflow-artifacts for artifact storage.

DAG 構成を確認する

Data Agent Kit Orchestrator は、宣言型 YAML 構成を使用してパイプラインを定義し、Apache Airflow にデプロイします。これにより、定義をバージョン管理し、CI/CD を介してデプロイできます。

IDE エクスプローラ ペインで、エージェントがワークスペースのルートに生成した 2 つのパイプライン ファイルを確認します。

  1. deployment.yaml: このファイルを開きます。これは環境レジストリとして機能します。論理 dev パイプラインを cymbal-airflow 環境にマッピングし、実行リージョン(us-central1)を設定して、コンパイルされた DAG と依存関係がステージングされる artifact_storage バケットを定義します。
  2. fraud_analysis_pipeline.yaml: このファイルを開きます。これにより、実行グラフが定義されます。トリガー スケジュール(interval: '0 0 * * *')を指定し、actions ブロックで 3 つのステップを順序付けます。
    • Dataproc Serverless で実行されている 01_ingestion.ipynb の取り込み notebook アクション。
    • dbt_project ディレクトリをターゲットとする変換 pipeline アクション。dependsOn 依存関係は取り込みステップを指しています。
    • dbt ステップを指す dependsOn 依存関係を持つ 03_inference.ipynb の推論 notebook アクション。Spanner JAR プロパティをバンドルします。
  3. また、エージェントは、生成されたこれらのアーティファクトをエディタ ペインの [Walkthrough] タブにまとめ、実行された構成と検証の概要を示します。

インタラクティブ DAG 構成

Data Agent Kit は、パイプライン構成をインタラクティブなビジュアル グラフとしてレンダリングし、Airflow DAG プロパティの検査と編集を行います。

  1. IDE のアクティビティ バーで、[Google Cloud Data Agent Kit] パネルを開きます。
  2. [DATA ENGINEERING] で [Orchestration Pipelines] を開きます。
  3. fraud_analysis_pipeline.yaml をクリックして、メイン エディタでビジュアル DAG キャンバスを開きます。

オーケストレーション DAG のビジュアル キャンバス

  1. 上部の Schedule trigger ノードをクリックします。右側に構成フライアウトが開き、解析された Cron 文字列(0 0 * * *)が表示されます。また、バックフィルやキャッチアップなどのパラメータを調整することもできます。
  2. ノートブック タスクノード(取り込みステップや推論ステップなど)のいずれかをクリックします。フライアウトが更新され、特定の Dataproc Serverless 実行マッピングとコネクタ プロパティが表示されます。
  3. ノードブロック内のノートブック ファイル名のハイパーリンク(01_ingestion.ipynb など)に注目してください。クリックすると、ノートブックがエディタで直接開きます。
  4. 左側のサイドバーの [Orchestration Pipelines] の下にある Deployment configuration をクリックします。このビューには、ターゲットの dev 環境クラスタと出力 GCS バケット アーティファクトが表示されます。

セクションのまとめ: エージェントを使用してオーケストレーション パイプライン構成を生成し、インタラクティブなビジュアル キャンバスで取り込み、dbt、推論タスク間の依存関係を定義しました。

8. デプロイ、実行、モニタリング

ローカルで定義した DAG を使用して、設定時にプロビジョニングされた Managed Airflow 環境に接続し、パイプラインをデプロイします。

Managed Service for Apache Airflow を構成する

デプロイする前に、Data Agent Kit の設定で Scheduler 接続を構成して、拡張機能が Managed Airflow 環境をターゲットにするようにします。

  1. IDE のアクティビティ バーで、[Google Cloud Data Agent Kit] パネルを開きます。
  2. SETTINGS で、[設定] をクリックします。
  3. 左側のメニューから [スケジューラ] を選択します。
  4. 設定を構成します。
    • プロジェクト ID: アクティブなプロジェクト ID を選択します。
    • リージョン: us-central1 を選択します。
    • 環境: cymbal-airflow を選択します。
  5. [保存] をクリックします。

Managed Service for Apache Airflow の設定

DAG をデプロイする

構成したパイプラインを、ビジュアル キャンバスから Managed Airflow 環境に直接デプロイします。

  1. [Google Cloud Data Agent Kit] サイドバーで、DATA ENGINEERING > Orchestration Pipelines を開いて fraud_analysis_pipeline.yaml をクリックし、ビジュアル DAG キャンバスを開きます。
  2. キャンバスのツールバーの右上にある青色の [パイプラインを実行] ボタンをクリックします。
  3. 環境プルダウン選択ツールで、dev を選択します。
  4. 下部のステータス領域(Running pipeline: Building pipeline locally...)で進行状況の通知を確認します。拡張機能により、DAG が自動的にコンパイルされ、ノートブックと dbt アセットがパッケージ化されて、Managed Airflow 環境の GCS バケットにアップロードされます(完了まで約 3 ~ 4 分かかります)。

ビジュアル キャンバスからパイプラインをデプロイする

実行をモニタリングする

ローカル コンパイルが完了し、ポップアップ通知で Triggered a new run for pipeline... successfully が確認されたら、ライブ実行をモニタリングします。

  1. Google Cloud Data Agent Kit のサイドバーで、DATA ENGINEERING > Orchestration Pipelines を開きます。
  2. [パイプライン管理] をクリックします。
  3. [パイプライン管理] テーブルで、fraud_analysis_pipeline をクリックして実行履歴を開きます。

パイプライン管理の概要

  1. [実行履歴] ビューで、カレンダーからアクティブな実行を選択します。
  2. 各パイプライン タスク(取り込み、dbt 変換、推論)の実行が進むにつれて、ステータス インジケーターが更新され、タスクの所要時間が入力されます。タスクをクリックすると、ライブ実行の出力と Airflow DAG ログを確認できます。

パイプラインの実行履歴とタスクの詳細

セクションのまとめ: Airflow スケジューラの接続を構成し、エンドツーエンドの分析パイプラインを Managed Airflow にデプロイして、ライブ実行をモニタリングし、未加工のログから最終的な Cloud Spanner 予測までシステムを検証しました。

9. クリーンアップ

この Codelab で使用したリソースに対して Google Cloud プロジェクトで継続的に課金されないようにするには、自動スクリプトを使用して環境を破棄します。

  1. [ターミナル] パネル(または Cloud Shell)で、scripts ディレクトリに移動して実行します。
cd ~/devrel-demos/codelabs/agentic-data-labs/data-science/scripts
chmod +x teardown.sh
./teardown.sh
  1. スクリプトは、削除する予定のすべてのリソースを一覧表示し、確認を求めます。
    • Managed Airflow 環境(cymbal-airflow)
    • Cloud Spanner インスタンス(cymbal-fraud)
    • BigQuery データセット(transactions_dataset_evals)
    • Cloud Storage バケット(gs://${PROJECT_ID}-fin-clearing-raw と gs://${PROJECT_ID}-models)
    • ワーカー サービス アカウント(composer-worker-sa)
  2. 確定するには、「y」と入力します。破棄スクリプトは、プロビジョニングされたすべての GCP サービスを削除し、ローカルファイルをクリーンアップします。

10. 完了

Cloud Storage、BigQuery、Managed Service for Apache Spark(Spark Serverless)、dbt、Cloud Spanner、Managed Service for Apache Airflow にまたがるエンドツーエンドの不正行為検出パイプラインを構築し、Antigravity IDE 内の Google Cloud Data Agent Kit を使用してペア プログラミングを行いました。

学習した内容

  1. 📥 Managed Service for Apache Spark と Data Agent Kit を使用して、取り込んだ未加工のトランザクション ログを BigQuery テーブルに格納しました。
  2. 🧹 データ品質テストを含む dbt プロジェクトを作成して、データを重複除去して正規化します。
  3. 🤖 RandomForestClassifier を使用して分散ランダム フォレスト モデルをトレーニングし、トレーニング済みモデルを Cloud Storage にエクスポートしました。
  4. ⚡ 受信したトランザクションに対してバッチ推論を実行し、リスクの高いレコードを監査レビューのために Cloud Spanner に転送しました。
  5. 🔄 Managed Service for Apache Airflow と IDE の DAG 管理ツールを使用して、スケジュールされた Airflow DAG としてワークフローをオーケストレート、デプロイ、モニタリングしました。

主なコンセプト

コンセプト

学習した内容

Data Agent Kit

自然言語を使用して IDE 内でペア プログラミングを行い、PySpark ノートブックの生成、dbt モデルの構成、Airflow DAG の定義を行う

BigQuery

分析 SQL、dbt 変換、ML トレーニング用のスケーラブルな表形式ストレージ

Spark Serverless

分散 PySpark データ読み込みとランダム フォレスト ML トレーニングのサーバーレス実行

Cloud Spanner コネクタ

バッチ Spark 推論予測をオペレーショナル データベースのレビューキューに直接書き込む

YAML DAG 宣言

宣言型パイプライン定義が、IDE でインタラクティブな Airflow ビジュアル グラフとしてレンダリングされる

DAG の視覚的な管理

パイプラインの依存関係の検査、Managed Airflow へのデプロイ、IDE 内でのライブタスクの実行履歴のモニタリング

次のステップ