Pipeline de détection des fraudes avec Data Agent Kit et Antigravity IDE

1. Introduction

Imaginez que vous êtes data scientist chez Cymbal Financial, un processeur de paiements à fort volume. Une vague de retards de paiements s'est produite, et l'équipe chargée de la conformité soupçonne une fraude coordonnée. Vous devez créer un pipeline pour ingérer les journaux bruts des transactions de la chambre de compensation, nettoyer les données, entraîner un modèle de machine learning, exécuter l'inférence par lot et transférer les transactions à haut risque dans une file d'attente d'examen Cloud Spanner pour un audit manuel.

Normalement, cela nécessite des jours d'écriture de code de configuration répétitif (notebooks Spark, configurations dbt, scripts d'entraînement, DAG Airflow) et un changement de contexte constant entre les interfaces de console et les éditeurs.

Dans cet atelier de programmation, vous allez programmer en binôme avec un agent à l'aide du Google Cloud Data Agent Kit (DAK) dans l'IDE Antigravity. À l'aide du langage naturel conversationnel, l'agent vous aidera à générer des notebooks Spark, à compiler un projet dbt, à construire une boucle d'inférence et à orchestrer le workflow à l'aide de Managed Service pour Apache Airflow.

Objectifs de l'atelier

  • Ingérez les journaux de la chambre de compensation depuis Cloud Storage à l'aide de Managed Service pour Apache Spark (Spark sans serveur) dans une table BigQuery.
  • Supprimez les doublons et normalisez les transactions à l'aide de dbt pour établir des couches de données propres (brutes, intermédiaires, enrichies).
  • Entraîner un modèle de classification de forêt aléatoire distribué (RandomForestClassifier) sur Spark sans serveur.
  • Exécutez l'inférence par lot sur les nouvelles transactions et écrivez les alertes à haut risque directement dans Cloud Spanner.
  • Orchestrez, configurez visuellement et déployez l'intégralité du pipeline à l'aide de Managed Service pour Apache Airflow et de la surveillance interactive des DAG dans l'IDE.

Prérequis

  • Un navigateur Web (par exemple, Chrome)
  • Un projet Google Cloud avec facturation activée (nous vous recommandons d'utiliser un nouveau projet dédié pour les ateliers pratiques).
  • Connaître les bases de SQL, Python et PySpark
  • IDE Antigravity avec un abonnement Google AI Pro (recommandé)

Les ressources créées dans cet atelier de programmation devraient coûter moins de 5 $. Veillez à suivre les instructions de la section Nettoyer à la fin de l'atelier pour supprimer les ressources provisionnées.

2. Configuration de l'environnement

Pour commencer l'atelier, vous allez exécuter un script d'amorçage. Ce script active automatiquement les API GCP requises, crée un bucket Cloud Storage pour l'ingestion, génère des ensembles de données de transactions et d'annuaires fictifs, charge les annuaires de référence dans BigQuery et lance le provisionnement en arrière-plan de Cloud Spanner et de Managed Service pour Apache Airflow (anciennement Cloud Composer).

Sélectionner ou créer un projet

Choisissez un projet existant ou créez-en un dans la console Google Cloud.

Valider la facturation

Vérifiez que la facturation est activée pour votre projet Google Cloud. Pour en savoir plus, consultez ce guide.

Exécuter le script de configuration

Vous utiliserez Google Cloud Shell (ou votre interface système locale configurée avec la Google Cloud CLI) pour lancer la configuration de l'environnement.

  1. Ouvrez la console Google Cloud.
  2. Cliquez sur Activer Cloud Shell dans la barre d'outils en haut à droite.

Ouvrir Cloud Shell

  1. Dans le terminal Cloud Shell, configurez votre projet actif :
gcloud config set project <<YOUR_PROJECT_ID>>
export PROJECT_ID=$(gcloud config get-value project)
  1. Clonez le dépôt de l'atelier de programmation et accédez au dossier des 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. Exécutez le script de configuration d'amorçage pour déployer toutes les ressources sur us-central1 :
chmod +x setup.sh setup_spanner.sh setup_composer.sh
export REGION=us-central1
./setup.sh
  1. Une fois le script terminé, un récapitulatif s'affiche pour vous indiquer que votre ensemble de données BigQuery et votre bucket Cloud Storage sont prêts. En arrière-plan, Cloud Spanner (environ deux minutes) et Managed Airflow (environ 20 minutes) continueront de provisionner. Vous pouvez suivre leur progression à tout moment en exécutant la commande suivante :
tail -f /tmp/spanner_setup.log
tail -f /tmp/composer_setup.log

Ouvrir l'IDE Antigravity

  1. Téléchargez et installez l'IDE Antigravity depuis la page de téléchargement de Google Antigravity.
  2. Lancez l'IDE Antigravity.
  3. Créez un dossier vide sur votre machine locale (par exemple, agentic-data-labs) et ouvrez-le dans l'IDE en sélectionnant Ouvrir le dossier. Il servira d'espace de travail local pour l'atelier de programmation.

Configurer le dossier du projet Antigravity IDE

Installer l'extension Data Agent Kit

L'extension Google Cloud Data Agent Kit offre une intégration approfondie aux services de données Google Cloud directement dans votre éditeur. Vous pouvez ainsi interagir avec BigQuery, Cloud SQL, Cloud Storage et d'autres services sans changer de contexte.

  1. Dans l'IDE Antigravity, cliquez sur l'icône Extensions dans la barre d'activité, tout à gauche de l'écran (elle ressemble à quatre carrés).
  2. Dans la barre de recherche en haut du volet "Extensions", saisissez Google Cloud Data Agent Kit.
  3. Recherchez l'extension Google Cloud Data Agent Kit publiée par googlecloudtools.
  4. Cliquez sur le bouton Install (installer).
  5. Une invite peut s'afficher et vous demander si vous faites confiance à l'éditeur "googlecloudtools" et à ses extensions. Cliquez sur Faire confiance aux éditeurs et installer pour continuer.

Installer l'extension Data Agent Kit

Une fois installé, une nouvelle icône Google Cloud Data Agent Kit s'affiche dans la barre d'activité, tout à gauche de l'IDE Antigravity.

  1. Une page d'intégration intitulée "Bienvenue dans Google Cloud Data Agent Kit" devrait s'ouvrir automatiquement. Si vous n'êtes pas connecté à votre compte Cloud, suivez les instructions pour autoriser l'accès.
  2. Dans la section Résumé de la configuration, recherchez le champ "Projet". Cliquez sur le menu déroulant et sélectionnez votre projet Google Cloud. Définissez votre région sur us-central1. Sélectionnez ensuite Configurer les serveurs MCP.

Configuration initiale de l'extension Data Agent Kit

  1. Sélectionnez Configurer les serveurs MCP. Dans le volet Configuration MCP, assurez-vous d'activer les serveurs MCP distants suivants :
    • BigQuery
    • Spanner
    • Notebooks

Cliquez ensuite sur Commencer.

Configurer les serveurs MCP

Explorer les options de configuration

Une fois la configuration terminée, vous serez redirigé vers la page "Get started with Google Cloud Data Agent Kit" (Commencer à utiliser Google Cloud Data Agent Kit).

  1. Sous "Configuration", cliquez sur Premiers pas.
  2. Le panneau Configuration du kit d'agent de données s'ouvre. Explorez les onglets :
    • Projet et région : vérifiez l'ID du projet sélectionné et assurez-vous que le script de configuration a activé toutes les API requises (Compute Engine, Cloud Storage, BigQuery, Spanner, etc.).
    • BigQuery : configurez l'emplacement par défaut de vos requêtes BigQuery. Utilisez la région us-central1.
    • Configurer les serveurs MCP : affichez les serveurs MCP activés (BigQuery, Notebooks, Spanner, etc.) qui permettent aux agents IA d'interagir de manière sécurisée avec vos données.
    • Compétences : explorez les compétences prédéfinies qui offrent aux agents des capacités spécialisées pour les tâches de données complexes.

Panneau des paramètres Data Agent Kit

Récapitulatif de la section : vous avez exécuté le script d'amorçage pour créer des composants GCS et BigQuery, tandis que Spanner et Airflow se sont créés en arrière-plan. Vous avez ensuite ouvert le projet dans l'IDE Antigravity et activé l'extension Google Cloud Data Agent Kit. Vous êtes maintenant prêt à écrire votre premier notebook.

3. Ingérer des journaux bruts à l'aide de Spark sans serveur

Dans cette section, vous allez ingérer des journaux de transactions JSON bruts dans le lac de données. Managed Service pour Apache Spark (Spark sans serveur) se connecte directement au stockage natif de BigQuery. Vous utiliserez le connecteur BigQuery standard pour gérer les données tabulaires et permettre les requêtes et les analyses directes.

Explorer l'environnement d'exécution Spark sans serveur préconfiguré

Avant d'exécuter le code Spark, inspectez le modèle Serverless Runtime préconfiguré par le script de configuration. Ce modèle définit le backend de l'environnement d'exécution cible et regroupe les dépendances de connecteur nécessaires.

  1. Dans la barre d'activité de l'IDE, ouvrez le panneau Google Cloud Data Agent Kit.
  2. Développez le menu déroulant Apache Spark, puis Sans serveur.
  3. Effectuez un clic droit sur fraud-pipeline-runtime, puis sélectionnez Profil pour ouvrir la vue de configuration dans l'éditeur.
  4. Dans l'onglet Profil, faites défiler la page vers le bas et développez Propriétés pour inspecter les dépendances personnalisées associées à l'environnement :
    • spark.jars : contient gs://spark-lib/spanner/spark-3.5-spanner-1.4.0.jar, qui utilise le connecteur Spark Spanner pour permettre aux jobs Spark d'écrire les résultats d'inférence directement dans Cloud Spanner plus tard dans l'atelier. (Remarque : Dataproc sans serveur inclut le connecteur Spark BigQuery de Google Cloud par défaut. Aucune configuration de fichier JAR supplémentaire n'est requise pour lire et écrire des tables BigQuery.)

Explorer les propriétés d'exécution Spark sans serveur

  1. Notez l'onglet Sessions interactives à gauche. Elle est actuellement vide, car vous n'avez pas encore exécuté de code. Dès que vous exécuterez le notebook à l'étape suivante, une session de calcul sans serveur en direct sera provisionnée de manière dynamique et s'affichera ici.

Ingérer des données à l'aide de Data Agent Kit

Au lieu de configurer manuellement une session Spark ou d'écrire des scripts de chargement PySpark à partir de zéro, vous allez programmer en binôme avec un agent à l'aide de Data Agent Kit.

  1. Ouvrez le volet Chat de l'agent en cliquant sur l'icône Activer/Désactiver l'agent dans la barre d'outils en haut à droite.
  2. Collez la requête suivante dans le chat (en veillant à remplacer ${PROJECT_ID} par l'ID de votre projet Google Cloud) :
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. Si l'agent demande l'autorisation d'exécuter des commandes de vérification en arrière-plan (par exemple, Autoriser l'exécution de cette commande ?), examinez la commande proposée et sélectionnez Oui, autoriser cette fois (ou Oui, et toujours autoriser).
  2. Lorsque l'agent a terminé de générer le fichier, cliquez sur le bouton bleu Tout accepter (ou sur l'icône en forme de coche) en bas du panneau de discussion pour enregistrer notebooks/01_ingestion.ipynb dans votre espace de travail.

Agent générant le notebook d'ingestion

Examiner et exécuter le notebook

  1. Ouvrez le fichier notebooks/01_ingestion.ipynb nouvellement généré dans l'IDE.
  2. Examinez le code PySpark pour la logique d'écriture du connecteur BigQuery.
  3. Cliquez sur Tout exécuter dans la barre d'outils du notebook de l'IDE.
  4. Si c'est la première fois que vous exécutez un notebook Spark à distance, l'IDE peut vous inviter à installer des dépendances locales. Si vous y êtes invité, cliquez sur Install dependencies for Remote Spark Kernels (Installer les dépendances pour les noyaux Spark à distance), confirmez les boîtes de dialogue d'installation, puis cliquez à nouveau sur Run All (Tout exécuter).
  5. Dans le menu déroulant Sélectionner un noyau, sélectionnez Noyaux Spark à distance > fraud-pipeline-runtime sur Serverless Spark. (Conseil : Si votre modèle d'exécution préconfiguré ne s'affiche pas, cliquez sur l'icône Actualiser en haut à droite du menu déroulant du sélecteur de noyau pour recharger les noyaux distants disponibles.)
  6. Consultez la barre d'état en bas à gauche de l'éditeur. Le message Connecting to kernel: fraud-pipeline-runtime on Serverless Spark... s'affiche. Étant donné qu'il s'agit du lancement initial du backend du kernel d'exécution Spark sans serveur, le provisionnement et le démarrage prendront quelques minutes.
  7. Une fois le noyau connecté, le notebook commence automatiquement à exécuter toutes les cellules de manière séquentielle pour traiter les journaux de transactions brutes dans votre ensemble de données BigQuery.

Validation

Une fois l'exécution terminée, vérifiez le catalogue Data Agent Kit pour confirmer la création de la table :

Vérifier la table &quot;Raw&quot; dans l'explorateur de catalogue

  1. Dans la barre d'activité de l'IDE, ouvrez le panneau Google Cloud Data Agent Kit.
  2. Développez la section CATALOGUE.
  3. Développez l'ID de votre projet.
  4. Développez BigQuery.
  5. Développez l'ensemble de données transactions_dataset_evals.
  6. Cliquez sur la table raw_transactions pour ouvrir sa vue détaillée dans l'éditeur principal.
  7. Dans le panneau de navigation de gauche, explorez les onglets Données, Schéma et Détails pour inspecter les enregistrements et les métadonnées ingérés.

Récapitulatif de la section : vous avez utilisé le langage naturel dans le chat de l'agent pour générer une charge de travail Spark sans serveur complète. Vous l'avez ensuite exécutée pour traiter les journaux JSON non structurés dans une table BigQuery (brute).

4. Dédupliquer et normaliser avec dbt

Avant d'entraîner le modèle de ML, vous allez appliquer la qualité des données en supprimant les journaux de streaming en double, en isolant les enregistrements incorrects (tels que les ID de transaction vides) et en joignant les données dimensionnelles (payeurs et bénéficiaires). Ce processus nécessite des transformations SQL idempotentes et fiables, ce qui fait de dbt (data build tool) un outil idéal.

Créer la structure du pipeline dbt

Utilisez l'agent pour générer un projet dbt sur l'ensemble de données BigQuery :

  1. Revenez au volet Chat de l'agent.
  2. Fournissez l'instruction suivante pour générer le projet 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'agent présente un artefact de plan d'implémentation dans le volet de l'éditeur principal. Examinez la structure de fichier et la logique SQL proposées.
  2. Cliquez sur Continuer (puis sur Tout accepter) pour autoriser l'agent à générer les fichiers dans votre espace de travail.

Plan de mise en œuvre avec bouton &quot;Continuer&quot;

  1. Une fois la génération terminée, l'agent affiche un tutoriel récapitulant les nouveaux composants. Acceptez toutes les modifications si vous y êtes invité.

Accepter tous les fichiers générés dans le panneau Chat

Compilation et tests

Bien que l'agent ait exécuté automatiquement dbt compile pour s'assurer que le code SQL généré était syntaxiquement valide, vous allez maintenant matérialiser ces vues et tables dans BigQuery, puis exécuter les tests de qualité des données pour la vérification locale. (Remarque : Plus loin dans l'atelier, vous automatiserez cette étape dbt dans un DAG Airflow de bout en bout.)

  1. Dans la barre d'activité située tout à gauche, cliquez sur l'icône Explorateur (ou appuyez sur Cmd/Ctrl+Shift+E).
  2. Développez dbt_project > models pour inspecter les modèles SQL générés. Cliquez sur enriched_transactions.sql pour ouvrir et examiner la logique de la fonctionnalité de transformation et de détection de la fraude dans l'éditeur.
  3. Dans l'explorateur de fichiers, effectuez un clic droit sur le dossier dbt_project, puis sélectionnez Ouvrir dans le terminal intégré. Cela ouvre automatiquement un volet de terminal directement défini sur le répertoire de travail dbt_project requis.
  4. Si dbt n'est pas encore installé, créez un environnement virtuel en dehors de dbt_project/ (à la racine de votre répertoire personnel ou de votre espace de travail) et installez l'adaptateur BigQuery :
python3 -m venv ~/.venv/dbt
source ~/.venv/dbt/bin/activate
pip install dbt-bigquery
  1. Exécutez les modèles dbt et leurs tests de qualité des données associés :
dbt build
  1. Observez la sortie du terminal. dbt compilera le code SQL, matérialisera les tables intermédiaires et enrichies dans BigQuery, et exécutera les tests de données.

Compiler et tester le projet dbt dans le terminal intégré

  1. Une fois la compilation terminée, fermez le volet du terminal pour libérer de l'espace à l'écran pour les étapes restantes.

Récapitulatif de la section : vous avez généré un projet dbt avec l'agent, exécuté des tests de qualité des données et transformé les enregistrements bruts en tables BigQuery intermédiaires et enrichies.

5. Entraîner un modèle de détection de fraude distribué avec Random Forest

Une fois les transactions enrichies matérialisées dans BigQuery, vous allez créer un modèle de machine learning pour classer les événements frauduleux. Random Forest est une méthode d'apprentissage par ensemble qui convient parfaitement aux données de classification tabulaires. L'exécution d'un RandomForestClassifier sur Spark Serverless distribue l'entraînement de modèle sur les nœuds de calcul sans que vous ayez à gérer l'infrastructure.

Dans cette étape, vous allez utiliser l'agent pour générer le pipeline d'entraînement Spark ML.

Générer le notebook d'entraînement au ML

  1. Ouvrez le volet Agent Chat.
  2. Fournissez la requête suivante pour concevoir la séquence d'entraînement du modèle (n'oubliez pas de remplacer ${PROJECT_ID} par l'ID de votre projet actif) :
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. Examinez le plan ou le code généré par l'agent, puis cliquez sur Continuer / Tout accepter pour enregistrer notebooks/02_training.ipynb dans votre espace de travail.

Agent générant le notebook d'entraînement

Examiner et exécuter le notebook

  1. Ouvrez notebooks/02_training.ipynb dans l'éditeur.
  2. Examinez les étapes du pipeline PySpark ML pour l'encodage des caractéristiques, l'assemblage des vecteurs et la logique de classification Random Forest.
  3. Cliquez sur Tout exécuter dans la barre d'outils du notebook de l'IDE.
  4. Lorsque le sélecteur de liste déroulante Sélectionner un noyau s'ouvre, sélectionnez fraud-pipeline-runtime sur Serverless Spark.

Sélectionner le noyau Serverless Spark pour le notebook d'entraînement

Validation

Une fois l'exécution terminée, vérifiez que le modèle a été entraîné et exporté correctement :

  1. Examinez les résultats de la cellule d'évaluation en bas du notebook pour vérifier le score AUC (Area Under ROC) indiqué.
  2. Pour vous assurer que les artefacts du modèle ont bien été enregistrés dans GCS, développez le volet de l'explorateur STORAGE dans la barre latérale Data Agent Kit.
  3. Localisez le bucket se terminant par -models (associé à votre ID de projet actif), développez-le, puis accédez aux détails pour vérifier que le répertoire fraud_model et ses étapes de pipeline existent.

Vérifier que le modèle est enregistré dans GCS

Récapitulatif de la section : vous avez utilisé l'agent pour créer un pipeline d'entraînement PySpark ML, entraîné un modèle Random Forest sur votre table BigQuery enrichie et exporté le modèle vers Cloud Storage.

6. Inférence par lot et écriture Cloud Spanner

Avec un modèle prédictif entraîné stocké dans Cloud Storage, vous allez exécuter l'inférence par lot sur les nouvelles transactions qui transitent par BigQuery. Les transactions à haut risque doivent être acheminées vers un système opérationnel afin qu'une équipe de conformité puisse les examiner. Cloud Spanner fournit une base de données transactionnelle évolutive pour cette file d'attente d'examen.

Générer le notebook d'inférence par lot

Utilisez l'agent pour créer un notebook d'inférence connectant BigQuery, Cloud Storage et Cloud Spanner :

  1. Ouvrez le volet Agent Chat.
  2. Fournissez le prompt suivant :
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. Acceptez le notebook généré pour enregistrer notebooks/03_inference.ipynb dans votre espace de travail.

Agent générant le notebook d'inférence

Examiner et exécuter le notebook

  1. Ouvrez le fichier notebooks/03_inference.ipynb nouvellement généré dans l'éditeur.
  2. Examinez la séquence d'inférence PySpark :
    • Dépendances : le modèle d'exécution sans serveur fournit les dépendances JAR cloud-spanner requises pour l'exécution de Spark.
    • Mise en forme des données : le script supprime les colonnes de vecteurs Spark ML complexes (telles que les caractéristiques brutes et les probabilités) avant l'écriture pour correspondre au schéma de table Spanner.
    • Connecteur Spanner : il écrit les lignes signalées à l'aide de .format("cloud-spanner") pour les ajouter directement à la file d'attente d'examen.
  3. Cliquez sur Tout exécuter dans la barre d'outils du notebook de l'IDE.
  4. Lorsque vous êtes invité à sélectionner un kernel, sélectionnez fraud-pipeline-runtime sur Serverless Spark.

Validation

Une fois le notebook d'inférence traité, vous pouvez interroger votre base de données Spanner opérationnelle directement dans l'IDE :

  1. Dans la barre d'activité de l'IDE, ouvrez le panneau Google Cloud Data Agent Kit.
  2. Développez la section CATALOGUE.
  3. Développez l'ID de votre projet, puis Spanner.
  4. Accédez à cymbal-fraud > fraud-db > Tables > SparkEvalFraudReviewQueue.
  5. Effectuez un clic droit sur la table, puis sélectionnez Interroger la table et exécutez la requête :
SELECT *
FROM `SparkEvalFraudReviewQueue`
LIMIT 100;
  1. Dans le volet Résultats de la requête ci-dessous, vous devriez voir les lignes nouvellement insérées représentant les transactions à haut risque signalées pour examen manuel.

Vérifier les lignes dans Cloud Spanner

Récapitulatif de la section : vous avez utilisé l'agent pour créer un notebook d'inférence par lot, vous avez attribué un score aux enregistrements BigQuery non libellés avec votre modèle entraîné et vous avez écrit les transactions à haut risque directement dans Cloud Spanner.

7. Échafauder et orchestrer avec Managed Airflow

Votre pipeline se compose actuellement d'étapes distinctes : un notebook d'ingestion, un projet de transformation dbt et un notebook d'inférence par lot. Pour que ce processus soit prêt pour la production, vous allez les assembler dans un graphique de dépendances planifié.

Managed Service pour Apache Airflow (anciennement Cloud Composer) fournit un moteur d'orchestration géré pour ce workflow. Le Data Agent Kit inclut une fonctionnalité Orchestration Pipelines qui traduit les définitions de pipeline YAML déclaratives directement en DAG Airflow.

Définir le pipeline

Utilisez l'agent pour générer la configuration du pipeline d'orchestration :

  1. Dans Agent Chat (Chat de l'agent), saisissez le prompt suivant (en remplaçant ${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.

Examiner la configuration du DAG

L'orchestrateur du Data Agent Kit utilise des configurations YAML déclaratives pour définir et déployer des pipelines sur Apache Airflow. Les définitions peuvent ainsi être contrôlées par version et déployées via CI/CD.

Dans le volet Explorateur de l'IDE, examinez les deux fichiers de pipeline générés par l'agent à la racine de votre espace de travail :

  1. deployment.yaml : ouvrez ce fichier. Il sert de registre d'environnement. Il mappe votre pipeline logique dev à l'environnement cymbal-airflow, définit la région d'exécution (us-central1) et définit le bucket artifact_storage dans lequel les DAG et les dépendances compilés sont mis en scène.
  2. fraud_analysis_pipeline.yaml : ouvrez ce fichier. Cela définit le graphique d'exécution. Il spécifie la programmation du déclencheur (interval: '0 0 * * *') et séquence les trois étapes sous le bloc actions :
    • Action d'ingestion notebook pour 01_ingestion.ipynb s'exécutant sur Dataproc sans serveur.
    • Action de transformation pipeline ciblant le répertoire dbt_project, avec une dépendance dependsOn pointant vers l'étape d'ingestion.
    • Action notebook d'inférence pour 03_inference.ipynb avec une dépendance dependsOn pointant vers l'étape dbt, regroupant la propriété JAR Spanner.
  3. L'agent résumera également ces artefacts générés dans un onglet Procédure pas à pas du volet de l'éditeur, en décrivant les configurations et les validations effectuées.

Configuration interactive des DAG

Le Data Agent Kit affiche la configuration de votre pipeline sous la forme d'un graphique visuel interactif permettant d'inspecter et de modifier les propriétés des DAG Airflow.

  1. Dans la barre d'activité de l'IDE, ouvrez le panneau Google Cloud Data Agent Kit.
  2. Sous DATA ENGINEERING, développez Orchestration Pipelines.
  3. Cliquez sur fraud_analysis_pipeline.yaml pour ouvrir le canevas DAG visuel dans l'éditeur principal.

Canevas visuel du DAG d'orchestration

  1. Cliquez sur le nœud Schedule trigger en haut de l'écran. Un menu déroulant de configuration s'ouvre à droite, affichant la chaîne Cron analysée (0 0 * * *) et vous permettant d'ajuster des paramètres tels que le remplissage et la rattrapage.
  2. Cliquez sur un nœud de tâche de notebook (par exemple, l'étape d'ingestion ou d'inférence). Le menu volant est mis à jour pour afficher les mappages d'exécution Dataproc sans serveur et les propriétés du connecteur spécifiques.
  3. Notez le lien hypertexte vers le nom de fichier du notebook (par exemple, 01_ingestion.ipynb) à l'intérieur du bloc de nœud. En cliquant dessus, vous ouvrez le notebook directement dans votre éditeur.
  4. Dans la barre latérale de gauche, sous "Orchestration Pipelines", cliquez sur Deployment configuration. Cette vue affiche le cluster d'environnement dev cible et les artefacts du bucket GCS de sortie.

Récapitulatif de la section : vous avez généré une configuration de pipeline d'orchestration avec l'agent, en définissant les dépendances entre les tâches d'ingestion, dbt et d'inférence dans un canevas visuel interactif.

8. Déployer, exécuter et surveiller

Une fois le DAG défini en local, vous vous connecterez à l'environnement Managed Airflow provisionné lors de la configuration et déploirez le pipeline.

Configurer Managed Service pour Apache Airflow

Avant le déploiement, configurez la connexion Scheduler dans les paramètres du kit Data Agent afin que l'extension cible votre environnement Managed Airflow :

  1. Dans la barre d'activité de l'IDE, ouvrez le panneau Google Cloud Data Agent Kit.
  2. Sous SETTINGS, cliquez sur Paramètres.
  3. Sélectionnez Planificateur dans le menu de gauche.
  4. Configurez les paramètres :
    • ID du projet : sélectionnez l'ID de votre projet actif.
    • Région : sélectionnez us-central1.
    • Environnement : sélectionnez cymbal-airflow.
  5. Cliquez sur Enregistrer.

Paramètres de Managed Service pour Apache Airflow

Déployer le DAG

Vous allez maintenant déployer le pipeline configuré directement dans votre environnement Managed Airflow à partir du canevas visuel :

  1. Dans la barre latérale Google Cloud Data Agent Kit, développez DATA ENGINEERING > Orchestration Pipelines, puis cliquez sur fraud_analysis_pipeline.yaml pour ouvrir le canevas visuel du DAG.
  2. En haut à droite de la barre d'outils du canevas, cliquez sur le bouton bleu Exécuter le pipeline.
  3. Dans le sélecteur de menu déroulant de l'environnement, sélectionnez dev.
  4. Observez la notification de progression dans la zone d'état en bas de l'écran (Running pipeline: Building pipeline locally...). L'extension compilera automatiquement votre DAG, empaquettera le notebook et les éléments dbt, puis les importera dans le bucket GCS de votre environnement Managed Airflow (cette opération prend environ trois à quatre minutes).

Déployer le pipeline depuis le canevas visuel

Surveiller l'exécution

Une fois la compilation locale terminée et la notification pop-up confirmant Triggered a new run for pipeline... successfully, surveillez l'exécution en direct :

  1. Dans la barre latérale Google Cloud Data Agent Kit, développez DATA ENGINEERING > Orchestration Pipelines.
  2. Cliquez sur Gestion des pipelines.
  3. Dans le tableau "Gestion des pipelines", cliquez sur fraud_analysis_pipeline pour ouvrir l'historique d'exécution.

Présentation de la gestion des pipelines

  1. Dans la vue Historique des exécutions, sélectionnez l'exécution active dans le calendrier.
  2. À mesure que l'exécution progresse dans chaque tâche du pipeline (ingestion, transformation dbt et inférence), les indicateurs d'état sont mis à jour et les durées des tâches sont renseignées. Cliquez sur une tâche pour inspecter sa sortie d'exécution en direct et les journaux DAG Airflow.

Historique d'exécution des pipelines en direct et détails des tâches

Récapitulatif de la section : vous avez configuré la connexion Airflow Scheduler, déployé votre pipeline analytique de bout en bout sur Managed Airflow et surveillé une exécution en direct, en vérifiant le système des journaux bruts aux prédictions finales de Cloud Spanner.

9. Effectuer un nettoyage

Pour éviter que les ressources utilisées dans cet atelier de programmation ne soient facturées sur votre projet Google Cloud, supprimez l'environnement à l'aide du script automatisé.

  1. Dans le panneau Terminal (ou dans Cloud Shell), accédez au répertoire des scripts et exécutez la commande suivante :
cd ~/devrel-demos/codelabs/agentic-data-labs/data-science/scripts
chmod +x teardown.sh
./teardown.sh
  1. Le script liste toutes les ressources qu'il prévoit de supprimer et demande une confirmation :
    • Environnement Airflow géré (cymbal-airflow)
    • Instance Cloud Spanner (cymbal-fraud)
    • Ensemble de données BigQuery (transactions_dataset_evals)
    • Buckets Cloud Storage (gs://${PROJECT_ID}-fin-clearing-raw et gs://${PROJECT_ID}-models)
    • Compte de service Worker (composer-worker-sa)
  2. Saisissez y pour confirmer. Le script de suppression supprimera tous les services GCP provisionnés et nettoiera les fichiers locaux.

10. Félicitations !

Vous avez créé un pipeline de détection des fraudes de bout en bout couvrant Cloud Storage, BigQuery, Managed Service pour Apache Spark (Spark Serverless), dbt, Cloud Spanner et Managed Service pour Apache Airflow, en programmation en binôme avec le Google Cloud Data Agent Kit dans l'IDE Antigravity.

Ce que vous avez accompli

  1. 📥 Ingestion des journaux de transactions bruts dans une table BigQuery à l'aide de Managed Service pour Apache Spark et du Data Agent Kit.
  2. 🧹 Dédoublonnez et normalisez les données en créant un projet dbt avec des tests de qualité des données.
  3. 🤖 Entraînement d'un modèle Random Forest distribué à l'aide de RandomForestClassifier et exportation du modèle entraîné vers Cloud Storage.
  4. ⚡ Exécution d'une inférence par lot sur les transactions entrantes et routage des enregistrements à haut risque vers Cloud Spanner pour examen d'audit.
  5. 🔄 Orchestrez, déployez et surveillez le workflow en tant que DAG Airflow planifié à l'aide de Managed Service pour Apache Airflow et des outils de gestion visuelle des DAG de l'IDE.

Concepts clés

Concept

Ce que vous avez appris

Data Agent Kit

Programmation en binôme dans l'IDE à l'aide du langage naturel pour générer des notebooks PySpark, configurer des modèles dbt et définir des DAG Airflow

BigQuery

Stockage tabulaire évolutif pour l'analyse SQL, les transformations dbt et l'entraînement ML

Spark Serverless

Exécution sans serveur pour l'entraînement ML Random Forest et le chargement de données PySpark distribuées

Connecteur Cloud Spanner

Écrire des prédictions d'inférence Spark par lot directement dans les files d'attente d'examen des bases de données opérationnelles

Déclarations de DAG YAML

Définitions de pipeline déclaratives affichées sous forme de graphiques visuels Airflow interactifs dans l'IDE

Gestion visuelle des DAG

Inspecter les dépendances du pipeline, déployer sur Managed Airflow et surveiller l'historique d'exécution des tâches en direct dans l'IDE

Étapes suivantes