Potok wykrywania oszustw z użyciem pakietu Data Agent Kit i IDE Antigravity

1. Wprowadzenie

Wyobraź sobie, że jesteś badaczem danych w firmie Cymbal Financial, która przetwarza dużą liczbę płatności. Wystąpiła fala opóźnień w rozliczeniach, a zespół ds. zgodności podejrzewa skoordynowane oszustwo. Musisz utworzyć potok do pozyskiwania nieprzetworzonych logów transakcji z centrum rozliczeniowego, oczyszczania danych, trenowania modelu uczenia maszynowego, przeprowadzania wnioskowania wsadowego i przekazywania transakcji o wysokim ryzyku do kolejki weryfikacji Cloud Spanner w celu ręcznego sprawdzenia.

Zwykle wymaga to wielu dni pisania powtarzalnego kodu konfiguracji (notatniki Spark, konfiguracje dbt, skrypty szkoleniowe, DAG Airflow) i ciągłego przełączania się między interfejsami konsoli a edytorami.

W tym ćwiczeniu będziesz programować w parze z agentem przy użyciu pakietu Google Cloud Data Agent Kit (DAK) w środowisku Antigravity IDE. Za pomocą konwersacyjnego języka naturalnego agent pomoże Ci wygenerować notatniki Spark, skompilować projekt dbt, utworzyć pętlę wnioskowania i zorganizować przepływ pracy za pomocą usługi zarządzanej dla Apache Airflow.

Jakie zadania wykonasz

  • Pozyskiwanie dzienników centrum rozliczeniowego z Cloud Storage za pomocą usługi zarządzanej dla Apache Spark (Spark Serverless) do tabeli BigQuery.
  • Usuwaj duplikaty transakcji i normalizuj je za pomocą dbt, aby tworzyć czyste warstwy danych (nieprzetworzone, tymczasowe, wzbogacone).
  • Trenowanie rozproszonego modelu klasyfikacji Random Forest (RandomForestClassifier) w Serverless Spark.
  • Uruchamiaj wnioskowanie wsadowe w przypadku nowych transakcji i zapisuj alerty o wysokim ryzyku bezpośrednio w Cloud Spanner.
  • Orkiestruj, wizualnie konfiguruj i wdrażaj cały potok za pomocą usługi zarządzanej dla Apache Airflow i interaktywnego monitorowania DAG w IDE.

Czego potrzebujesz

  • przeglądarka, np. Chrome;
  • Projekt Google Cloud z włączonym rozliczeniem (w przypadku praktycznych laboratoriów zalecamy użycie nowego, dedykowanego projektu).
  • Podstawowa znajomość SQL, Pythona i PySparka.
  • Antigravity IDE z subskrypcją Google AI Pro (zalecane)

Zasoby utworzone w tym module powinny kosztować mniej niż 5 USD. Po zakończeniu modułu wykonaj instrukcje z sekcji Czyszczenie, aby usunąć udostępnione zasoby.

2. Konfigurowanie środowiska

Aby rozpocząć laboratorium, uruchom skrypt bootstrap. Ten skrypt automatycznie włącza wymagane interfejsy API GCP, tworzy zasobnik Cloud Storage na potrzeby pozyskiwania danych, generuje przykładowe zbiory danych transakcji i katalogów, wczytuje katalogi referencyjne do BigQuery i uruchamia w tle udostępnianie Cloud Spanner i usługi zarządzanej dla Apache Airflow (dawniej Cloud Composer).

Wybierz lub utwórz projekt

Wybierz istniejący projekt lub utwórz nowy w konsoli Google Cloud.

Weryfikacja rozliczeń

Sprawdź, czy w projekcie Google Cloud włączone są płatności. Więcej informacji o tym, jak to zrobić, znajdziesz w tym przewodniku.

Uruchamianie skryptu konfiguracji

Do uruchomienia konfiguracji środowiska użyjesz Google Cloud Shell (lub lokalnej powłoki skonfigurowanej za pomocą Google Cloud CLI).

  1. Otwórz konsolę Google Cloud.
  2. Na pasku narzędzi w prawym górnym rogu kliknij Aktywuj Cloud Shell.

Otwórz Cloud Shell

  1. W terminalu Cloud Shell skonfiguruj aktywny projekt:
gcloud config set project <<YOUR_PROJECT_ID>>
export PROJECT_ID=$(gcloud config get-value project)
  1. Sklonuj repozytorium z codelabem i przejdź do folderu skryptów:
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. Uruchom skrypt konfiguracji wczytywania, aby wdrożyć wszystkie zasoby w us-central1:
chmod +x setup.sh setup_spanner.sh setup_composer.sh
export REGION=us-central1
./setup.sh
  1. Po zakończeniu działania skryptu zobaczysz podsumowanie wskazujące, że zbiór danych BigQuery i zasobnik Cloud Storage są gotowe. W tle będzie kontynuowane udostępnianie usług Cloud Spanner (zajmie to około 2 minut) i Managed Airflow (zajmie to około 20 minut). Postępy możesz śledzić w dowolnym momencie, uruchamiając to polecenie:
tail -f /tmp/spanner_setup.log
tail -f /tmp/composer_setup.log

Otwieranie Antigravity IDE

  1. Pobierz i zainstaluj Antigravity IDE ze strony pobierania Google Antigravity.
  2. Uruchom Antigravity IDE.
  3. Utwórz nowy, pusty folder na komputerze lokalnym (np.o nazwie agentic-data-labs) i otwórz go w środowisku IDE, klikając Otwórz folder. Będzie to Twój lokalny obszar roboczy na potrzeby ćwiczenia.

Konfigurowanie folderu projektu Antigravity IDE

Instalowanie rozszerzenia Data Agent Kit

Rozszerzenie Google Cloud Data Agent Kit zapewnia głęboką integrację z usługami danych Google Cloud bezpośrednio w edytorze, co umożliwia interakcję z BigQuery, Cloud SQL, Cloud Storage i innymi usługami bez przełączania kontekstu.

  1. W środowisku IDE Antigravity kliknij ikonę Rozszerzenia na pasku aktywności po lewej stronie ekranu (wygląda jak 4 kwadraty).
  2. Na pasku wyszukiwania u góry panelu Rozszerzenia wpisz Google Cloud Data Agent Kit.
  3. Znajdź rozszerzenie o nazwie Google Cloud Data Agent Kit opublikowane przez googlecloudtools.
  4. Kliknij przycisk Zainstaluj.
  5. Może pojawić się pytanie: „Czy ufasz wydawcy „googlecloudtools” i jego rozszerzeniom?”. Aby kontynuować, kliknij Zaufaj wydawcom i zainstaluj.

Instalowanie rozszerzenia Data Agent Kit

Po zainstalowaniu w pasku aktywności po lewej stronie środowiska Antigravity IDE pojawi się nowa ikona Google Cloud Data Agent Kit.

  1. Powinna się automatycznie otworzyć strona wprowadzająca o nazwie „Welcome to Google Cloud Data Agent Kit” (Witamy w zestawie narzędzi agenta danych Google Cloud). Jeśli nie jesteś zalogowany(-a) na konto Cloud, postępuj zgodnie z instrukcjami, aby zezwolić na dostęp.
  2. W sekcji Podsumowanie konfiguracji znajdź pole projektu. Kliknij menu i wybierz projekt w chmurze Google Cloud. Ustaw region jako us-central1. Następnie kliknij Skonfiguruj serwery MCP.

Początkowa konfiguracja rozszerzenia Data Agent Kit

  1. Kliknij Skonfiguruj serwery MCP. W panelu Konfiguracja MCP włącz te zdalne serwery MCP:
    • BigQuery
    • Spanner
    • Notatniki

Następnie kliknij Rozpocznij.

Konfigurowanie serwerów MCP

Poznaj opcje konfiguracji

Po zakończeniu konfiguracji otworzy się strona „Rozpocznij korzystanie z zestawu narzędzi Google Cloud Data Agent”.

  1. W sekcji „Konfiguracja” kliknij Rozpocznij.
  2. Otworzy się panel Konfiguracja pakietu agenta danych. Przejrzyj karty:
    • Projekt i region: sprawdź wybrany identyfikator projektu i upewnij się, że skrypt konfiguracji włączył wszystkie wymagane interfejsy API (Compute Engine, Cloud Storage, BigQuery, Spanner itp.).
    • BigQuery: skonfiguruj domyślną lokalizację zapytań BigQuery. Użyj regionu us-central1.
    • Konfigurowanie serwerów MCP: wyświetl włączone serwery MCP (BigQuery, Notebooks, Spanner itp.), które umożliwiają agentom AI bezpieczną interakcję z Twoimi danymi.
    • Umiejętności: poznaj gotowe umiejętności, które zapewniają agentom specjalistyczne możliwości wykonywania złożonych zadań związanych z danymi.

Panel ustawień Data Agent Kit

Podsumowanie sekcji: uruchomiono skrypt początkowy, aby utworzyć zasoby GCS i BigQuery, a Spanner i Airflow są tworzone w tle. Następnie otworzono projekt w IDE Antigravity i aktywowano rozszerzenie Google Cloud Data Agent Kit. Możesz teraz napisać pierwszy notatnik.

3. Przetwarzanie nieprzetworzonych logów za pomocą Spark Serverless

W tej sekcji zaimportujesz do jeziora danych nieprzetworzone dzienniki transakcji w formacie JSON. Usługa zarządzana dla Apache Spark (Spark Serverless) łączy się bezpośrednio z natywnym miejscem na dane BigQuery. Do zarządzania danymi tabelarycznymi oraz umożliwiania bezpośredniego wykonywania zapytań i analizowania danych będziesz używać standardowego oprogramowania sprzęgającego BigQuery.

Poznawanie wstępnie skonfigurowanego środowiska wykonawczego Serverless Spark

Przed wykonaniem kodu Sparka sprawdź szablon środowiska wykonawczego bezserwerowego, który został wstępnie skonfigurowany przez skrypt konfiguracji. Ten szablon określa docelowe środowisko wykonawcze backendu i zawiera niezbędne zależności łącznika.

  1. Na pasku działań IDE otwórz panel Google Cloud Data Agent Kit.
  2. Rozwiń menu Apache Spark, a następnie Serverless.
  3. Kliknij prawym przyciskiem myszy fraud-pipeline-runtime i wybierz Profil, aby otworzyć widok konfiguracji w edytorze.
  4. Na karcie Profil przewiń w dół i rozwiń sekcję Właściwości, aby sprawdzić niestandardowe zależności dołączone do środowiska:
    • spark.jars: zawiera gs://spark-lib/spanner/spark-3.5-spanner-1.4.0.jar, który używa konektora Spark Spanner, aby umożliwić zadaniom Spark zapisywanie wyników wnioskowania bezpośrednio w Cloud Spanner w dalszej części modułu. (Uwaga: usługa Dataproc Serverless domyślnie zawiera oprogramowanie sprzęgające Spark BigQuery Google Cloud, które nie wymaga dodatkowej konfiguracji pliku JAR do odczytywania i zapisywania tabel BigQuery).

Przeglądanie właściwości środowiska wykonawczego Serverless Spark

  1. Po lewej stronie zobaczysz kartę Sesje interaktywne. Jest ona obecnie pusta, ponieważ nie wykonano jeszcze żadnego kodu. Gdy w następnym kroku uruchomisz notatnik, dynamicznie zostanie udostępniona sesja obliczeniowa bez serwera, która pojawi się tutaj.

Pozyskiwanie danych za pomocą pakietu Data Agent Kit

Zamiast ręcznie konfigurować sesję Spark lub pisać od zera skrypty wczytywania PySpark, będziesz programować w parach z agentem przy użyciu Data Agent Kit.

  1. Otwórz panel Czat z agentem, klikając ikonę Przełącz agenta na pasku narzędzi w prawym górnym rogu.
  2. Wklej w czacie ten prompt (pamiętaj, aby zastąpić ${PROJECT_ID} identyfikatorem projektu 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. Jeśli agent poprosi o zezwolenie na wykonanie poleceń weryfikacji w tle (np. „Zezwolić na uruchomienie tego polecenia?”), sprawdź proponowane polecenie i kliknij Tak, zezwól tym razem (lub Tak, i zawsze zezwalaj).
  2. Gdy agent zakończy generowanie pliku, kliknij niebieski przycisk Zaakceptuj wszystko (lub ikonę zaznaczenia) u dołu panelu czatu, aby zapisać notebooks/01_ingestion.ipynb w obszarze roboczym.

Agent generujący notatnik do pozyskiwania danych

Sprawdzanie i uruchamianie notatnika

  1. Otwórz nowo wygenerowany plik notebooks/01_ingestion.ipynb w IDE.
  2. Sprawdź kod PySpark dla logiki zapisu oprogramowania sprzęgającego BigQuery.
  3. Na pasku narzędzi notatnika w IDE kliknij Uruchom wszystko.
  4. Jeśli uruchamiasz zdalny notatnik Spark po raz pierwszy, IDE może wyświetlić prośbę o zainstalowanie lokalnych zależności. Jeśli pojawi się odpowiedni komunikat, kliknij Install dependencies for Remote Spark Kernels (Zainstaluj zależności dla zdalnych jąder Sparka) i potwierdź okna instalacji, a następnie ponownie kliknij Run All (Uruchom wszystko).
  5. W menu Select Kernel (Wybierz jądro) wybierz Remote Spark Kernels (Zdalne jądra Sparka) –> fraud-pipeline-runtime on Serverless Spark (fraud-pipeline-runtime w Serverless Spark). (Wskazówka: jeśli nie widzisz na liście wstępnie skonfigurowanego szablonu środowiska wykonawczego, kliknij ikonę odświeżania w prawym górnym rogu menu wyboru jądra, aby ponownie załadować dostępne zdalne jądra).
  6. Spójrz na pasek stanu w lewym dolnym rogu edytora. Zobaczysz Connecting to kernel: fraud-pipeline-runtime on Serverless Spark.... Ponieważ jest to pierwsze uruchomienie środowiska wykonawczego Spark Serverless, jego udostępnienie i uruchomienie zajmie kilka minut.
  7. Gdy jądro zakończy łączenie, notatnik automatycznie zacznie wykonywać wszystkie komórki po kolei, aby przetworzyć surowe dzienniki transakcji w zbiór danych BigQuery.

Weryfikacja

Po zakończeniu wykonania sprawdź katalog Data Agent Kit, aby potwierdzić utworzenie tabeli:

Weryfikowanie tabeli Raw w eksploratorze katalogu

  1. Na pasku działań IDE otwórz panel Google Cloud Data Agent Kit.
  2. Rozwiń sekcję KATALOG.
  3. Rozwiń identyfikator projektu.
  4. Rozwiń BigQuery.
  5. Rozwiń zbiór danych transactions_dataset_evals.
  6. Kliknij tabelę raw_transactions, aby otworzyć jej widok szczegółowy w głównym edytorze.
  7. W panelu nawigacyjnym po lewej stronie otwórz karty Dane, Schemat i Szczegóły, aby sprawdzić przetworzone rekordy i metadane.

Podsumowanie sekcji: w agencie czatu użyto języka naturalnego do wygenerowania kompletnego obciążenia Spark Serverless. Następnie wykonujesz go, aby przetworzyć nieustrukturyzowane logi JSON w tabeli BigQuery (surowej).

4. Usuwanie duplikatów i normalizacja za pomocą dbt

Przed trenowaniem modelu ML wymusisz jakość danych, usuwając zduplikowane dzienniki przesyłania strumieniowego, wyodrębniając nieprawidłowe rekordy (np. puste identyfikatory transakcji) i łącząc dane wymiarowe (płatników i odbiorców płatności). Ten proces wymaga idempotentnych i niezawodnych przekształceń SQL, dlatego doskonale sprawdza się w nim dbt (data build tool).

Tworzenie szkieletu potoku dbt

Użyj agenta, aby wygenerować projekt dbt na podstawie zbioru danych BigQuery:

  1. Wróć do panelu Czat z agentem.
  2. Aby wygenerować projekt dbt, wykonaj te instrukcje:
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. W głównym panelu edytora agent wyświetli artefakt Plan wdrożenia. Sprawdź proponowaną strukturę pliku i logikę SQL.
  2. Kliknij Dalej (a potem Zaakceptuj wszystko), aby zezwolić agentowi na wygenerowanie plików w obszarze roboczym.

Plan wdrożenia z przyciskiem Dalej

  1. Po zakończeniu generowania agent wyświetli przewodnik podsumowujący nowe komponenty. Jeśli pojawi się taka prośba, zaakceptuj wszystkie zmiany.

Akceptowanie wszystkich wygenerowanych plików w panelu Google Chat

Tworzenie i testowanie

Chociaż agent automatycznie uruchomił dbt compile, aby upewnić się, że wygenerowany kod SQL jest poprawny pod względem składni, teraz zmaterializujesz te widoki i tabele w BigQuery i uruchomisz testy jakości danych na potrzeby lokalnej weryfikacji. (Uwaga: w dalszej części tego modułu zautomatyzujesz ten krok dbt w ramach kompleksowego DAG Airflow).

  1. Na pasku działań po lewej stronie kliknij ikonę Eksplorator (lub naciśnij Cmd/Ctrl+Shift+E).
  2. Rozwiń dbt_project –> models, aby sprawdzić wygenerowane modele SQL. Kliknij enriched_transactions.sql, aby otworzyć i sprawdzić logikę funkcji przekształcania i wykrywania oszustw w edytorze.
  3. W Eksploratorze plików kliknij prawym przyciskiem myszy folder dbt_project i wybierz Open in Integrated Terminal (Otwórz w zintegrowanym terminalu). Spowoduje to automatyczne otwarcie panelu terminala ustawionego bezpośrednio w wymaganym dbt_project katalogu roboczym.
  4. Jeśli nie masz jeszcze zainstalowanego pakietu dbt, utwórz środowisko wirtualne poza dbt_project/ (w katalogu głównym lub obszarze roboczym) i zainstaluj adapter BigQuery:
python3 -m venv ~/.venv/dbt
source ~/.venv/dbt/bin/activate
pip install dbt-bigquery
  1. Uruchom modele dbt i powiązane z nimi testy jakości danych:
dbt build
  1. Obserwuj dane wyjściowe terminala. dbt skompiluje SQL, zmaterializuje tabele tymczasowe i wzbogacone w BigQuery oraz wykona testy danych.

Tworzenie i testowanie projektu dbt w zintegrowanym terminalu

  1. Po zakończeniu kompilacji zamknij panel terminala, aby zwolnić miejsce na ekranie na potrzeby pozostałych kroków.

Podsumowanie sekcji: wygenerowano projekt dbt za pomocą agenta, przeprowadzono testy jakości danych i przekształcono surowe rekordy w tabele BigQuery z danymi tymczasowymi i wzbogaconymi.

5. Trenowanie rozproszonego modelu wykrywania oszustw za pomocą algorytmu Random Forest

Po zmaterializowaniu wzbogaconych transakcji w BigQuery utworzysz model uczenia maszynowego do klasyfikowania oszukańczych zdarzeń. Random Forest to metoda uczenia zespołowego, która dobrze sprawdza się w przypadku tabelarycznych danych klasyfikacji. Uruchomienie RandomForestClassifier w usłudze Spark Serverless powoduje rozdzielenie trenowania modelu między węzły robocze bez konieczności zarządzania infrastrukturą.

W tym kroku użyjesz agenta do wygenerowania potoku trenowania Spark ML.

Generowanie notatnika do trenowania ML

  1. Otwórz panel Czat z agentem.
  2. Podaj ten prompt, aby zaprojektować sekwencję trenowania modelu (pamiętaj, aby zastąpić ${PROJECT_ID} identyfikatorem projektu):
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. Sprawdź plan agenta lub wygenerowany kod i kliknij Dalej / Zaakceptuj wszystko, aby zapisać notebooks/02_training.ipynb w przestrzeni roboczej.

Agent generujący notatnik do trenowania

Sprawdzanie i uruchamianie notatnika

  1. Otwórz plik notebooks/02_training.ipynb w edytorze.
  2. Sprawdź etapy potoku PySpark ML pod kątem kodowania funkcji, składania wektorów i logiki klasyfikacji Random Forest.
  3. Na pasku narzędzi notatnika w IDE kliknij Uruchom wszystko.
  4. Gdy otworzy się menu Select Kernel (Wybierz jądro), wybierz fraud-pipeline-runtime on Serverless Spark (fraud-pipeline-runtime w Serverless Spark).

Wybieranie jądra bezserwerowej usługi Spark dla notatnika trenowania

Weryfikacja

Po zakończeniu wykonania sprawdź, czy model został prawidłowo wytrenowany i wyeksportowany:

  1. Sprawdź wyniki komórki oceny u dołu notatnika, aby zweryfikować podany wynik obszaru pod krzywą ROC (AUC).
  2. Aby sprawdzić, czy artefakty modelu zostały zapisane w GCS, rozwiń panel eksploratora STORAGE na pasku bocznym Data Agent Kit.
  3. Znajdź zasobnik kończący się na -models (powiązany z aktywnym identyfikatorem projektu), rozwiń go i sprawdź, czy istnieją katalog fraud_model i jego etapy potoku.

Sprawdź, czy model został zapisany w GCS

Podsumowanie sekcji: za pomocą agenta utworzono potok trenowania ML w PySpark, wytrenowano model Random Forest na wzbogaconej tabeli BigQuery i wyeksportowano model do Cloud Storage.

6. Wnioskowanie wsadowe i zapisywanie w Cloud Spanner

Za pomocą wytrenowanego modelu predykcyjnego przechowywanego w Cloud Storage przeprowadzisz wnioskowanie wsadowe na nowych transakcjach przepływających przez BigQuery. Transakcje o wysokim ryzyku muszą być kierowane do systemu operacyjnego, aby zespół ds. zgodności mógł je sprawdzić. Cloud Spanner zapewnia skalowalną bazę danych transakcyjnych dla tej kolejki opinii.

Generowanie notatnika wnioskowania zbiorczego

Użyj agenta, aby utworzyć notatnik wnioskowania łączący BigQuery, Cloud Storage i Cloud Spanner:

  1. Otwórz panel Czat z agentem.
  2. Wpisz ten 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. Zaakceptuj wygenerowany notatnik, aby zapisać notebooks/03_inference.ipynb w obszarze roboczym.

Agent generujący notatnik wnioskowania

Sprawdzanie i uruchamianie notatnika

  1. Otwórz nowo wygenerowany plik notebooks/03_inference.ipynb w edytorze.
  2. Sprawdź sekwencję wnioskowania PySpark:
    • Zależności: szablon środowiska wykonawczego bezserwerowego udostępnia wymagane zależności cloud-spanner JAR do wykonywania Sparka.
    • Formatowanie danych: skrypt usuwa złożone kolumny wektorowe Spark ML (takie jak surowe cechy i prawdopodobieństwa) przed zapisaniem, aby dopasować je do schematu tabeli Spannera.
    • Łącznik Spanner: zapisuje oznaczone wiersze za pomocą .format("cloud-spanner"), aby bezpośrednio dołączyć je do kolejki opinii.
  3. Na pasku narzędzi notatnika w IDE kliknij Uruchom wszystko.
  4. Gdy pojawi się prośba o wybranie jądra, kliknij fraud-pipeline-runtime on Serverless Spark (fraud-pipeline-runtime w Serverless Spark).

Weryfikacja

Gdy notatnik wnioskowania zakończy przetwarzanie, możesz wysyłać zapytania do operacyjnej bazy danych Spanner bezpośrednio w środowisku IDE:

  1. Na pasku działań IDE otwórz panel Google Cloud Data Agent Kit.
  2. Rozwiń sekcję KATALOG.
  3. Rozwiń identyfikator projektu, a następnie Spanner.
  4. Kliknij kolejno cymbal-fraud –> fraud-db –> Tables –> SparkEvalFraudReviewQueue.
  5. Kliknij tabelę prawym przyciskiem myszy i wybierz Zapytanie do tabeli, a następnie wykonaj zapytanie:
SELECT *
FROM `SparkEvalFraudReviewQueue`
LIMIT 100;
  1. W panelu Wyniki zapytania poniżej powinny być widoczne nowo wstawione wiersze reprezentujące transakcje wysokiego ryzyka oznaczone do ręcznego sprawdzenia.

Weryfikowanie wierszy w Cloud Spanner

Podsumowanie sekcji: za pomocą agenta utworzono notatnik wnioskowania wsadowego, oceniono nieoznaczone rekordy BigQuery za pomocą wytrenowanego modelu i zapisano transakcje o wysokim ryzyku bezpośrednio w Cloud Spanner.

7. Tworzenie szkieletu i aranżowanie za pomocą usługi Managed Airflow

Twój potok składa się obecnie z osobnych kroków: notatnika do pozyskiwania danych, projektu przekształceń dbt i notatnika wnioskowania wsadowego. Aby przygotować go do użytku produkcyjnego, połączysz go w wykres zależności zaplanowanych.

Usługa zarządzana dla Apache Airflow (wcześniej Cloud Composer) udostępnia zarządzany silnik orkiestracji tego przepływu pracy. Zestaw Data Agent Kit zawiera funkcję Orchestration Pipelines, która tłumaczy deklaratywne definicje potoków YAML bezpośrednio na DAG-i Airflow.

Definiowanie potoku

Użyj agenta, aby wygenerować konfigurację potoku administracji:

  1. W Agent Chat wpisz ten prompt (pamiętaj, aby zastąpić ${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.

Sprawdź konfigurację DAG-a

Orchestrator pakietu Data Agent Kit używa deklaratywnych konfiguracji YAML do definiowania i wdrażania potoków w Apache Airflow, co umożliwia kontrolowanie wersji definicji i wdrażanie ich za pomocą CI/CD.

W panelu Eksploratora IDE sprawdź 2 pliki potoku wygenerowane przez agenta w katalogu głównym obszaru roboczego:

  1. deployment.yaml: otwórz ten plik. Będzie to rejestr środowisk. Mapuje logiczny dev potok na cymbal-airflow środowisko, ustawia region wykonania (us-central1) i definiuje artifact_storage zasobnik, w którym są przechowywane skompilowane DAG-i i zależności.
  2. fraud_analysis_pipeline.yaml: otwórz ten plik. Określa on wykres wykonania. Określa harmonogram wyzwalania (interval: '0 0 * * *') i sekwencjonuje 3 kroki w bloku actions:
    • działanie związane z przetwarzaniem notebook w przypadku 01_ingestion.ipynb uruchomionego w Dataproc Serverless;
    • Działanie przekształcenia pipeline kierowane na katalog dbt_project z zależnością dependsOn wskazującą krok pozyskiwania.
    • Działanie wnioskowania notebook dla 03_inference.ipynb z zależnością dependsOn wskazującą krok dbt, zawierające właściwość JAR Spanner.
  3. Agent podsumuje też wygenerowane artefakty na karcie Przewodnik w panelu edytora, przedstawiając wykonane konfiguracje i weryfikacje.

Interaktywna konfiguracja DAG

Zestaw Data Agent Kit renderuje konfigurację potoku jako interaktywny wykres wizualny, który umożliwia sprawdzanie i edytowanie właściwości DAG-ów Airflow.

  1. Na pasku działań IDE otwórz panel Google Cloud Data Agent Kit.
  2. W sekcji DATA ENGINEERING rozwiń Orchestration Pipelines.
  3. Kliknij fraud_analysis_pipeline.yaml, aby otworzyć wizualny obszar DAG w głównym edytorze.

Obszar roboczy wizualizacji DAG-a administracji

  1. U góry kliknij węzeł Schedule trigger. Po prawej stronie otworzy się wysuwane menu konfiguracji z przeanalizowanym ciągiem Cron (0 0 * * *), w którym możesz dostosować parametry takie jak uzupełnianie wsteczne i nadrobienie zaległości.
  2. Kliknij węzeł zadania notatnika (np. krok pozyskiwania lub wnioskowania). Wyskakujące okienko zostanie zaktualizowane, aby wyświetlać konkretne mapowania wykonania Dataproc Serverless i właściwości łącznika.
  3. Zwróć uwagę na hiperlink do nazwy pliku notatnika (np. 01_ingestion.ipynb) w bloku węzła. Kliknięcie go spowoduje otwarcie notatnika bezpośrednio w edytorze.
  4. Na pasku bocznym po lewej stronie w sekcji Orchestration Pipelines (Potoki orkiestracji) kliknij Deployment configuration. Ten widok pokazuje docelowy klaster środowiska dev i artefakty wyjściowego zasobnika GCS.

Podsumowanie sekcji: za pomocą agenta wygenerowano konfigurację potoku orkiestracji, definiując zależności między zadaniami dotyczącymi pozyskiwania, dbt i wnioskowania na interaktywnej platformie wizualnej.

8. Wdrażanie, wykonywanie i monitorowanie

Po zdefiniowaniu DAG-a lokalnie połączysz się ze środowiskiem Managed Airflow udostępnionym podczas konfiguracji i wdrożysz potok.

Konfigurowanie usługi zarządzanej dla Apache Airflow

Przed wdrożeniem skonfiguruj połączenie z harmonogramem w ustawieniach Data Agent Kit, aby rozszerzenie było kierowane na środowisko Managed Airflow:

  1. Na pasku działań IDE otwórz panel Google Cloud Data Agent Kit.
  2. W sekcji SETTINGS kliknij Ustawienia.
  3. W menu po lewej stronie wybierz Harmonogram.
  4. Skonfiguruj ustawienia:
    • Identyfikator projektu: wybierz aktywny identyfikator projektu.
    • Region: wybierz us-central1.
    • Środowisko: kliknij cymbal-airflow.
  5. Kliknij Zapisz.

Ustawienia usługi zarządzanej dla Apache Airflow

Wdrażanie DAG-a

Skonfigurowany potok wdrożysz teraz bezpośrednio w środowisku Managed Airflow z poziomu wizualnego obszaru roboczego:

  1. Na pasku bocznym Google Cloud Data Agent Kit rozwiń kolejno DATA ENGINEERING > Orchestration Pipelines i kliknij fraud_analysis_pipeline.yaml, aby otworzyć wizualne miejsce robocze DAG.
  2. W prawym górnym rogu paska narzędzi obszaru roboczego kliknij niebieski przycisk Uruchom potok.
  3. W menu środowiska kliknij dev.
  4. Obserwuj powiadomienie o postępach w obszarze stanu u dołu (Running pipeline: Building pipeline locally...). Rozszerzenie automatycznie skompiluje DAG, spakuje notatnik i zasoby dbt oraz prześle je do zasobnika GCS w środowisku Managed Airflow (zajmie to około 3–4 minut).

Wdrażanie potoku z poziomu wizualnego obszaru roboczego

Monitorowanie uruchomienia

Po zakończeniu kompilacji lokalnej i wyświetleniu powiadomienia potwierdzającego Triggered a new run for pipeline... successfully monitoruj wykonanie na żywo:

  1. Na pasku bocznym Google Cloud Data Agent Kit rozwiń kolejno DATA ENGINEERING i Orchestration Pipelines.
  2. Kliknij Zarządzanie potokami.
  3. W tabeli Zarządzanie potokami kliknij fraud_analysis_pipeline, aby otworzyć historię wykonania.

Omówienie zarządzania potokami

  1. W widoku Historia wykonywania wybierz z kalendarza aktywne uruchomienie.
  2. W miarę postępu wykonywania poszczególnych zadań potoku (pozyskiwanie, przekształcanie dbt i wnioskowanie) aktualizują się wskaźniki stanu i wypełniają się czasy trwania zadań. Kliknij dowolne zadanie, aby sprawdzić jego dane wyjściowe wykonania na żywo i dzienniki DAG Airflow.

Historia wykonywania potoku na żywo i szczegóły zadań

Podsumowanie sekcji: skonfigurowano połączenie z harmonogramem Airflow, wdrożono kompleksową potok analityczny w usłudze Managed Airflow i monitorowano wykonanie na żywo, weryfikując system od surowych logów po końcowe prognozy Cloud Spanner.

9. Czyszczenie danych

Aby uniknąć obciążania projektu w Google Cloud bieżącymi opłatami za zasoby użyte w tym laboratorium, zamknij środowisko za pomocą automatycznego skryptu.

  1. W panelu Terminal (lub w Cloud Shell) przejdź do katalogu skryptów i wykonaj to polecenie:
cd ~/devrel-demos/codelabs/agentic-data-labs/data-science/scripts
chmod +x teardown.sh
./teardown.sh
  1. Skrypt wyświetli listę wszystkich zasobów, które planuje usunąć, i poprosi o potwierdzenie:
    • Zarządzane środowisko Airflow (cymbal-airflow)
    • Instancja Cloud Spanner (cymbal-fraud)
    • Zbiór danych BigQuery (transactions_dataset_evals)
    • Zasobniki Cloud Storage (gs://${PROJECT_ID}-fin-clearing-raw i gs://${PROJECT_ID}-models)
    • Konto usługi zasobu roboczego (composer-worker-sa)
  2. Wpisz y, aby potwierdzić. Skrypt zamykający usunie wszystkie udostępnione usługi GCP i zwolni miejsce na pliki lokalne.

10. Gratulacje!

Udało Ci się utworzyć kompleksowy potok wykrywania oszustw obejmujący Cloud Storage, BigQuery, usługę zarządzaną dla Apache Spark (Spark Serverless), dbt, Cloud Spanner i usługę zarządzaną dla Apache Airflow. Pracujesz w parze z zestawem narzędzi Google Cloud Data Agent Kit w środowisku Antigravity IDE.

Co udało Ci się osiągnąć

  1. 📥 Pozyskane surowe logi transakcji w tabeli BigQuery za pomocą usługi zarządzanej dla Apache Spark i zestawu Data Agent Kit.
  2. 🧹 Usuwanie duplikatów i normalizowanie danych przez utworzenie projektu dbt z testami jakości danych.
  3. 🤖 Wytrenowano rozproszony model lasu losowego za pomocą RandomForestClassifier i wyeksportowano wytrenowany model do Cloud Storage.
  4. ⚡ Przeprowadzono wnioskowanie wsadowe na podstawie przychodzących transakcji i przekierowano rekordy o wysokim ryzyku do Cloud Spanner w celu sprawdzenia.
  5. 🔄 Orkiestracja, wdrażanie i monitorowanie przepływu pracy jako zaplanowanego DAG-a Airflow przy użyciu usługi zarządzanej dla Apache Airflow i wizualnych narzędzi IDE do zarządzania DAG-ami.

Kluczowych pojęć

Pomysł

Czego się dowiedziałeś(-aś)

Data Agent Kit

Programowanie w parach w IDE z użyciem języka naturalnego do generowania notatników PySpark, konfigurowania modeli dbt i definiowania DAG-ów Airflow.

BigQuery

Skalowalna pamięć tabelaryczna do analitycznego SQL, transformacji dbt i trenowania ML

Spark Serverless

Bezserwerowe wykonywanie rozproszonego wczytywania danych PySpark i trenowania modelu ML Random Forest

Cloud Spanner Connector

Zapisywanie prognoz wnioskowania wsadowego Spark bezpośrednio w kolejkach opinii w operacyjnej bazie danych

Deklaracje DAG w YAML

Deklaratywne definicje potoków renderowane jako interaktywne wykresy wizualne Airflow w środowisku IDE

Wizualne zarządzanie DAG-iem

Sprawdzanie zależności potoku, wdrażanie w Managed Airflow i monitorowanie historii wykonywania zadań na żywo w IDE

Dalsze kroki