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).
- Otwórz konsolę Google Cloud.
- Na pasku narzędzi w prawym górnym rogu kliknij Aktywuj Cloud Shell.

- W terminalu Cloud Shell skonfiguruj aktywny projekt:
gcloud config set project <<YOUR_PROJECT_ID>>
export PROJECT_ID=$(gcloud config get-value project)
- 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
- 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
- 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
- Pobierz i zainstaluj Antigravity IDE ze strony pobierania Google Antigravity.
- Uruchom Antigravity IDE.
- 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.

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

Po zainstalowaniu w pasku aktywności po lewej stronie środowiska Antigravity IDE pojawi się nowa ikona Google Cloud Data Agent Kit.
- 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.
- 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.

- Kliknij Skonfiguruj serwery MCP. W panelu Konfiguracja MCP włącz te zdalne serwery MCP:
- BigQuery
- Spanner
- Notatniki
Następnie kliknij Rozpocznij.

Poznaj opcje konfiguracji
Po zakończeniu konfiguracji otworzy się strona „Rozpocznij korzystanie z zestawu narzędzi Google Cloud Data Agent”.
- W sekcji „Konfiguracja” kliknij Rozpocznij.
- 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.

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.
- Na pasku działań IDE otwórz panel Google Cloud Data Agent Kit.
- Rozwiń menu Apache Spark, a następnie Serverless.
- Kliknij prawym przyciskiem myszy
fraud-pipeline-runtimei wybierz Profil, aby otworzyć widok konfiguracji w edytorze. - 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: zawierags://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).

- 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.
- Otwórz panel Czat z agentem, klikając ikonę Przełącz agenta na pasku narzędzi w prawym górnym rogu.
- 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.
- 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).
- Gdy agent zakończy generowanie pliku, kliknij niebieski przycisk Zaakceptuj wszystko (lub ikonę zaznaczenia) u dołu panelu czatu, aby zapisać
notebooks/01_ingestion.ipynbw obszarze roboczym.

Sprawdzanie i uruchamianie notatnika
- Otwórz nowo wygenerowany plik
notebooks/01_ingestion.ipynbw IDE. - Sprawdź kod PySpark dla logiki zapisu oprogramowania sprzęgającego BigQuery.
- Na pasku narzędzi notatnika w IDE kliknij Uruchom wszystko.
- 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).
- 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).
- 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. - 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:

- Na pasku działań IDE otwórz panel Google Cloud Data Agent Kit.
- Rozwiń sekcję KATALOG.
- Rozwiń identyfikator projektu.
- Rozwiń BigQuery.
- Rozwiń zbiór danych
transactions_dataset_evals. - Kliknij tabelę
raw_transactions, aby otworzyć jej widok szczegółowy w głównym edytorze. - 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:
- Wróć do panelu Czat z agentem.
- 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.
- W głównym panelu edytora agent wyświetli artefakt Plan wdrożenia. Sprawdź proponowaną strukturę pliku i logikę SQL.
- Kliknij Dalej (a potem Zaakceptuj wszystko), aby zezwolić agentowi na wygenerowanie plików w obszarze roboczym.

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

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).
- Na pasku działań po lewej stronie kliknij ikonę Eksplorator (lub naciśnij
Cmd/Ctrl+Shift+E). - Rozwiń
dbt_project–>models, aby sprawdzić wygenerowane modele SQL. Kliknijenriched_transactions.sql, aby otworzyć i sprawdzić logikę funkcji przekształcania i wykrywania oszustw w edytorze. - W Eksploratorze plików kliknij prawym przyciskiem myszy folder
dbt_projecti wybierz Open in Integrated Terminal (Otwórz w zintegrowanym terminalu). Spowoduje to automatyczne otwarcie panelu terminala ustawionego bezpośrednio w wymaganymdbt_projectkatalogu roboczym. - Jeśli nie masz jeszcze zainstalowanego pakietu
dbt, utwórz środowisko wirtualne pozadbt_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
- Uruchom modele dbt i powiązane z nimi testy jakości danych:
dbt build
- Obserwuj dane wyjściowe terminala. dbt skompiluje SQL, zmaterializuje tabele tymczasowe i wzbogacone w BigQuery oraz wykona testy danych.

- 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
- Otwórz panel Czat z agentem.
- 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.
- Sprawdź plan agenta lub wygenerowany kod i kliknij Dalej / Zaakceptuj wszystko, aby zapisać
notebooks/02_training.ipynbw przestrzeni roboczej.

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

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

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:
- Otwórz panel Czat z agentem.
- 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.
- Zaakceptuj wygenerowany notatnik, aby zapisać
notebooks/03_inference.ipynbw obszarze roboczym.

Sprawdzanie i uruchamianie notatnika
- Otwórz nowo wygenerowany plik
notebooks/03_inference.ipynbw edytorze. - Sprawdź sekwencję wnioskowania PySpark:
- Zależności: szablon środowiska wykonawczego bezserwerowego udostępnia wymagane zależności
cloud-spannerJAR 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.
- Zależności: szablon środowiska wykonawczego bezserwerowego udostępnia wymagane zależności
- Na pasku narzędzi notatnika w IDE kliknij Uruchom wszystko.
- 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:
- Na pasku działań IDE otwórz panel Google Cloud Data Agent Kit.
- Rozwiń sekcję KATALOG.
- Rozwiń identyfikator projektu, a następnie Spanner.
- Kliknij kolejno
cymbal-fraud–>fraud-db–>Tables–>SparkEvalFraudReviewQueue. - Kliknij tabelę prawym przyciskiem myszy i wybierz Zapytanie do tabeli, a następnie wykonaj zapytanie:
SELECT *
FROM `SparkEvalFraudReviewQueue`
LIMIT 100;
- W panelu Wyniki zapytania poniżej powinny być widoczne nowo wstawione wiersze reprezentujące transakcje wysokiego ryzyka oznaczone do ręcznego sprawdzenia.

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:
- 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:
deployment.yaml: otwórz ten plik. Będzie to rejestr środowisk. Mapuje logicznydevpotok nacymbal-airflowśrodowisko, ustawia region wykonania (us-central1) i definiujeartifact_storagezasobnik, w którym są przechowywane skompilowane DAG-i i zależności.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 blokuactions:- działanie związane z przetwarzaniem
notebookw przypadku01_ingestion.ipynburuchomionego w Dataproc Serverless; - Działanie przekształcenia
pipelinekierowane na katalogdbt_projectz zależnościądependsOnwskazującą krok pozyskiwania. - Działanie wnioskowania
notebookdla03_inference.ipynbz zależnościądependsOnwskazującą krok dbt, zawierające właściwość JAR Spanner.
- działanie związane z przetwarzaniem
- 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.
- Na pasku działań IDE otwórz panel Google Cloud Data Agent Kit.
- W sekcji
DATA ENGINEERINGrozwińOrchestration Pipelines. - Kliknij
fraud_analysis_pipeline.yaml, aby otworzyć wizualny obszar DAG w głównym edytorze.

- 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. - 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.
- 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. - Na pasku bocznym po lewej stronie w sekcji Orchestration Pipelines (Potoki orkiestracji) kliknij
Deployment configuration. Ten widok pokazuje docelowy klaster środowiskadevi 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:
- Na pasku działań IDE otwórz panel Google Cloud Data Agent Kit.
- W sekcji
SETTINGSkliknij Ustawienia. - W menu po lewej stronie wybierz Harmonogram.
- Skonfiguruj ustawienia:
- Identyfikator projektu: wybierz aktywny identyfikator projektu.
- Region: wybierz
us-central1. - Środowisko: kliknij
cymbal-airflow.
- Kliknij Zapisz.

Wdrażanie DAG-a
Skonfigurowany potok wdrożysz teraz bezpośrednio w środowisku Managed Airflow z poziomu wizualnego obszaru roboczego:
- Na pasku bocznym Google Cloud Data Agent Kit rozwiń kolejno
DATA ENGINEERING>Orchestration Pipelinesi kliknijfraud_analysis_pipeline.yaml, aby otworzyć wizualne miejsce robocze DAG. - W prawym górnym rogu paska narzędzi obszaru roboczego kliknij niebieski przycisk Uruchom potok.
- W menu środowiska kliknij
dev. - 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).

Monitorowanie uruchomienia
Po zakończeniu kompilacji lokalnej i wyświetleniu powiadomienia potwierdzającego Triggered a new run for pipeline... successfully monitoruj wykonanie na żywo:
- Na pasku bocznym Google Cloud Data Agent Kit rozwiń kolejno
DATA ENGINEERINGiOrchestration Pipelines. - Kliknij Zarządzanie potokami.
- W tabeli Zarządzanie potokami kliknij
fraud_analysis_pipeline, aby otworzyć historię wykonania.

- W widoku Historia wykonywania wybierz z kalendarza aktywne uruchomienie.
- 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.

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.
- 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
- 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-rawigs://${PROJECT_ID}-models) - Konto usługi zasobu roboczego (
composer-worker-sa)
- Zarządzane środowisko Airflow (
- 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ąć
- 📥 Pozyskane surowe logi transakcji w tabeli BigQuery za pomocą usługi zarządzanej dla Apache Spark i zestawu Data Agent Kit.
- 🧹 Usuwanie duplikatów i normalizowanie danych przez utworzenie projektu dbt z testami jakości danych.
- 🤖 Wytrenowano rozproszony model lasu losowego za pomocą
RandomForestClassifieri wyeksportowano wytrenowany model do Cloud Storage. - ⚡ Przeprowadzono wnioskowanie wsadowe na podstawie przychodzących transakcji i przekierowano rekordy o wysokim ryzyku do Cloud Spanner w celu sprawdzenia.
- 🔄 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ś) |
Programowanie w parach w IDE z użyciem języka naturalnego do generowania notatników PySpark, konfigurowania modeli dbt i definiowania DAG-ów Airflow. | |
Skalowalna pamięć tabelaryczna do analitycznego SQL, transformacji dbt i trenowania ML | |
Bezserwerowe wykonywanie rozproszonego wczytywania danych PySpark i trenowania modelu ML Random Forest | |
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
- Zapoznaj się z dokumentacją Google Cloud Data Agent Kit
- Więcej informacji o usłudze zarządzanej dla Apache Spark
- Więcej informacji o usłudze zarządzanej dla Apache Airflow
- Twórz własne potoki wielu usług za pomocą środowiska IDE Antigravity