使用 Data Agent Kit 和 Antigravity IDE 建構詐欺偵測管道

1. 簡介

假設您是 Cymbal Financial 的資料科學家,這家公司是交易量極高的付款處理方。近期發生一連串結算延遲事件,法規遵循團隊懷疑是有人聯手詐欺。您需要建構管道,擷取原始結算所交易記錄、清理資料、訓練機器學習模型、執行批次推論,並將高風險交易匯入 Cloud Spanner 審查佇列,以進行人工稽核。

一般來說,這需要花費數天時間編寫重複的設定程式碼 (Spark 筆記本、dbt 設定、訓練指令碼、Airflow DAG),並在控制台介面和編輯器之間不斷切換。

在本程式碼實驗室中,您將在 Antigravity IDE 中,使用 Google Cloud Data Agent Kit (DAK) 與代理程式進行配對程式設計。服務專員會使用對話式自然語言,協助您產生 Spark 筆記本、編譯 dbt 專案、建構推論迴圈,以及使用 Managed Service for Apache Airflow 自動調度管理工作流程。

學習內容

軟硬體需求

  • 網路瀏覽器,例如 Chrome
  • 已啟用計費功能的 Google Cloud 雲端專案 (建議使用新的專屬專案進行實作實驗室)。
  • 對 SQL、Python 和 PySpark 有基本的瞭解。
  • Antigravity IDE (建議使用 Google AI Pro 訂閱)

本程式碼研究室建立的資源費用應不到 $5 美元。請務必按照實驗室結尾的「清除」指示,刪除已佈建的資源。

2. 環境設定

如要啟動實驗室,請執行啟動程序指令碼。這個指令碼會自動啟用必要的 GCP API、建立擷取 Cloud Storage bucket、產生模擬交易和目錄資料集、將參照目錄載入 BigQuery,並啟動 Cloud Spanner 和 Managed Service for Apache Airflow (原名為 Cloud Composer) 的背景佈建作業。

選取或建立專案

在 Google Cloud 控制台中選擇現有專案,或建立新專案。

驗證帳單

請確認 Google Cloud 專案已啟用計費功能。如要進一步瞭解如何操作,請參閱本指南。

執行設定指令碼

您將使用 Google Cloud Shell (或已設定 Google Cloud CLI 的本機殼層) 啟動環境設定。

  1. 開啟 Google Cloud 控制台。
  2. 點選右上工具列中的「啟用 Cloud Shell」。

開啟 Cloud Shell

  1. 在 Cloud Shell 終端機中設定現有專案:
gcloud config set project <<YOUR_PROJECT_ID>>
export PROJECT_ID=$(gcloud config get-value project)
  1. 複製程式碼研究室存放區,然後前往指令碼資料夾:
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),然後選擇「Open Folder」(開啟資料夾),在 IDE 中開啟該資料夾。這會做為本程式碼研究室的本機工作區。

設定 Antigravity IDE 專案資料夾

安裝 Data Agent Kit 擴充功能

Google Cloud Data Agent Kit 擴充功能可直接在編輯器中與 Google Cloud 資料服務深度整合,讓您與 BigQuery、Cloud SQL、Cloud Storage 等服務互動,不必切換環境。

  1. 在 Antigravity IDE 中,點選畫面最左側活動列的「擴充功能」圖示 (看起來像四個正方形)。
  2. 在「擴充功能」窗格頂端的搜尋列中,輸入 Google Cloud Data Agent Kit。
  3. 找出由 googlecloudtools 發布的「Google Cloud Data Agent Kit」擴充功能。
  4. 按一下「Install」按鈕。
  5. 系統可能會顯示提示,詢問「您是否信任發布者『googlecloudtools』及其擴充功能?」。按一下「信任發布者並安裝」繼續操作。

安裝 Data Agent Kit 擴充功能

安裝完成後,Antigravity IDE 最左側的活動列會顯示新的「Google Cloud Data Agent Kit」圖示。

  1. 系統應會自動開啟名為「Welcome to Google Cloud Data Agent Kit」(歡迎使用 Google Cloud 資料代理套件) 的入門頁面。如果未登入 Cloud 帳戶,請按照提示允許存取。
  2. 在「設定摘要」部分,找到專案欄位。按一下下拉式選單,然後選取 Google Cloud 專案。將區域設為 us-central1。然後選取「設定 MCP 伺服器」。

Data Agent Kit 擴充功能的初始設定

  1. 選取「設定 MCP 伺服器」。在「MCP Configuration」窗格下方,請務必啟用下列遠端 MCP 伺服器:
    • BigQuery
    • Spanner
    • 筆記本

然後按一下「開始使用」。

設定 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 伺服器:查看已啟用的 MCP 伺服器 (BigQuery、Notebooks、Spanner 等),讓 AI 代理安全地與您的資料互動。
    • 技能:探索預先建構的技能,為代理程式提供專門功能,處理複雜的資料工作。

Data Agent Kit 設定面板

章節回顧:您執行了啟動程序指令碼來建立 GCS 和 BigQuery 資產,而 Spanner 和 Airflow 則在背景中建構。接著,您在 Antigravity IDE 中開啟專案,並啟用了 Google Cloud Data Agent Kit 擴充功能。現在可以開始撰寫第一個筆記本了。

3. 使用 Spark Serverless 擷取原始記錄檔

在本節中,您會將原始 JSON 交易記錄擷取至資料湖泊。Managed Service for Apache Spark (Spark 無伺服器) 可直接連線至 BigQuery 的原生儲存空間。您將使用標準 BigQuery 連接器管理表格資料,並啟用直接查詢和分析功能。

探索預先設定的 Spark Serverless 執行階段

執行 Spark 程式碼前,請檢查設定指令碼預先設定的無伺服器執行階段範本。這個範本會定義目標執行環境後端,並將必要的連接器依附元件組合在一起。

  1. 在 IDE 活動列中,開啟「Google Cloud Data Agent Kit」面板。
  2. 展開「Apache Spark」下拉式選單,然後展開「無伺服器」。
  3. 在 fraud-pipeline-runtime 上按一下滑鼠右鍵,然後選取「設定檔」,即可在編輯器中開啟設定檢視畫面。
  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 連接器,因此不需要額外的 JAR 設定,即可讀取及寫入 BigQuery 資料表)。

探索 Spark 無伺服器執行階段屬性

  1. 請注意左側的「互動式工作階段」分頁標籤。目前為空白,因為您尚未執行任何程式碼。在下一個步驟中執行筆記本後,系統會動態佈建即時無伺服器運算工作階段,並顯示在這裡!

使用 Data Agent Kit 擷取資料

您將使用 Data Agent Kit 與代理程式進行配對程式設計,不必手動設定 Spark 工作階段,也不必從頭編寫 PySpark 載入指令碼。

  1. 按一下右上工具列的「切換代理程式」圖示,開啟「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. 當 AI 代理程式生成檔案後,按一下對話窗格底部的藍色「接受全部」按鈕 (或勾號圖示),即可將 notebooks/01_ingestion.ipynb 儲存至工作區。

代理程式生成擷取筆記本

查看並執行筆記本

  1. 在 IDE 中開啟新產生的 notebooks/01_ingestion.ipynb。
  2. 查看 BigQuery 連接器寫入邏輯的 PySpark 程式碼。
  3. 按一下 IDE 筆記本工具列中的「Run All」(全部執行)。
  4. 如果是第一次執行遠端 Spark 筆記本,IDE 可能會提示您安裝本機依附元件。如果出現提示,請按一下「Install dependencies for Remote Spark Kernels」,確認安裝對話方塊,然後再次點選「Run All」。
  5. 在「Select Kernel」(選取核心) 下拉式選單中,依序選擇「Remote Spark Kernels」(遠端 Spark 核心) ->「fraud-pipeline-runtime on Serverless Spark」(無伺服器 Spark 上的 fraud-pipeline-runtime)。(提示:如果沒有看到預先設定的執行階段範本,請按一下核心挑選器下拉式選單右上方的重新整理圖示,重新載入可用的遠端核心)。
  6. 查看編輯器左下方的狀態列。畫面會顯示 Connecting to kernel: fraud-pipeline-runtime on Serverless Spark...。由於這是 Spark 無伺服器執行階段核心後端的首次啟動,因此需要幾分鐘才能完成佈建和啟動。
  7. 核心完成連線後,筆記本會自動依序執行所有儲存格,將原始交易記錄處理至 BigQuery 資料集。

驗證

執行完畢後,請檢查 Data Agent Kit 目錄,確認資料表已建立:

在目錄探索工具中驗證原始資料表

  1. 在 IDE 活動列中,開啟「Google Cloud Data Agent Kit」面板。
  2. 展開「目錄」部分。
  3. 展開專案 ID。
  4. 展開「BigQuery」BigQuery。
  5. 展開 transactions_dataset_evals 資料集。
  6. 按一下 raw_transactions 資料表,即可在主要編輯器中開啟詳細資料檢視畫面。
  7. 在左側導覽面板中,瀏覽「資料」、「結構定義」和「詳細資料」分頁,檢查擷取的記錄和中繼資料。

本節重點:您在 Agent Chat 中使用自然語言,產生完整的 Spark 無伺服器工作負載。接著執行該管道,將非結構化 JSON 記錄檔處理成 BigQuery (原始) 資料表。

4. 使用 dbt 簡化及正規化資料

訓練機器學習模型前,您會先移除重複的串流記錄、隔離不良記錄 (例如空白交易 ID),並加入維度資料 (付款人和收款人),以確保資料品質。這個程序需要具備等冪性且可靠的 SQL 轉換,因此 dbt (資料建構工具) 非常適合。

架構 dbt 管道

使用代理程式,根據 BigQuery 資料集產生 dbt 專案:

  1. 返回「Agent Chat」(代理程式對話) 窗格。
  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. 按一下「繼續」 (然後按一下「全部接受」),允許代理程式在工作區中生成檔案。

顯示「繼續」按鈕的實作計畫

  1. 生成完成後,代理程式會顯示「導覽」,總結新元件。如果系統提示,請接受所有變更。

接受 Chat 窗格中所有生成的檔案

建構與測試

雖然代理程式已自動執行 dbt compile,確保產生的 SQL 語法有效,但您現在要將這些檢視區塊和資料表具體化到 BigQuery 中,並執行資料品質測試以進行本機驗證。(注意:在本實驗室的後續步驟中,您會自動執行這個 dbt 步驟,做為端對端 Airflow DAG 的一部分)。

  1. 在最左側的活動列中,按一下「Explorer」圖示 (或按下 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 中具體化經過擴增的交易後,您將建構機器學習模型,對詐欺事件進行分類。隨機森林是一種適合表格型分類資料的集成學習方法。在 Spark Serverless 上執行 RandomForestClassifier 時,模型訓練作業會分散到各個工作節點,您不必管理基礎架構。

在本步驟中,您將使用代理生成 Spark ML 訓練 pipeline。

生成機器學習訓練筆記本

  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. 查看 AI 代理程式的企劃書或生成的程式碼,然後按一下「繼續」/「全部接受」,將 notebooks/02_training.ipynb 儲存至工作區。

代理程式生成訓練筆記本

查看並執行筆記本

  1. 在編輯器中開啟 notebooks/02_training.ipynb。
  2. 查看 PySpark 機器學習管道階段,瞭解特徵編碼、向量組裝和隨機森林分類邏輯。
  3. 按一下 IDE 筆記本工具列中的「Run All」(全部執行)。
  4. 「Select Kernel」(選取核心) 下拉式選單開啟後,選取「fraud-pipeline-runtime on Serverless Spark」(無伺服器 Spark 上的詐欺管道執行階段)。

為訓練筆記本選取無伺服器 Spark 核心

驗證

執行完畢後,請確認模型已正確訓練及匯出:

  1. 查看筆記本底部的評估儲存格輸出內容,確認回報的 ROC 曲線下面積 (AUC) 分數。
  2. 如要確保模型構件已成功儲存至 GCS,請展開 Data Agent Kit 側欄中的「STORAGE」探索器窗格。
  3. 找到以 -models 結尾的 bucket (與有效專案 ID 相關聯),展開該 bucket,然後向下鑽研,確認 fraud_model 目錄及其管道階段是否存在。

確認模型已儲存在 GCS 中

本節重點回顧:您使用代理程式建立 PySpark ML 訓練管道、在經過擴充的 BigQuery 資料表上訓練 Random Forest 模型,並將模型匯出至 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 Runtime 範本提供 Spark 執行作業所需的 cloud-spanner JAR 依附元件。
    • 資料格式:指令碼會捨棄複雜的 Spark ML 向量資料欄 (例如原始特徵和機率),然後再寫入資料,以符合 Spanner 資料表結構定義。
    • Spanner 連接器:使用 .format("cloud-spanner") 寫入標記的資料列,直接附加至審查佇列。
  3. 按一下 IDE 筆記本工具列中的「Run All」(全部執行)。
  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 包含 Orchestration Pipelines 功能,可將宣告式 YAML 管道定義直接轉換為 Airflow DAG。

定義管道

使用代理程式產生自動調度管理管道設定:

  1. 在 Agent Chat 中提供下列提示 (請記得將 ${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 探索工具窗格中,查看代理程式在工作區根目錄中產生的兩個管道檔案:

  1. deployment.yaml:開啟這個檔案。這會做為環境登錄。這會將邏輯 dev 管道對應至 cymbal-airflow 環境、設定執行區域 (us-central1),並定義編譯後的 DAG 和依附元件的暫存 artifact_storage 值區。
  2. fraud_analysis_pipeline.yaml:開啟這個檔案。這會定義執行圖。這會指定觸發排程 (interval: '0 0 * * *'),並在 actions 區塊下依序執行三個步驟:
    • 在 Dataproc Serverless 上執行的擷取 notebook 動作 01_ingestion.ipynb。
    • 以 dbt_project 目錄為目標的轉換 pipeline 動作,以及指向擷取步驟的 dependsOn 依附元件。
    • 03_inference.ipynb 的推論 notebook 動作,具有指向 dbt 步驟的 dependsOn 依附元件,並將 Spanner JAR 屬性組合在一起。
  3. 代理程式也會將這些產生的構件歸納到編輯器窗格的「逐步解說」分頁中,列出執行的設定和驗證。

互動式 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。
    • 「Region」(地區):選取 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 bucket (完成這項作業約需 3 到 4 分鐘)。

從視覺畫布部署管道

監控執行作業

完成本機編譯後,彈出式通知會確認 Triggered a new run for pipeline... successfully,請監控即時執行情況:

  1. 在「Google Cloud Data Agent Kit」側欄中,依序展開 DATA ENGINEERING > Orchestration Pipelines。
  2. 按一下「Pipelines management」(管道管理)。
  3. 在「Pipelines Management」(管道管理) 表格中,按一下 fraud_analysis_pipeline 開啟執行記錄。

管道管理總覽

  1. 在「執行記錄」檢視畫面中,從日曆選取有效執行作業。
  2. 隨著執行作業在各個管道工作 (擷取、dbt 轉換和推論) 中進行,狀態指標會更新,並填入工作持續時間。按一下任一工作,即可檢查即時執行輸出內容和 Airflow DAG 記錄。

管道執行記錄和工作詳細資料

本節重點:您已設定 Airflow 排程器連線、將端對端分析管道部署至 Managed Airflow,並監控即時執行作業,從原始記錄到最終的 Cloud Spanner 預測結果,驗證整個系統。

9. 清理

如要避免系統持續向您的 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 bucket (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 轉換和機器學習訓練

Spark 無伺服器

以無伺服器方式執行分散式 PySpark 資料載入和隨機森林機器學習訓練

Cloud Spanner 連接器

將批次 Spark 推論預測結果直接寫入營運資料庫審查佇列

YAML DAG 宣告

在 IDE 中,以互動式 Airflow 視覺化圖表呈現宣告式管道定義

以視覺化方式管理 DAG

在 IDE 中檢查管道依附元件、部署至 Managed Airflow,以及監控即時工作執行記錄

後續步驟