1. Einführung
Stellen Sie sich vor, Sie arbeiten als Data Scientist bei Cymbal Financial, einem Zahlungsabwickler mit hohem Transaktionsvolumen. Es ist eine Welle von Abrechnungsverzögerungen aufgetreten und das Compliance-Team vermutet koordinierten Betrug. Sie müssen eine Pipeline erstellen, um Rohdaten aus Clearinghouse-Transaktionsprotokollen aufzunehmen, die Daten zu bereinigen, ein Modell für maschinelles Lernen zu trainieren, Batch-Inferenz auszuführen und Transaktionen mit hohem Risiko in eine Cloud Spanner-Warteschlange für die manuelle Prüfung zu übertragen.
Normalerweise erfordert dies tagelanges Schreiben von sich wiederholendem Einrichtungscode (Spark-Notebooks, dbt-Konfigurationen, Trainingsskripts, Airflow-DAGs) und ständiges Wechseln zwischen Konsolenschnittstellen und Editoren.
In diesem Codelab arbeiten Sie mit einem Agenten zusammen, der das Google Cloud Data Agent Kit (DAK) in der Antigravity IDE verwendet. Mithilfe von konversationeller natürlicher Sprache hilft Ihnen der Agent, Spark-Notebooks zu erstellen, ein dbt-Projekt zu kompilieren, eine Inferenzschleife zu erstellen und den Workflow mit dem Managed Service for Apache Airflow zu orchestrieren.
Aufgaben
- Clearinghouse-Logs aus Cloud Storage mit Managed Service for Apache Spark (Spark Serverless) in eine BigQuery-Tabelle aufnehmen.
- Transaktionen deduplizieren und normalisieren mit dbt, um saubere Datenschichten (Rohdaten, Staging, angereichert) zu erstellen.
- Verteiltes Random Forest-Klassifikationsmodell trainieren (
RandomForestClassifier) in Spark Serverless. - Batchinferenz für neue Transaktionen ausführen und Benachrichtigungen zu hohem Risiko direkt in Cloud Spanner schreiben.
- Die gesamte Pipeline orchestrieren, visuell konfigurieren und bereitstellen – mit Managed Service for Apache Airflow und interaktiver DAG-Überwachung in der IDE.
Voraussetzungen
- Ein Webbrowser wie Chrome
- Ein Google Cloud-Projekt mit aktivierter Abrechnung (wir empfehlen, für die praktischen Labs ein neues, dediziertes Projekt zu verwenden).
- Grundkenntnisse in SQL, Python und PySpark
- Antigravity IDE mit einem Google AI Pro-Abo (empfohlen)
Die in diesem Codelab erstellten Ressourcen sollten weniger als 5 $ kosten. Folgen Sie am Ende des Labs der Anleitung zum Bereinigen, um bereitgestellte Ressourcen zu löschen.
2. Umgebung einrichten
Um das Lab zu starten, führen Sie ein Bootstrap-Skript aus. Mit diesem Skript werden automatisch die erforderlichen GCP-APIs aktiviert, ein Cloud Storage-Bucket für die Aufnahme erstellt, Mock-Transaktions- und ‑Verzeichnis-Datasets generiert, Referenzverzeichnisse in BigQuery geladen und die Hintergrundbereitstellung von Cloud Spanner und Managed Service for Apache Airflow (ehemals Cloud Composer) gestartet.
Projekt auswählen oder erstellen
Wählen Sie in der Google Cloud Console ein vorhandenes Projekt aus oder erstellen Sie ein neues.
Abrechnung bestätigen
Die Abrechnung für das Google Cloud-Projekt muss aktiviert sein. Weitere Informationen
Setupskript ausführen
Sie verwenden Google Cloud Shell (oder Ihre lokale Shell, die mit der Google Cloud CLI konfiguriert ist), um die Einrichtung der Umgebung zu starten.
- Öffnen Sie die Google Cloud Console.
- Klicken Sie in der Symbolleiste rechts oben auf Cloud Shell aktivieren.

- Konfigurieren Sie im Cloud Shell-Terminal Ihr aktives Projekt:
gcloud config set project <<YOUR_PROJECT_ID>>
export PROJECT_ID=$(gcloud config get-value project)
- Klonen Sie das Codelab-Repository und wechseln Sie in den Ordner „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
- Führen Sie das Bootstrap-Setupscript aus, um alle Ressourcen in
us-central1bereitzustellen:
chmod +x setup.sh setup_spanner.sh setup_composer.sh
export REGION=us-central1
./setup.sh
- Wenn das Skript fertig ist, wird eine Zusammenfassung ausgegeben, die angibt, dass Ihr BigQuery-Dataset und Ihr Cloud Storage-Bucket bereit sind. Im Hintergrund werden Cloud Spanner (dauert ca. 2 Minuten) und Managed Airflow (dauert ca. 20 Minuten) weiterhin bereitgestellt. Sie können den Fortschritt jederzeit mit folgendem Befehl verfolgen:
tail -f /tmp/spanner_setup.log
tail -f /tmp/composer_setup.log
Antigravity IDE öffnen
- Laden Sie die Antigravity IDE von der Google Antigravity-Downloadseite herunter und installieren Sie sie.
- Starten Sie die Antigravity IDE.
- Erstellen Sie einen neuen, leeren Ordner auf Ihrem lokalen Computer (z. B. mit dem Namen
agentic-data-labs) und öffnen Sie ihn in der IDE, indem Sie Ordner öffnen auswählen. Dieser Ordner dient als Ihr lokaler Arbeitsbereich für das Codelab.

Data Agent Kit-Erweiterung installieren
Die Erweiterung „Google Cloud Data Agent Kit“ bietet eine enge Integration mit Google Cloud-Datendiensten direkt in Ihrem Editor. So können Sie mit BigQuery, Cloud SQL, Cloud Storage und anderen Diensten interagieren, ohne den Kontext wechseln zu müssen.
- Klicken Sie in der Antigravity IDE in der Aktivitätsleiste ganz links auf dem Bildschirm auf das Symbol Erweiterungen (es sieht aus wie vier Quadrate).
- Geben Sie oben im Bereich „Erweiterungen“ in der Suchleiste
Google Cloud Data Agent Kitein. - Suchen Sie nach der Erweiterung Google Cloud Data Agent Kit, die von
googlecloudtoolsveröffentlicht wurde. - Klicken Sie auf die Schaltfläche Installieren.
- Möglicherweise werden Sie gefragt, ob Sie dem Herausgeber „googlecloudtools“ und seinen Erweiterungen vertrauen. Klicken Sie auf Verlagen und Webpublishern vertrauen und installieren, um fortzufahren.

Nach der Installation wird in der Aktivitätsleiste ganz links in der Antigravity IDE ein neues Symbol für das Google Cloud Data Agent Kit angezeigt.
- Eine Onboarding-Seite mit dem Titel „Willkommen im Google Cloud Data Agent Kit“ sollte automatisch geöffnet werden. Wenn Sie nicht in Ihrem Cloud-Konto angemeldet sind, folgen Sie der Anleitung, um den Zugriff zu gewähren.
- Suchen Sie im Abschnitt Konfigurationszusammenfassung nach dem Projektfeld. Klicken Sie auf das Drop-down-Menü und wählen Sie Ihr Google Cloud-Projekt aus. Legen Sie als Region
us-central1fest. Wählen Sie dann MCP-Server konfigurieren aus.

- Wählen Sie MCP-Server konfigurieren aus. Achten Sie darauf, dass im Bereich MCP Configuration (MCP-Konfiguration) die folgenden Remote-MCP-Server aktiviert sind:
- BigQuery
- Spanner
- Notebooks
Klicken Sie dann auf Jetzt starten.

Konfigurationsoptionen ansehen
Nach Abschluss der Einrichtung werden Sie zur Seite „Erste Schritte mit dem Google Cloud Data Agent Kit“ weitergeleitet.
- Klicken Sie unter „Einrichtung und Konfiguration“ auf Jetzt starten.
- Dadurch wird der Bereich Data Agent Kit Configuration (Konfiguration des Data Agent Kit) geöffnet. Tabs ansehen:
- Projekt und Region:Prüfen Sie die ausgewählte Projekt-ID und bestätigen Sie, dass das Setupscript alle erforderlichen APIs (Compute Engine, Cloud Storage, BigQuery, Spanner usw.) aktiviert hat.
- BigQuery:Konfigurieren Sie den standardmäßigen Standort für Ihre BigQuery-Abfragen. Verwenden Sie die Region
us-central1. - MCP-Server konfigurieren:Hier sehen Sie die aktivierten MCP-Server (BigQuery, Notebooks, Spanner usw.), die es KI-Agenten ermöglichen, sicher mit Ihren Daten zu interagieren.
- Skills:Vordefinierte Skills bieten Agents spezielle Funktionen für komplexe Datenaufgaben.

Zusammenfassung des Abschnitts:Sie haben das Bootstrap-Skript ausgeführt, um GCS- und BigQuery-Assets zu erstellen, während Spanner und Airflow im Hintergrund erstellt wurden. Sie haben das Projekt dann in der Antigravity IDE geöffnet und die Erweiterung „Google Cloud Data Agent Kit“ aktiviert. Jetzt können Sie Ihr erstes Notebook schreiben.
3. Rohlogs mit Spark Serverless aufnehmen
In diesem Abschnitt nehmen Sie rohe JSON-Transaktionsprotokolle in den Data Lake auf. Managed Service for Apache Spark (Spark Serverless) stellt eine direkte Verbindung zum nativen Speicher von BigQuery her. Sie verwenden den Standard-BigQuery-Connector, um tabellarische Daten zu verwalten und direkte Abfragen und Analysen zu ermöglichen.
Vorkonfigurierte Spark Serverless-Laufzeit kennenlernen
Bevor Sie Spark-Code ausführen, sollten Sie die Serverless Runtime-Vorlage prüfen, die vom Setupscript vorkonfiguriert wurde. In dieser Vorlage werden das Backend der Zielausführungsumgebung und die erforderlichen Connector-Abhängigkeiten definiert.
- Öffnen Sie in der Aktivitätsleiste der IDE den Bereich Google Cloud Data Agent Kit.
- Maximieren Sie das Drop-down-Menü Apache Spark und dann Serverless.
- Klicken Sie mit der rechten Maustaste auf
fraud-pipeline-runtimeund wählen Sie Profil aus, um die Konfigurationsansicht im Editor zu öffnen. - Scrollen Sie auf dem Tab Profil nach unten und maximieren Sie Eigenschaften, um die benutzerdefinierten Abhängigkeiten zu prüfen, die an die Umgebung angehängt sind:
spark.jars: Enthältgs://spark-lib/spanner/spark-3.5-spanner-1.4.0.jar, das den Spark Spanner-Connector verwendet, damit Spark-Jobs später im Lab Inferenz-Ergebnisse direkt in Cloud Spanner schreiben können. Hinweis: Dataproc Serverless enthält standardmäßig den Spark BigQuery-Connector von Google Cloud. Es ist also keine zusätzliche JAR-Konfiguration zum Lesen und Schreiben von BigQuery-Tabellen erforderlich.

- Beachten Sie den Tab Interaktive Sitzungen auf der linken Seite. Er ist derzeit leer, da Sie noch keinen Code ausgeführt haben. Sobald Sie das Notebook im nächsten Schritt ausführen, wird dynamisch eine aktive serverlose Compute-Sitzung bereitgestellt, die hier angezeigt wird.
Daten mit dem Data Agent Kit aufnehmen
Anstatt eine Spark-Sitzung manuell zu konfigurieren oder PySpark-Ladeskripts von Grund auf neu zu schreiben, arbeiten Sie mit einem Agenten zusammen, der das Data Agent Kit verwendet.
- Öffnen Sie den Bereich Agent Chat, indem Sie oben rechts in der Symbolleiste auf das Symbol Toggle Agent (Agent ein-/ausblenden) klicken.
- Fügen Sie den folgenden Prompt in den Chat ein und ersetzen Sie dabei
${PROJECT_ID}durch Ihre tatsächliche Google Cloud-Projekt-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.
- Wenn der KI-Agent um die Erlaubnis bittet, Hintergrundüberprüfungsbefehle auszuführen (z.B. „Allow running this command?“), prüfen Sie den vorgeschlagenen Befehl und wählen Sie Yes, allow this time (oder Yes, and always allow) aus.
- Wenn der Agent die Datei fertig generiert hat, klicken Sie unten im Chatbereich auf die blaue Schaltfläche Alle akzeptieren (oder das Häkchensymbol), um
notebooks/01_ingestion.ipynbin Ihrem Arbeitsbereich zu speichern.

Notebook überprüfen und ausführen
- Öffnen Sie die neu generierte
notebooks/01_ingestion.ipynbin der IDE. - Sehen Sie sich den PySpark-Code für die Schreiblogik des BigQuery-Connectors an.
- Klicken Sie in der Notebook-Symbolleiste der IDE auf Alle ausführen.
- Wenn Sie zum ersten Mal ein Spark-Notebook remote ausführen, werden Sie in der IDE möglicherweise aufgefordert, lokale Abhängigkeiten zu installieren. Klicken Sie bei Aufforderung auf Abhängigkeiten für Remote-Spark-Kernel installieren, bestätigen Sie die Installationsdialogfelder und klicken Sie dann noch einmal auf Alle ausführen.
- Wählen Sie im Drop-down-Menü Select Kernel (Kernel auswählen) die Option Remote Spark Kernels (Remote-Spark-Kernel) -> fraud-pipeline-runtime on Serverless Spark (fraud-pipeline-runtime auf Serverless Spark) aus. Tipp: Wenn Ihre vorkonfigurierte Laufzeitvorlage nicht aufgeführt ist, klicken Sie rechts oben im Drop-down-Menü für die Kernelauswahl auf das Aktualisierungssymbol, um die verfügbaren Remote-Kernel neu zu laden.
- Sehen Sie sich die Statusleiste unten links im Editor an. Sie sehen
Connecting to kernel: fraud-pipeline-runtime on Serverless Spark.... Da dies der erste Start des Spark Serverless-Laufzeit-Kernel-Backends ist, dauert die Bereitstellung und das Hochfahren einige Minuten. - Sobald die Verbindung des Kernels hergestellt ist, werden alle Zellen im Notebook automatisch sequenziell ausgeführt, um die Rohdaten der Transaktionsprotokolle in Ihrem BigQuery-Dataset zu verarbeiten.
Bestätigung
Prüfen Sie nach Abschluss der Ausführung den Data Agent Kit-Katalog, um die Tabellenerstellung zu bestätigen:

- Öffnen Sie in der Aktivitätsleiste der IDE den Bereich Google Cloud Data Agent Kit.
- Maximieren Sie den Bereich KATALOG.
- Maximieren Sie Ihre Projekt-ID.
- Maximieren Sie BigQuery.
- Maximieren Sie das Dataset
transactions_dataset_evals. - Klicken Sie auf die Tabelle
raw_transactions, um die Detailansicht im Haupteditor zu öffnen. - Sehen Sie sich in der linken Navigationsleiste die Tabs Daten, Schema und Details an, um die aufgenommenen Datensätze und Metadaten zu prüfen.
Zusammenfassung des Abschnitts:Sie haben natürliche Sprache im Agent-Chat verwendet, um eine vollständige serverlose Spark-Arbeitslast zu generieren. Anschließend haben Sie sie ausgeführt, um unstrukturierte JSON-Logs in eine BigQuery-Rohdatentabelle zu verarbeiten.
4. Mit dbt deduplizieren und normalisieren
Vor dem Training des ML-Modells sorgen Sie für Datenqualität, indem Sie doppelte Streamingprotokolle entfernen, fehlerhafte Datensätze (z. B. leere Transaktions-IDs) isolieren und dimensionale Daten (Zahler und Zahlungsempfänger) zusammenführen. Für diesen Prozess sind idempotente, zuverlässige SQL-Transformationen erforderlich, weshalb sich dbt (data build tool) hervorragend eignet.
dbt-Pipeline erstellen
Verwenden Sie den Agent, um ein dbt-Projekt für das BigQuery-Dataset zu generieren:
- Kehren Sie zum Bereich Agent Chat zurück.
- Geben Sie den folgenden Befehl ein, um das dbt-Projekt zu generieren:
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.
- Der Agent präsentiert im Haupteditorbereich ein Implementierungsplan-Artefakt. Prüfen Sie die vorgeschlagene Dateistruktur und SQL-Logik.
- Klicken Sie auf Weiter und dann auf Alle akzeptieren, damit der Agent die Dateien in Ihrem Arbeitsbereich generieren kann.

- Nachdem der Agent die Generierung abgeschlossen hat, wird eine Kurzanleitung mit einer Zusammenfassung der neuen Komponenten angezeigt. Akzeptieren Sie alle Änderungen, wenn Sie dazu aufgefordert werden.

Erstellen und testen
Der Agent hat dbt compile automatisch ausgeführt, um sicherzustellen, dass der generierte SQL-Code syntaktisch gültig ist. Sie müssen diese Ansichten und Tabellen jetzt in BigQuery materialisieren und die Datenqualitätstests zur lokalen Überprüfung ausführen. Hinweis: Später in diesem Lab automatisieren Sie diesen dbt-Schritt als Teil eines End-to-End-Airflow-DAG.
- Klicken Sie in der Aktivitätsleiste ganz links auf das Symbol Explorer (oder drücken Sie
Cmd/Ctrl+Shift+E). - Maximieren Sie
dbt_project->models, um die generierten SQL-Modelle zu prüfen. Klicken Sie aufenriched_transactions.sql, um die Logik der Transformations- und Betrugsfunktionen im Editor zu öffnen und zu prüfen. - Klicken Sie im Datei-Explorer mit der rechten Maustaste auf den Ordner
dbt_projectund wählen Sie Im integrierten Terminal öffnen aus. Dadurch wird automatisch ein Terminalbereich geöffnet, der direkt auf das erforderliche Arbeitsverzeichnisdbt_projecteingestellt ist. - Wenn Sie
dbtnoch nicht installiert haben, erstellen Sie eine virtuelle Umgebung außerhalb vondbt_project/(im Stammverzeichnis Ihres privaten oder geschäftlichen Arbeitsbereichs) und installieren Sie den BigQuery-Adapter:
python3 -m venv ~/.venv/dbt
source ~/.venv/dbt/bin/activate
pip install dbt-bigquery
- Führen Sie die dbt-Modelle und die zugehörigen Datenqualitätstests aus:
dbt build
- Sehen Sie sich die Terminalausgabe an. dbt kompiliert den SQL-Code, materialisiert die Staging- und angereicherten Tabellen in BigQuery und führt die Datentests aus.

- Schließen Sie nach Abschluss des Builds das Terminalfenster, um Platz auf dem Bildschirm für die verbleibenden Schritte zu schaffen.
Zusammenfassung des Abschnitts:Sie haben mit dem Agent ein dbt-Projekt generiert, Datenqualitätstests ausgeführt und die Rohdatensätze in Staging- und angereicherte BigQuery-Tabellen transformiert.
5. Verteiltes Betrugserkennungsmodell mit Random Forest trainieren
Anschließend erstellen Sie mit den angereicherten Transaktionen in BigQuery ein Machine-Learning-Modell, um betrügerische Ereignisse zu klassifizieren. Random Forest ist eine Ensemble-Lernmethode, die sich gut für tabellarische Klassifikationsdaten eignet. Wenn Sie eine RandomForestClassifier in Spark Serverless ausführen, wird das Modelltraining auf Worker-Knoten verteilt, ohne dass Sie die Infrastruktur verwalten müssen.
In diesem Schritt verwenden Sie den Agent, um die Spark ML-Trainingspipeline zu generieren.
ML-Trainings-Notebook generieren
- Öffnen Sie den Bereich Agent Chat.
- Geben Sie den folgenden Prompt ein, um die Modelltrainingssequenz zu entwerfen. Denken Sie daran,
${PROJECT_ID}durch Ihre aktive Projekt-ID zu ersetzen:
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.
- Sehen Sie sich den Plan oder den generierten Code des Agents an und klicken Sie auf Weiter / Alle akzeptieren, um
notebooks/02_training.ipynbin Ihrem Arbeitsbereich zu speichern.

Notebook überprüfen und ausführen
- Öffnen Sie
notebooks/02_training.ipynbim Editor. - Sehen Sie sich die PySpark ML-Pipelinestufen für die Feature-Codierung, die Vektor-Assembly und die Random Forest-Klassifizierungslogik an.
- Klicken Sie in der Notebook-Symbolleiste der IDE auf Alle ausführen.
- Wenn die Drop-down-Auswahl Kernel auswählen geöffnet wird, wählen Sie fraud-pipeline-runtime on Serverless Spark aus.

Bestätigung
Prüfen Sie nach Abschluss der Ausführung, ob das Modell richtig trainiert und exportiert wurde:
- Sehen Sie sich die Ausgaben der Auswertungszelle unten im Notebook an, um den gemeldeten AUC-Wert (Area Under ROC) zu überprüfen.
- Wenn Sie prüfen möchten, ob die Modellartefakte erfolgreich in GCS gespeichert wurden, maximieren Sie in der Seitenleiste des Data Agent Kit den Bereich STORAGE.
- Suchen Sie den Bucket, der mit
-modelsendet (mit Ihrer aktiven Projekt-ID verknüpft), maximieren Sie ihn und führen Sie einen Drilldown durch, um zu prüfen, ob das Verzeichnisfraud_modelund seine Pipelinestufen vorhanden sind.

Zusammenfassung des Abschnitts:Sie haben den Agent verwendet, um eine PySpark ML-Trainingspipeline zu erstellen, ein Random Forest-Modell für Ihre angereicherte BigQuery-Tabelle trainiert und das Modell in Cloud Storage exportiert.
6. Batch-Inferenz und Cloud Spanner-Schreibvorgang
Mit einem trainierten Vorhersagemodell, das in Cloud Storage gespeichert ist, führen Sie die Batchinferenz für neue Transaktionen durch, die über BigQuery eingehen. Transaktionen mit hohem Risiko müssen an ein operatives System weitergeleitet werden, damit sie von einem Compliance-Team überprüft werden können. Cloud Spanner bietet eine skalierbare Transaktionsdatenbank für diese Warteschlange.
Batchinferenz-Notebook generieren
Verwenden Sie den Agent, um ein Inferenz-Notebook zu erstellen, das BigQuery, Cloud Storage und Cloud Spanner verbindet:
- Öffnen Sie den Bereich Agent Chat.
- Geben Sie den folgenden Prompt ein:
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.
- Akzeptieren Sie das generierte Notebook, um
notebooks/03_inference.ipynbin Ihrem Arbeitsbereich zu speichern.

Notebook überprüfen und ausführen
- Öffnen Sie die neu generierte Datei
notebooks/03_inference.ipynbim Editor. - Sehen Sie sich die PySpark-Inferenzsequenz an:
- Abhängigkeiten:Die Vorlage für die serverlose Laufzeit enthält die erforderlichen
cloud-spanner-JAR-Abhängigkeiten für die Spark-Ausführung. - Datenformatierung:Im Skript werden komplexe Spark ML-Vektorspalten (z. B. Roh-Features und Wahrscheinlichkeiten) vor dem Schreiben gelöscht, damit sie dem Spanner-Tabellenschema entsprechen.
- Spanner-Connector:Er schreibt die gekennzeichneten Zeilen mit
.format("cloud-spanner"), um sie direkt an die Warteschlange für die Überprüfung anzuhängen.
- Abhängigkeiten:Die Vorlage für die serverlose Laufzeit enthält die erforderlichen
- Klicken Sie in der Notebook-Symbolleiste der IDE auf Alle ausführen.
- Wenn Sie aufgefordert werden, einen Kernel auszuwählen, wählen Sie fraud-pipeline-runtime on Serverless Spark aus.
Bestätigung
Sobald das Inferenz-Notebook die Verarbeitung abgeschlossen hat, können Sie Ihre operative Spanner-Datenbank direkt in der IDE abfragen:
- Öffnen Sie in der Aktivitätsleiste der IDE den Bereich Google Cloud Data Agent Kit.
- Maximieren Sie den Bereich KATALOG.
- Maximieren Sie Ihre Projekt-ID und dann Spanner.
- Gehen Sie zu
cymbal-fraud->fraud-db->Tables->SparkEvalFraudReviewQueue. - Klicken Sie mit der rechten Maustaste auf die Tabelle und wählen Sie Tabelle abfragen aus. Führen Sie dann die Abfrage aus:
SELECT *
FROM `SparkEvalFraudReviewQueue`
LIMIT 100;
- Im Bereich Abfrageergebnisse unten sollten neu eingefügte Zeilen mit Transaktionen mit hohem Risiko angezeigt werden, die zur manuellen Überprüfung gekennzeichnet sind.

Zusammenfassung des Abschnitts:Sie haben den Agent verwendet, um ein Notebook für die Batchinferenz zu erstellen, ungelabelte BigQuery-Datensätze mit Ihrem trainierten Modell bewertet und Transaktionen mit hohem Risiko direkt in Cloud Spanner geschrieben.
7. Mit Managed Airflow Gerüste erstellen und orchestrieren
Ihre Pipeline besteht derzeit aus einzelnen Schritten: einem Notebook für die Aufnahme, einem dbt-Transformationsprojekt und einem Notebook für die Batch-Inferenz. Damit das Ganze produktionsreif ist, müssen Sie die einzelnen Schritte in einem geplanten Abhängigkeitsdiagramm zusammenfügen.
Managed Service for Apache Airflow (ehemals Cloud Composer) bietet eine verwaltete Orchestrierungs-Engine für diesen Workflow. Das Data Agent Kit enthält die Funktion Orchestration Pipelines, mit der deklarative YAML-Pipelinedefinitionen direkt in Airflow-DAGs übersetzt werden.
Pipeline definieren
Verwenden Sie den Agent, um die Konfiguration der Orchestrierungspipeline zu generieren:
- Geben Sie im Agent-Chat den folgenden Prompt ein (denken Sie daran,
${PROJECT_ID}zu ersetzen):
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-Konfiguration überprüfen
Der Data Agent Kit Orchestrator verwendet deklarative YAML-Konfigurationen, um Pipelines für Apache Airflow zu definieren und bereitzustellen. So können Definitionen per CI/CD versionsverwaltet und bereitgestellt werden.
Sehen Sie sich im Bereich „IDE Explorer“ die beiden Pipeline-Dateien an, die der Agent im Stammverzeichnis Ihres Arbeitsbereichs generiert hat:
deployment.yaml: Öffnen Sie diese Datei. Dies dient als Ihre Umgebungsregistrierung. Sie ordnet Ihre logischedev-Pipeline dercymbal-airflow-Umgebung zu, legt die Ausführungsregion (us-central1) fest und definiert denartifact_storage-Bucket, in dem kompilierte DAGs und Abhängigkeiten bereitgestellt werden.fraud_analysis_pipeline.yaml: Öffnen Sie diese Datei. Damit wird die Ausführungsgrafik definiert. Er gibt den Triggerzeitplan (interval: '0 0 * * *') an und ordnet die drei Schritte unter dem Blockactionsan:- Eine
notebook-Aufnahmeaktion für01_ingestion.ipynb, die in Dataproc Serverless ausgeführt wird. - Eine Transformationsaktion
pipeline, die auf das Verzeichnisdbt_projectausgerichtet ist, mit einerdependsOn-Abhängigkeit, die auf den Erfassungsschritt verweist. - Eine
notebook-Inferenzaktion für03_inference.ipynbmit einerdependsOn-Abhängigkeit, die auf den dbt-Schritt verweist und die Spanner-JAR-Property bündelt.
- Eine
- Der Agent fasst diese generierten Artefakte auch auf dem Tab Walkthrough (Übersicht) im Bearbeitungsbereich zusammen und beschreibt die durchgeführten Konfigurationen und Validierungen.
Interaktive DAG-Konfiguration
Das Data Agent Kit rendert Ihre Pipelinekonfiguration als interaktives visuelles Diagramm, mit dem Sie Airflow-DAG-Eigenschaften prüfen und bearbeiten können.
- Öffnen Sie in der Aktivitätsleiste der IDE den Bereich Google Cloud Data Agent Kit.
- Erweitern Sie unter
DATA ENGINEERINGdie OptionOrchestration Pipelines. - Klicken Sie auf
fraud_analysis_pipeline.yaml, um den visuellen DAG-Arbeitsbereich im Haupteditor zu öffnen.

- Klicken Sie oben auf den Knoten
Schedule trigger. Rechts wird ein Konfigurations-Flyout geöffnet, in dem der geparste Cron-String (0 0 * * *) angezeigt wird. Sie können Parameter wie „Backfill“ und „Catchup“ anpassen. - Klicken Sie auf einen Notebook-Aufgabenknoten, z. B. den Ingestions- oder Inferenzschritt. Im Flyout werden die spezifischen Ausführungszuordnungen und Connector-Attribute für Dataproc Serverless angezeigt.
- Beachten Sie den Hyperlink zum Notebook-Dateinamen (z. B.
01_ingestion.ipynb) im Knotenblock. Wenn Sie darauf klicken, wird das Notebook direkt im Editor geöffnet. - Klicken Sie in der linken Seitenleiste unter „Orchestration Pipelines“ auf
Deployment configuration. In dieser Ansicht sehen Sie den Zielcluster derdev-Umgebung und die Ausgabeartefakte des GCS-Buckets.
Zusammenfassung des Abschnitts:Sie haben mit dem Agent eine Orchestrierungspipeline-Konfiguration erstellt und Abhängigkeiten zwischen Aufgaben für Aufnahme, dbt und Inferenz in einem interaktiven visuellen Arbeitsbereich definiert.
8. Bereitstellen, ausführen und überwachen
Nachdem Sie den DAG lokal definiert haben, stellen Sie eine Verbindung zur Managed Airflow-Umgebung her, die während der Einrichtung bereitgestellt wurde, und stellen die Pipeline bereit.
Managed Service for Apache Airflow konfigurieren
Konfigurieren Sie vor der Bereitstellung die Scheduler-Verbindung in den Data Agent Kit-Einstellungen so, dass die Erweiterung auf Ihre Managed Airflow-Umgebung ausgerichtet ist:
- Öffnen Sie in der Aktivitätsleiste der IDE den Bereich Google Cloud Data Agent Kit.
- Klicken Sie unter
SETTINGSauf Einstellungen. - Wählen Sie im Menü auf der linken Seite Scheduler aus.
- Konfigurieren Sie die Einstellungen:
- Projekt-ID: Wählen Sie Ihre aktive Projekt-ID aus.
- Region: Wählen Sie
us-central1aus. - Umgebung: Wählen Sie
cymbal-airflowaus.
- Klicken Sie auf Speichern.

DAG bereitstellen
Sie stellen die konfigurierte Pipeline jetzt direkt über den visuellen Arbeitsbereich in Ihrer Managed Airflow-Umgebung bereit:
- Maximieren Sie in der Seitenleiste Google Cloud Data Agent Kit die Elemente
DATA ENGINEERING>Orchestration Pipelinesund klicken Sie auffraud_analysis_pipeline.yaml, um den visuellen DAG-Arbeitsbereich zu öffnen. - Klicken Sie rechts oben in der Arbeitsbereich-Symbolleiste auf den blauen Button Pipeline ausführen.
- Wählen Sie im Drop-down-Menü für die Umgebung
devaus. - Beobachten Sie die Fortschrittsbenachrichtigung im unteren Statusbereich (
Running pipeline: Building pipeline locally...). Die Erweiterung kompiliert automatisch Ihren DAG, packt das Notebook und die dbt-Assets und lädt sie in den GCS-Bucket Ihrer Managed Airflow-Umgebung hoch. Dieser Vorgang dauert etwa 3 bis 4 Minuten.

Lauf überwachen
Sobald die lokale Kompilierung abgeschlossen ist und die Pop-up-Benachrichtigung Triggered a new run for pipeline... successfully bestätigt, überwachen Sie die Live-Ausführung:
- Erweitern Sie in der Seitenleiste des Google Cloud Data Agent Kit die Optionen
DATA ENGINEERING>Orchestration Pipelines. - Klicken Sie auf Pipelines verwalten.
- Klicken Sie in der Tabelle „Pipelines verwalten“ auf
fraud_analysis_pipeline, um den Ausführungsverlauf zu öffnen.

- Wählen Sie in der Ansicht Ausführungsverlauf den aktiven Lauf im Kalender aus.
- Während die Ausführung der einzelnen Pipeline-Aufgaben (Erfassung, dbt-Transformation und Inferenz) fortschreitet, werden die Statusanzeigen aktualisiert und die Aufgabendauern ausgefüllt. Klicken Sie auf eine beliebige Aufgabe, um die Ausgabe der Live-Ausführung und die Airflow-DAG-Logs zu prüfen.

Zusammenfassung des Abschnitts:Sie haben die Airflow Scheduler-Verbindung konfiguriert, Ihre End-to-End-Analysepipeline in Managed Airflow bereitgestellt und eine Live-Ausführung überwacht. Dabei haben Sie das System von Rohlogs bis hin zu finalen Cloud Spanner-Vorhersagen geprüft.
9. Bereinigen
Damit Ihrem Google Cloud-Projekt für die in diesem Codelab verwendeten Ressourcen keine laufenden Gebühren in Rechnung gestellt werden, müssen Sie die Umgebung mit dem automatisierten Skript abbauen.
- Wechseln Sie im Bereich Terminal (oder in Cloud Shell) zum Verzeichnis „scripts“ und führen Sie Folgendes aus:
cd ~/devrel-demos/codelabs/agentic-data-labs/data-science/scripts
chmod +x teardown.sh
./teardown.sh
- Im Skript werden alle Ressourcen aufgeführt, die gelöscht werden sollen, und Sie werden zur Bestätigung aufgefordert:
- Managed Airflow-Umgebung (
cymbal-airflow) - Cloud Spanner-Instanz (
cymbal-fraud) - BigQuery-Dataset (
transactions_dataset_evals) - Cloud Storage-Buckets (
gs://${PROJECT_ID}-fin-clearing-rawundgs://${PROJECT_ID}-models) - Worker-Dienstkonto (
composer-worker-sa)
- Managed Airflow-Umgebung (
- Geben Sie zur Bestätigung
yein. Mit dem Teardown-Skript werden alle bereitgestellten GCP-Dienste entfernt und lokale Dateien bereinigt.
10. Glückwunsch!
Sie haben eine End-to-End-Pipeline zur Betrugserkennung erstellt, die Cloud Storage, BigQuery, Managed Service for Apache Spark (Spark Serverless), dbt, Cloud Spanner und Managed Service for Apache Airflow umfasst. Dabei haben Sie in der Antigravity IDE mit dem Google Cloud Data Agent Kit zusammengearbeitet.
Ihre Erfolge
- 📥 Rohdaten-Transaktionslogs mit Managed Service for Apache Spark und dem Data Agent Kit in eine BigQuery-Tabelle aufgenommen.
- 🧹 Doppelte Daten entfernt und Daten normalisiert, indem ein dbt-Projekt mit Datenqualitätsprüfungen erstellt wurde.
- 🤖 Ein verteiltes Random Forest-Modell mit
RandomForestClassifiertrainiert und das trainierte Modell in Cloud Storage exportiert. - ⚡ Batchinferenz ausgeführt für eingehende Transaktionen und Datensätze mit hohem Risiko zur Überprüfung in Cloud Spanner weitergeleitet.
- 🔄 Workflow als geplanter Airflow-DAG orchestriert, bereitgestellt und überwacht mit Managed Service for Apache Airflow und den visuellen DAG-Verwaltungstools der IDE.
Wichtige Konzepte
Konzept | Das haben Sie gelernt |
Pair-Programming in der IDE mit natürlicher Sprache zum Generieren von PySpark-Notebooks, Konfigurieren von dbt-Modellen und Definieren von Airflow-DAGs | |
Skalierbarer tabellarischer Speicher für analytisches SQL, dbt-Transformationen und ML-Training | |
Serverlose Ausführung für das verteilte Laden von PySpark-Daten und das ML-Training mit Random Forest | |
Batch-Spark-Inferenzvorhersagen direkt in die Warteschlangen für die Überprüfung der Betriebsdatenbank schreiben | |
YAML-DAG-Deklarationen | Deklarative Pipeline-Definitionen, die als interaktive Airflow-Diagramme in der IDE gerendert werden |
Visuelle DAG-Verwaltung | Pipelineabhängigkeiten prüfen, in Managed Airflow bereitstellen und den Live-Aufgabenausführungsverlauf in der IDE überwachen |
Nächste Schritte
- Dokumentation zum Google Cloud Data Agent Kit
- Weitere Informationen zu Managed Service for Apache Spark
- Weitere Informationen zu Managed Service for Apache Airflow
- Eigene Pipelines für mehrere Dienste mit der Antigravity IDE erstellen