Pipeline di rilevamento delle frodi con Data Agent Kit e Antigravity IDE

1. Introduzione

Immagina di essere un data scientist presso Cymbal Financial, un elaboratore dei pagamenti ad alto volume. Si è verificata un'ondata di ritardi nei pagamenti e il team di conformità sospetta una frode coordinata. Devi creare una pipeline per importare i log delle transazioni della camera di compensazione non elaborate, pulire i dati, addestrare un modello di machine learning, eseguire l'inferenza batch e inserire le transazioni ad alto rischio in una coda di revisione Cloud Spanner per il controllo manuale.

Normalmente, ciò richiede giorni di scrittura di codice di configurazione ripetitivo (notebook Spark, configurazioni dbt, script di addestramento, DAG di Airflow) e un cambio di contesto costante tra interfacce della console ed editor.

In questo codelab, programmerai in coppia con un agente utilizzando Google Cloud Data Agent Kit (DAK) all'interno dell'IDE Antigravity. Utilizzando il linguaggio naturale conversazionale, l'agente ti aiuterà a generare notebook Spark, compilare un progetto dbt, costruire un ciclo di inferenza e orchestrare il flusso di lavoro utilizzando Managed Service for Apache Airflow.

In questo lab proverai a:

  • Importa i log della clearing house da Cloud Storage utilizzando Managed Service for Apache Spark (Spark serverless) in una tabella BigQuery.
  • Rimuovi i duplicati e normalizza le transazioni utilizzando dbt per stabilire livelli di dati puliti (raw, staging, arricchiti).
  • Addestra un modello di classificazione Random Forest distribuito (RandomForestClassifier) su Spark Serverless.
  • Esegui l'inferenza batch sulle nuove transazioni e scrivi avvisi ad alto rischio direttamente in Cloud Spanner.
  • Orchestra, configura visivamente e implementa l'intera pipeline utilizzando Managed Service for Apache Airflow e il monitoraggio interattivo dei DAG all'interno dell'IDE.

Che cosa ti serve

  • Un browser web come Chrome
  • Un progetto Google Cloud con la fatturazione abilitata (ti consigliamo di utilizzare un nuovo progetto dedicato per i laboratori pratici).
  • Conoscenza di base di SQL, Python e PySpark.
  • Antigravity IDE con un abbonamento a Google AI Pro (consigliato)

Le risorse create in questo codelab dovrebbero costare meno di 5 $. Assicurati di seguire le istruzioni per la pulizia alla fine del lab per eliminare le risorse di cui è stato eseguito il provisioning.

2. Configurazione dell'ambiente

Per iniziare il lab, esegui uno script di bootstrap. Questo script attiva automaticamente le API GCP richieste, crea un bucket Cloud Storage per l'importazione, genera set di dati di transazioni e directory simulati, carica le directory di riferimento in BigQuery e avvia il provisioning in background di Cloud Spanner e del servizio gestito per Apache Airflow (precedentemente noto come Cloud Composer).

Seleziona o crea un progetto

Scegli un progetto esistente o creane uno nuovo nella console Google Cloud.

Verifica della fatturazione

Verifica che la fatturazione sia attivata per il tuo progetto Google Cloud. Per scoprire di più su come fare, segui questa guida.

Esegui lo script di configurazione

Utilizzerai Google Cloud Shell (o la shell locale configurata con Google Cloud CLI) per avviare la configurazione dell'ambiente.

  1. Apri la console Google Cloud.
  2. Fai clic su Attiva Cloud Shell nella barra degli strumenti in alto a destra.

Apri Cloud Shell

  1. Nel terminale Cloud Shell, configura il progetto attivo:
gcloud config set project <<YOUR_PROJECT_ID>>
export PROJECT_ID=$(gcloud config get-value project)
  1. Clona il repository del codelab e vai alla cartella degli script:
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. Esegui lo script di configurazione del bootstrap per eseguire il deployment di tutte le risorse in us-central1:
chmod +x setup.sh setup_spanner.sh setup_composer.sh
export REGION=us-central1
./setup.sh
  1. Al termine dello script, verrà visualizzato un output di riepilogo che indica che il set di dati BigQuery e il bucket Cloud Storage sono pronti. In background, continueranno il provisioning di Cloud Spanner (richiede circa 2 minuti) e Managed Airflow (richiede circa 20 minuti). Puoi monitorare i progressi in qualsiasi momento eseguendo:
tail -f /tmp/spanner_setup.log
tail -f /tmp/composer_setup.log

Apri l'IDE Antigravity

  1. Scarica e installa l'IDE Antigravity dalla pagina di download di Google Antigravity.
  2. Avvia l'IDE Antigravity.
  3. Crea una nuova cartella vuota sul computer locale (ad es.denominata agentic-data-labs) e aprila nell'IDE scegliendo Apri cartella. che fungerà da spazio di lavoro locale per il codelab.

Configura la cartella del progetto Antigravity IDE

Installare l'estensione Data Agent Kit

L'estensione Google Cloud Data Agent Kit fornisce una profonda integrazione con i servizi di dati Google Cloud direttamente nell'editor, consentendoti di interagire con BigQuery, Cloud SQL, Cloud Storage e altro ancora senza cambiare contesto.

  1. Nell'IDE Antigravity, fai clic sull'icona Estensioni nella barra delle attività all'estrema sinistra dello schermo (simile a quattro quadrati).
  2. Nella barra di ricerca nella parte superiore del riquadro Estensioni, digita Google Cloud Data Agent Kit.
  3. Individua l'estensione denominata Google Cloud Data Agent Kit pubblicata da googlecloudtools.
  4. Fai clic sul pulsante Installa.
  5. Potrebbe essere visualizzato un messaggio che chiede: "Ritieni attendibile l'editore "googlecloudtools" e le relative estensioni?". Fai clic su Considera attendibili gli editori e installa per continuare.

Installare l'estensione Data Agent Kit

Una volta installato, vedrai una nuova icona Google Cloud Data Agent Kit nella barra delle attività all'estrema sinistra dell'IDE Antigravity.

  1. Dovrebbe aprirsi automaticamente una pagina di onboarding intitolata "Benvenuto in Google Cloud Data Agent Kit". Se non hai eseguito l'accesso al tuo account Cloud, segui le istruzioni per consentire l'accesso.
  2. Nella sezione Riepilogo configurazione, individua il campo del progetto. Fai clic sul menu a discesa e seleziona il tuo progetto Google Cloud. Imposta la tua regione su us-central1. Poi seleziona Configura server MCP.

Configurazione iniziale dell'estensione Data Agent Kit

  1. Seleziona Configura server MCP. Nel riquadro Configurazione MCP, assicurati di abilitare i seguenti server MCP remoti:
    • BigQuery
    • Spanner
    • Notebook

Poi fai clic su Inizia.

Configura i server MCP

Esplora le opzioni di configurazione

Al termine della configurazione, verrà visualizzata la pagina "Inizia a utilizzare Google Cloud Data Agent Kit".

  1. Nella sezione "Configurazione", fai clic su Inizia.
  2. Si aprirà il riquadro Configurazione del Data Agent Kit. Esplora le schede:
    • Progetto e regione:verifica l'ID progetto selezionato e conferma che lo script di configurazione abbia abilitato tutte le API necessarie (Compute Engine, Cloud Storage, BigQuery, Spanner e così via).
    • BigQuery:configura la località predefinita per le query BigQuery. Utilizza la regione us-central1.
    • Configura i server MCP:visualizza i server MCP abilitati (BigQuery, Notebook, Spanner e così via) che consentono agli agenti AI di interagire in modo sicuro con i tuoi dati.
    • Competenze:esplora le competenze predefinite che forniscono agli agenti funzionalità specializzate per attività complesse sui dati.

Riquadro Impostazioni di Data Agent Kit

Riepilogo della sezione:hai eseguito lo script di bootstrap per creare asset GCS e BigQuery mentre Spanner e Airflow vengono creati in background. Poi hai aperto il progetto nell'IDE Antigravity e hai attivato l'estensione Google Cloud Data Agent Kit. Ora puoi scrivere il tuo primo notebook.

3. Importare log non elaborati utilizzando Spark Serverless

In questa sezione, importerai i log delle transazioni JSON non elaborati nel data lake. Managed Service for Apache Spark (Spark serverless) si connette direttamente allo spazio di archiviazione nativo di BigQuery. Utilizzerai il connettore BigQuery standard per gestire i dati tabellari e attivare query e analisi dirette.

Esplora il runtime Spark serverless preconfigurato

Prima di eseguire il codice Spark, esamina il modello Serverless Runtime preconfigurato dallo script di configurazione. Questo modello definisce il backend dell'ambiente di esecuzione di destinazione e raggruppa le dipendenze del connettore necessarie.

  1. Nella barra delle attività dell'IDE, apri il riquadro Google Cloud Data Agent Kit.
  2. Espandi il menu a discesa Apache Spark, poi espandi Serverless.
  3. Fai clic con il tasto destro del mouse su fraud-pipeline-runtime e seleziona Profilo per aprire la relativa visualizzazione della configurazione nell'editor.
  4. Nella scheda Profilo, scorri verso il basso ed espandi Proprietà per esaminare le dipendenze personalizzate associate all'ambiente:
    • spark.jars: contiene gs://spark-lib/spanner/spark-3.5-spanner-1.4.0.jar, che utilizza il connettore Spark Spanner per consentire ai job Spark di scrivere i risultati dell'inferenza direttamente in Cloud Spanner più avanti nel lab. (Nota: Dataproc Serverless include il connettore Spark BigQuery di Google Cloud per impostazione predefinita, senza richiedere alcuna configurazione jar aggiuntiva per leggere e scrivere tabelle BigQuery).

Esplora le proprietà del runtime Spark serverless

  1. Nota la scheda Sessioni interattive a sinistra. Al momento è vuoto perché non hai ancora eseguito alcun codice. Non appena esegui il notebook nel passaggio successivo, una sessione di calcolo serverless live verrà eseguita il provisioning in modo dinamico e verrà visualizzata qui.

Importare dati utilizzando Data Agent Kit

Anziché configurare manualmente una sessione Spark o scrivere da zero script di caricamento PySpark, programmerai in coppia con un agente utilizzando Data Agent Kit.

  1. Apri il riquadro Chat con l'agente facendo clic sull'icona Attiva/disattiva agente nella barra degli strumenti in alto a destra.
  2. Incolla il seguente prompt nella chat (assicurati di sostituire ${PROJECT_ID} con il tuo ID progetto Google Cloud effettivo):
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. Se l'agente chiede l'autorizzazione per eseguire comandi di verifica in background (ad es. "Consentire l'esecuzione di questo comando?"), esamina il comando proposto e seleziona Sì, consenti questa volta (o Sì, consenti sempre).
  2. Quando l'agente termina la generazione del file, fai clic sul pulsante blu Accetta tutto (o sull'icona del segno di spunta) nella parte inferiore del riquadro della chat per salvare notebooks/01_ingestion.ipynb nel tuo spazio di lavoro.

Agente che genera il notebook di importazione

Rivedi ed esegui il notebook

  1. Apri il nuovo notebooks/01_ingestion.ipynb generato nell'IDE.
  2. Esamina il codice PySpark per la logica di scrittura del connettore BigQuery.
  3. Fai clic su Esegui tutto nella barra degli strumenti del notebook dell'IDE.
  4. Se è la prima volta che esegui un notebook Spark remoto, l'IDE potrebbe chiederti di installare le dipendenze locali. Se richiesto, fai clic su Install dependencies for Remote Spark Kernels (Installa le dipendenze per i kernel Spark remoti), conferma le finestre di dialogo di installazione e poi fai di nuovo clic su Run All (Esegui tutto).
  5. Nel menu a discesa Seleziona kernel, scegli Kernel Spark remoti -> fraud-pipeline-runtime su Serverless Spark. (Suggerimento: se non vedi il modello di runtime preconfigurato elencato, fai clic sull'icona di aggiornamento nel menu a discesa del selettore del kernel in alto a destra per ricaricare i kernel remoti disponibili).
  6. Guarda la barra di stato in basso a sinistra nell'editor. Visualizzerai Connecting to kernel: fraud-pipeline-runtime on Serverless Spark.... Poiché si tratta dell'avvio iniziale del backend del kernel del runtime Spark Serverless, il provisioning e l'avvio richiederanno alcuni minuti.
  7. Una volta completata la connessione del kernel, il blocco note inizierà automaticamente a eseguire tutte le celle in sequenza per elaborare i log delle transazioni non elaborate nel set di dati BigQuery.

Verifica

Una volta completata l'esecuzione, controlla il catalogo Data Agent Kit per verificare la creazione della tabella:

Verifica la tabella Raw in Esplora cataloghi

  1. Nella barra delle attività dell'IDE, apri il riquadro Google Cloud Data Agent Kit.
  2. Espandi la sezione CATALOGO.
  3. Espandi l'ID progetto.
  4. Espandi BigQuery.
  5. Espandi il set di dati transactions_dataset_evals.
  6. Fai clic sulla tabella raw_transactions per aprire la visualizzazione dettagliata nell'editor principale.
  7. Nel riquadro di navigazione a sinistra, esplora le schede Dati, Schema e Dettagli per esaminare i record e i metadati importati.

Riepilogo della sezione:hai utilizzato il linguaggio naturale nella chat con l'agente per generare un carico di lavoro Spark serverless completo. Poi l'hai eseguito per elaborare i log JSON non strutturati in una tabella BigQuery (non elaborata).

4. Deduplicare e normalizzare con dbt

Prima di addestrare il modello di ML, applicherai la qualità dei dati rimuovendo i log di streaming duplicati, isolando i record errati (ad esempio gli ID transazione vuoti) e unendo i dati dimensionali (autori dei pagamenti e beneficiari). Questo processo richiede trasformazioni SQL idempotenti e affidabili, il che rende dbt (data build tool) una soluzione ideale.

Crea lo scaffolding della pipeline dbt

Utilizza l'agente per generare un progetto dbt sul set di dati BigQuery:

  1. Torna al riquadro Chat con l'agente.
  2. Fornisci la seguente istruzione per generare il progetto 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. L'agente presenterà un artefatto Implementation Plan nel riquadro principale dell'editor. Esamina la struttura dei file e la logica SQL proposte.
  2. Fai clic su Procedi (e poi su Accetta tutto) per consentire all'agente di generare i file nel tuo spazio di lavoro.

Piano di implementazione con il pulsante Procedi

  1. Al termine della generazione, l'agente mostra una procedura dettagliata che riepiloga i nuovi componenti. Se richiesto, accetta tutte le modifiche.

Accettare tutti i file generati nel riquadro di Chat

Sviluppare e testare

Sebbene l'agente abbia eseguito automaticamente dbt compile per garantire che l'SQL generato fosse valido dal punto di vista sintattico, ora materializzerai queste viste e tabelle in BigQuery ed eseguirai i test di qualità dei dati per la verifica locale. Nota: più avanti nel lab, automatizzerai questo passaggio di dbt nell'ambito di un DAG Airflow end-to-end.

  1. Nella barra delle attività all'estrema sinistra, fai clic sull'icona Explorer (o premi Cmd/Ctrl+Shift+E).
  2. Espandi dbt_project -> models per esaminare i modelli SQL generati. Fai clic su enriched_transactions.sql per aprire e rivedere la logica della funzionalità di trasformazione e frode nell'editor.
  3. In Esplora file, fai clic con il tasto destro del mouse sulla cartella dbt_project e seleziona Apri nel terminale integrato. Si apre automaticamente un riquadro del terminale impostato direttamente sulla directory di lavoro dbt_project richiesta.
  4. Se non hai già installato dbt, crea un ambiente virtuale al di fuori di dbt_project/ (nella directory principale della tua home page o del tuo spazio di lavoro) e installa l'adattatore BigQuery:
python3 -m venv ~/.venv/dbt
source ~/.venv/dbt/bin/activate
pip install dbt-bigquery
  1. Esegui i modelli dbt e i relativi test di qualità dei dati:
dbt build
  1. Guarda l'output del terminale. dbt compilerà l'SQL, materializzerà le tabelle di staging e arricchite in BigQuery ed eseguirà i test dei dati.

Crea build e testa il progetto dbt nel terminale integrato

  1. Al termine della build, chiudi il riquadro del terminale per liberare spazio sullo schermo per i passaggi rimanenti.

Riepilogo della sezione:hai generato un progetto dbt con l'agente, eseguito test di qualità dei dati e trasformato i record non elaborati in tabelle BigQuery di staging e arricchite.

5. Addestra un modello di rilevamento delle frodi distribuito con Random Forest

Con le transazioni arricchite materializzate in BigQuery, creerai un modello di machine learning per classificare gli eventi fraudolenti. Random Forest è un metodo di apprendimento di insieme adatto ai dati di classificazione tabulare. L'esecuzione di un RandomForestClassifier su Spark Serverless distribuisce l'addestramento del modello tra i nodi di lavoro senza richiedere la gestione dell'infrastruttura.

In questo passaggio, utilizzerai l'agente per generare la pipeline di addestramento Spark ML.

Genera il notebook di addestramento ML

  1. Apri il riquadro Chat con l'agente.
  2. Fornisci il seguente prompt per progettare la sequenza di addestramento del modello (ricorda di sostituire ${PROJECT_ID} con il tuo ID progetto attivo):
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. Esamina il piano o il codice generato dell'agente e fai clic su Procedi / Accetta tutto per salvare notebooks/02_training.ipynb nel tuo spazio di lavoro.

Agente che genera il notebook di addestramento

Rivedi ed esegui il notebook

  1. Apri notebooks/02_training.ipynb nell'editor.
  2. Esamina le fasi della pipeline PySpark ML per la codifica delle funzionalità, l'assemblaggio dei vettori e la logica di classificazione Random Forest.
  3. Fai clic su Esegui tutto nella barra degli strumenti del notebook dell'IDE.
  4. Quando si apre il selettore del menu a discesa Seleziona kernel, seleziona fraud-pipeline-runtime su Serverless Spark.

Selezione del kernel Serverless Spark per il notebook di addestramento

Verifica

Al termine dell'esecuzione, verifica che il modello sia stato addestrato ed esportato correttamente:

  1. Esamina gli output della cella di valutazione vicino alla parte inferiore del notebook per verificare il punteggio dell'area sotto la curva ROC (AUC) riportato.
  2. Per assicurarti che gli artefatti del modello siano stati salvati correttamente in GCS, espandi il riquadro dell'explorer STORAGE nella barra laterale di Data Agent Kit.
  3. Individua il bucket che termina con -models (collegato al tuo ID progetto attivo), espandilo e fai il drill-down per verificare che esistano la directory fraud_model e le relative fasi della pipeline.

Verifica che il modello sia stato salvato in GCS

Riepilogo della sezione:hai utilizzato l'agente per creare una pipeline di addestramento ML PySpark, hai addestrato un modello Random Forest sulla tabella BigQuery arricchita e hai esportato il modello in Cloud Storage.

6. Inferenza batch e scrittura di Cloud Spanner

Con un modello predittivo addestrato archiviato in Cloud Storage, eseguirai l'inferenza batch sulle nuove transazioni che scorrono in BigQuery. Le transazioni ad alto rischio devono essere indirizzate a un sistema operativo in modo che un team di conformità possa esaminarle. Cloud Spanner fornisce un database transazionale scalabile per questa coda di revisione.

Genera il notebook di inferenza batch

Utilizza l'agente per creare un notebook di inferenza che colleghi BigQuery, Cloud Storage e Cloud Spanner:

  1. Apri il riquadro Chat con l'agente.
  2. Fornisci il seguente prompt:
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. Accetta il notebook generato per salvare notebooks/03_inference.ipynb nel tuo workspace.

Agente che genera il notebook di inferenza

Rivedi ed esegui il notebook

  1. Apri il nuovo notebooks/03_inference.ipynb generato nell'editor.
  2. Esamina la sequenza di inferenza PySpark:
    • Dipendenze:il modello di runtime serverless fornisce le dipendenze JAR cloud-spanner richieste per l'esecuzione di Spark.
    • Formattazione dei dati:lo script elimina le colonne vettoriali Spark ML complesse (come le funzionalità e le probabilità non elaborate) prima di scrivere in modo che corrispondano allo schema della tabella Spanner.
    • Connettore Spanner:scrive le righe segnalate utilizzando .format("cloud-spanner") per aggiungerle direttamente alla coda di revisione.
  3. Fai clic su Esegui tutto nella barra degli strumenti del notebook dell'IDE.
  4. Quando ti viene chiesto di selezionare un kernel, seleziona fraud-pipeline-runtime su Serverless Spark.

Verifica

Una volta terminata l'elaborazione del notebook di inferenza, puoi eseguire query sul tuo database operativo Spanner direttamente all'interno dell'IDE:

  1. Nella barra delle attività dell'IDE, apri il riquadro Google Cloud Data Agent Kit.
  2. Espandi la sezione CATALOGO.
  3. Espandi l'ID progetto, quindi espandi Spanner.
  4. Vai a cymbal-fraud -> fraud-db -> Tables -> SparkEvalFraudReviewQueue.
  5. Fai clic con il tasto destro del mouse sulla tabella e seleziona Query tabella, quindi esegui la query:
SELECT *
FROM `SparkEvalFraudReviewQueue`
LIMIT 100;
  1. Nel riquadro Risultati delle query di seguito, dovresti vedere le righe appena inserite che rappresentano le transazioni ad alto rischio segnalate per la revisione manuale.

Verifica le righe in Cloud Spanner

Riepilogo della sezione:hai utilizzato l'agente per creare un notebook di inferenza batch, hai assegnato un punteggio ai record BigQuery senza etichetta con il modello addestrato e hai scritto le transazioni ad alto rischio direttamente in Cloud Spanner.

7. Strutturare e orchestrare con Managed Airflow

La pipeline è attualmente costituita da passaggi discreti: un notebook di importazione, un progetto di trasformazione dbt e un notebook di inferenza batch. Per renderlo pronto per la produzione, li unirai in un grafico delle dipendenze pianificato.

Managed Service for Apache Airflow (in precedenza noto come Cloud Composer) fornisce un motore di orchestrazione gestito per questo flusso di lavoro. Data Agent Kit include una funzionalità Orchestration Pipelines che traduce le definizioni dichiarative delle pipeline YAML direttamente in DAG Airflow.

Definisci la pipeline

Utilizza l'agente per generare la configurazione della pipeline di orchestrazione:

  1. Nella chat dell'agente, fornisci il seguente prompt (ricordati di sostituire ${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.

Rivedi la configurazione del DAG

Data Agent Kit Orchestrator utilizza configurazioni YAML dichiarative per definire e implementare pipeline in Apache Airflow, consentendo il controllo delle versioni e l'implementazione tramite CI/CD.

Nel riquadro Explorer dell'IDE, esamina i due file della pipeline generati dall'agente nella directory principale dello spazio di lavoro:

  1. deployment.yaml: apri questo file. Questo funge da registro dell'ambiente. Mappa la pipeline logica dev all'ambiente cymbal-airflow, imposta la regione di esecuzione (us-central1) e definisce il bucket artifact_storage in cui vengono preparati i DAG compilati e le dipendenze.
  2. fraud_analysis_pipeline.yaml: apri questo file. Definisce il grafico di esecuzione. Specifica la pianificazione del trigger (interval: '0 0 * * *') e sequenzia i tre passaggi nel blocco actions:
    • Un'azione di importazione notebook per 01_ingestion.ipynb in esecuzione su Dataproc Serverless.
    • Un'azione di trasformazione pipeline che ha come target la directory dbt_project, con una dipendenza dependsOn che punta al passaggio di importazione.
    • Un'azione di inferenza notebook per 03_inference.ipynb con una dipendenza dependsOn che punta al passaggio dbt, raggruppando la proprietà JAR di Spanner.
  3. L'agente riepilogherà anche questi artefatti generati in una scheda Procedura dettagliata nel riquadro dell'editor, delineando le configurazioni e le convalide eseguite.

Configurazione interattiva del DAG

Data Agent Kit esegue il rendering della configurazione della pipeline come grafico visivo interattivo per l'ispezione e la modifica delle proprietà DAG di Airflow.

  1. Nella barra delle attività dell'IDE, apri il riquadro Google Cloud Data Agent Kit.
  2. In DATA ENGINEERING, espandi Orchestration Pipelines.
  3. Fai clic su fraud_analysis_pipeline.yaml per aprire il canvas DAG visivo nell'editor principale.

Canvas visivo del DAG di orchestrazione

  1. Fai clic sul nodo Schedule trigger in alto. A destra si apre un riquadro a comparsa di configurazione che mostra la stringa Cron analizzata (0 0 * * *) e consente di modificare parametri come il backfill e il recupero.
  2. Fai clic sul nodo dell'attività del notebook (ad esempio il passaggio di importazione o inferenza). Il riquadro a comparsa si aggiorna per mostrare le mappature di esecuzione di Dataproc Serverless e le proprietà del connettore specifiche.
  3. Nota l'hyperlink al nome file del notebook (ad es. 01_ingestion.ipynb) all'interno del blocco del nodo. Se fai clic, il notebook si apre direttamente nell'editor.
  4. Nella barra laterale sinistra, sotto Pipeline di orchestrazione, fai clic su Deployment configuration. Questa visualizzazione mostra il cluster dell'ambiente dev di destinazione e gli artefatti del bucket GCS di output.

Riepilogo della sezione: hai generato una configurazione della pipeline di orchestrazione con l'agente, definendo le dipendenze tra le attività di importazione, dbt e inferenza in un canvas visivo interattivo.

8. Esegui il deployment, l'esecuzione e il monitoraggio

Con il DAG definito localmente, ti connetterai all'ambiente Managed Airflow di cui è stato eseguito il provisioning durante la configurazione ed eseguirai il deployment della pipeline.

Configura Managed Service for Apache Airflow

Prima del deployment, configura la connessione dello scheduler nelle impostazioni di Data Agent Kit in modo che l'estensione abbia come target il tuo ambiente Managed Airflow:

  1. Nella barra delle attività dell'IDE, apri il riquadro Google Cloud Data Agent Kit.
  2. In SETTINGS, fai clic su Impostazioni.
  3. Seleziona Pianificatore dal menu a sinistra.
  4. Configura le impostazioni:
    • ID progetto: seleziona l'ID progetto attivo.
    • Regione: seleziona us-central1.
    • Ambiente: seleziona cymbal-airflow.
  5. Fai clic su Salva.

Impostazioni di Managed Service for Apache Airflow

Esegui il deployment del DAG

Ora esegui il deployment della pipeline configurata direttamente nell'ambiente Managed Airflow dal canvas visivo:

  1. Nella barra laterale Google Cloud Data Agent Kit, espandi DATA ENGINEERING > Orchestration Pipelines e fai clic su fraud_analysis_pipeline.yaml per aprire il canvas DAG visivo.
  2. Nell'angolo in alto a destra della barra degli strumenti del canvas, fai clic sul pulsante blu Esegui pipeline.
  3. Nel selettore a discesa dell'ambiente, seleziona dev.
  4. Osserva la notifica di avanzamento nell'area di stato in basso (Running pipeline: Building pipeline locally...). L'estensione compilerà automaticamente il DAG, comprimerà il notebook e gli asset dbt e li caricherà nel bucket GCS dell'ambiente Managed Airflow (il completamento dell'operazione richiede circa 3-4 minuti).

Deployment della pipeline dal canvas visivo

Monitorare l'esecuzione

Una volta completata la compilazione locale e la notifica popup conferma Triggered a new run for pipeline... successfully, monitora l'esecuzione live:

  1. Nella barra laterale Google Cloud Data Agent Kit, espandi DATA ENGINEERING > Orchestration Pipelines.
  2. Fai clic su Gestione pipeline.
  3. Nella tabella Gestione pipeline, fai clic su fraud_analysis_pipeline per aprire la cronologia di esecuzione.

Panoramica della gestione delle pipeline

  1. Nella visualizzazione Cronologia esecuzioni, seleziona l'esecuzione attiva dal calendario.
  2. Man mano che l'esecuzione procede in ogni attività della pipeline (importazione, trasformazione dbt e inferenza), gli indicatori di stato si aggiornano e le durate delle attività vengono compilate. Fai clic su un'attività per ispezionare l'output di esecuzione in tempo reale e i log DAG di Airflow.

Cronologia di esecuzione della pipeline live e dettagli delle attività

Riepilogo della sezione:hai configurato la connessione di Airflow Scheduler, hai eseguito il deployment della pipeline di analisi end-to-end in Managed Airflow e hai monitorato un'esecuzione in tempo reale, verificando il sistema dai log non elaborati alle previsioni finali di Cloud Spanner.

9. Esegui la pulizia

Per evitare che al tuo progetto Google Cloud vengano addebitati costi continui per le risorse utilizzate in questo codelab, smantella l'ambiente utilizzando lo script automatizzato.

  1. Nel riquadro Terminale (o in Cloud Shell), vai alla directory degli script ed esegui:
cd ~/devrel-demos/codelabs/agentic-data-labs/data-science/scripts
chmod +x teardown.sh
./teardown.sh
  1. Lo script elencherà tutte le risorse che prevede di eliminare e chiederà la conferma:
    • Ambiente Airflow gestito (cymbal-airflow)
    • Istanza Cloud Spanner (cymbal-fraud)
    • Set di dati BigQuery (transactions_dataset_evals)
    • Bucket Cloud Storage (gs://${PROJECT_ID}-fin-clearing-raw e gs://${PROJECT_ID}-models)
    • Service account worker (composer-worker-sa)
  2. Digita y per confermare. Lo script di rimozione rimuoverà tutti i servizi Google Cloud di cui è stato eseguito il provisioning e libererà spazio sui file locali.

10. Complimenti!

Hai creato una pipeline end-to-end di rilevamento delle frodi che comprende Cloud Storage, BigQuery, Managed Service for Apache Spark (Spark Serverless), dbt, Cloud Spanner e Managed Service for Apache Airflow, con la programmazione in coppia con Google Cloud Data Agent Kit all'interno dell'IDE Antigravity.

Che cosa hai realizzato

  1. 📥 Registri delle transazioni non elaborate inseriti in una tabella BigQuery utilizzando Managed Service for Apache Spark e Data Agent Kit.
  2. 🧹 Dati deduplicati e normalizzati creando un progetto dbt con test di qualità dei dati.
  3. 🤖 È stato addestrato un modello Random Forest distribuito utilizzando RandomForestClassifier ed è stato esportato in Cloud Storage.
  4. ⚡ Eseguita l'inferenza batch sulle transazioni in entrata e instradati i record ad alto rischio in Cloud Spanner per la revisione dell'audit.
  5. 🔄 Orchestra, esegui il deployment e monitora il flusso di lavoro come DAG di Airflow pianificato utilizzando Managed Service for Apache Airflow e gli strumenti di gestione visiva dei DAG dell'IDE.

Concetti fondamentali

Concetto

Che cosa hai imparato

Data Agent Kit

Programmazione in coppia all'interno dell'IDE utilizzando il linguaggio naturale per generare notebook PySpark, configurare modelli dbt e definire DAG Airflow

BigQuery

Spazio di archiviazione tabellare scalabile per SQL analitico, trasformazioni dbt e addestramento ML

Spark serverless

Esecuzione serverless per il caricamento distribuito dei dati PySpark e l'addestramento ML Random Forest

Cloud Spanner Connector

Scrittura delle previsioni di inferenza batch di Spark direttamente nelle code di revisione del database operativo

Dichiarazioni DAG YAML

Definizioni di pipeline dichiarative visualizzate come grafici visivi interattivi di Airflow nell'IDE

Gestione DAG visiva

Ispezione delle dipendenze della pipeline, deployment in Managed Airflow e monitoraggio della cronologia di esecuzione delle attività live all'interno dell'IDE

Passaggi successivi