צינור לזיהוי הונאות באמצעות Data Agent Kit ו-Antigravity IDE

1. מבוא

נניח שאתם מדעני נתונים ב-Cymbal Financial, ספק שירותי תשלומים עם נפח עסקאות גבוה. התרחש גל של עיכובים בסליקה, וצוות התאימות חושד בהונאה מתואמת. אתם צריכים ליצור צינור להעברה של יומני עסקאות גולמיים של מרכז הסליקה, לנקות את הנתונים, לאמן מודל ללמידת מכונה, להריץ הסקה של קבוצות נתונים ולהעביר עסקאות בסיכון גבוה לתור לבדיקה ב-Cloud Spanner לצורך ביקורת ידנית.

בדרך כלל, התהליך הזה דורש ימים של כתיבת קוד הגדרה חוזר (מחברות Spark, הגדרות dbt, סקריפטים להדרכה, Airflow DAG) ומעבר מתמיד בין ממשקי מסוף ועורכים.

בשיעור Codelab הזה תעבדו בתכנות זוגי עם סוכן באמצעות Google Cloud Data Agent Kit (DAK) בתוך Antigravity IDE. באמצעות שפה טבעית שיחתית, הסוכן יעזור לכם ליצור מחברות Spark, לקמפל פרויקט dbt, לבנות לולאת היקש ולתזמר את תהליך העבודה באמצעות Managed Service for Apache Airflow.

הפעולות שתבצעו:

  • הטמעה של יומני מרכז סליקה מ-Cloud Storage באמצעות Managed Service for Apache Spark (Spark Serverless) בטבלה ב-BigQuery.
  • ביטול כפילויות ונרמול של טרנזקציות באמצעות dbt כדי ליצור שכבות נתונים נקיות (גולמיות, זמניות, מועשרות).
  • אימון מודל סיווג של יער אקראי מבוזר (RandomForestClassifier) ב-Spark Serverless.
  • הרצת הסקה (inference) של נתונים בכמות גדולה על עסקאות חדשות וכתיבת התראות על סיכון גבוה ישירות ל-Cloud Spanner.
  • תזמור, הגדרה חזותית ופריסה של צינור העיבוד כולו באמצעות Managed Service for Apache Airflow וניטור אינטראקטיבי של DAG בתוך סביבת הפיתוח המשולבת (IDE).

הדרישות

  • דפדפן אינטרנט כמו Chrome
  • פרויקט בענן של Google עם חיוב מופעל (מומלץ להשתמש בפרויקט חדש וייעודי למעבדות מעשיות).
  • היכרות בסיסית עם SQL,‏ Python ו-PySpark.
  • ‫Antigravity IDE עם מינוי ל-Google AI Pro (מומלץ)

העלות של המשאבים שנוצרו ב-codelab הזה צריכה להיות פחות מ-5$. בסיום שיעור ה-Lab, חשוב לפעול לפי ההוראות שבקטע ניקוי כדי למחוק את המשאבים שהוקצו.

2. הגדרת הסביבה

כדי להתחיל את המעבדה, מריצים סקריפט bootstrap. הסקריפט הזה מפעיל באופן אוטומטי את ממשקי ה-API הנדרשים של GCP, יוצר קטגוריה של Cloud Storage להטמעת נתונים, יוצר קבוצות נתונים מדומה של טרנזקציות וספריות, טוען ספריות הפניה ל-BigQuery ומתחיל הקצאת משאבים ברקע של Cloud Spanner ושל Managed Service for Apache Airflow (לשעבר Cloud Composer).

יצירת פרויקט חדש או בחירה בפרויקט קיים

בוחרים פרויקט קיים או יוצרים פרויקט חדש במסוף Google Cloud.

אימות החיוב

הקפידו לוודא שהחיוב מופעל בפרויקט בענן שלכם ב-Google Cloud. מידע נוסף על האופן שבו עושים את זה זמין במדריך הזה.

הפעלת סקריפט ההגדרה

תשתמשו ב-Google Cloud Shell (או במעטפת המקומית שהגדרתם באמצעות Google Cloud CLI) כדי להפעיל את הגדרת הסביבה.

  1. פותחים את מסוף Google Cloud.
  2. בסרגל הכלים שבפינה השמאלית העליונה, לוחצים על Activate Cloud Shell (הפעלת Cloud Shell).

פתיחת Cloud Shell

  1. בטרמינל של Cloud Shell, מגדירים את הפרויקט הפעיל:
gcloud config set project <<YOUR_PROJECT_ID>>
export PROJECT_ID=$(gcloud config get-value project)
  1. משכפלים את מאגר ה-codelab ועוברים לתיקיית התסריטים:
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. מריצים את סקריפט ההגדרה של bootstrap כדי לפרוס את כל המשאבים אל us-central1:
chmod +x setup.sh setup_spanner.sh setup_composer.sh
export REGION=us-central1
./setup.sh
  1. כשהסקריפט יסיים לפעול, יוצג סיכום שמציין שמערך הנתונים ב-BigQuery והקטגוריה ב-Cloud Storage מוכנים. ברקע, ימשיכו הקצאת המשאבים של Cloud Spanner (ייקח בערך 2 דקות) ושל Managed Airflow (ייקח בערך 20 דקות). אתם יכולים לעקוב אחרי ההתקדמות שלהם בכל שלב באמצעות הפקודה:
tail -f /tmp/spanner_setup.log
tail -f /tmp/composer_setup.log

פתיחת Antigravity IDE

  1. מורידים ומתקינים את Antigravity IDE מדף ההורדה של Google Antigravity.
  2. מפעילים את Antigravity IDE.
  3. יוצרים תיקייה חדשה וריקה במחשב המקומי (למשל, בשם agentic-data-labs) ופותחים אותה בסביבת הפיתוח המשולבת (IDE) על ידי בחירה באפשרות פתיחת תיקייה. התיקייה הזו תשמש כסביבת העבודה המקומית שלכם ל-codelab.

הגדרת תיקיית הפרויקט ב-Antigravity IDE

התקנת התוסף Data Agent Kit

התוסף Google Cloud Data Agent Kit מספק שילוב עמוק עם שירותי הנתונים של Google Cloud ישירות בתוך העורך, ומאפשר לכם ליצור אינטראקציה עם BigQuery,‏ Cloud SQL,‏ Cloud Storage ועוד בלי להחליף הקשרים.

  1. ב-Antigravity IDE, לוחצים על סמל התוספים בסרגל הפעילות בצד ימין של המסך (הסמל נראה כמו ארבעה ריבועים).
  2. בסרגל החיפוש בחלק העליון של חלונית התוספים, מקלידים Google Cloud Data Agent Kit.
  3. מאתרים את התוסף בשם Google Cloud Data Agent Kit שפורסם על ידי googlecloudtools
  4. לוחצים על הלחצן התקנה.
  5. יכול להיות שתופיע הנחיה עם השאלה "האם אתה בוטח בבעל התוסף googlecloudtools ובתוספים שלו?". לוחצים על Trust Publishers & Install (הבעת אמון בבעלי האתרים והתקנה) כדי להמשיך.

התקנת התוסף Data Agent Kit

אחרי ההתקנה, יופיע סמל חדש של Google Cloud Data Agent Kit בסרגל הפעילות בצד ימין של Antigravity IDE.

  1. דף ההצטרפות 'ברוכים הבאים ל-Google Cloud Data Agent Kit' אמור להיפתח באופן אוטומטי. אם לא נכנסתם לחשבון Cloud, פועלים לפי ההנחיות כדי לאשר את הגישה.
  2. בקטע Configuration Summary (סיכום ההגדרה), מאתרים את שדה הפרויקט. לוחצים על התפריט הנפתח ובוחרים את הפרויקט ב-Google Cloud. מגדירים את האזור בתור us-central1. לאחר מכן בוחרים באפשרות Configure MCP Servers (הגדרת שרתי MCP).

הגדרה ראשונית של התוסף Data Agent Kit

  1. לוחצים על Configure MCP Servers (הגדרת שרתי MCP). בחלונית MCP Configuration, מוודאים שהפעלתם את שרתי ה-MCP המרוחקים הבאים:
    • BigQuery
    • Spanner
    • מחברות

לוחצים על שנתחיל?.

הגדרת שרתי MCP

אפשרויות ההגדרה

בסיום תהליך ההגדרה, תופנו לדף Get started with Google Cloud Data Agent Kit (תחילת העבודה עם ערכת כלי הסוכן של Google Cloud Data).

  1. בקטע 'הגדרה וקביעת תצורה', לוחצים על תחילת העבודה.
  2. תיפתח החלונית הגדרה של Data Agent Kit. מעיינים בכרטיסיות:
    • פרויקט ואזור: מוודאים את מזהה הפרויקט שנבחר ומאשרים שסקריפט ההגדרה הפעיל את כל ממשקי ה-API הנדרשים (Compute Engine,‏ Cloud Storage,‏ BigQuery,‏ Spanner וכו').
    • ‫BigQuery: הגדרת מיקום ברירת המחדל לשאילתות BigQuery. משתמשים באזור us-central1.
    • הגדרת שרתי MCP: אפשר לראות את שרתי ה-MCP המופעלים (BigQuery,‏ Notebooks,‏ Spanner וכו') שמאפשרים לסוכני AI אינטראקציה מאובטחת עם הנתונים שלכם.
    • מיומנויות: אפשר לעיין במיומנויות מוגדרות מראש שמספקות לסוכנים יכולות ייעודיות למשימות מורכבות שקשורות לנתונים.

חלונית ההגדרות של Data Agent Kit

סיכום הקטע: הפעלתם את סקריפט האתחול כדי ליצור נכסים ב-GCS וב-BigQuery, בזמן שבניית Spanner ו-Airflow מתבצעת ברקע. לאחר מכן פתחתם את הפרויקט ב-Antigravity IDE והפעלתם את התוסף Google Cloud Data Agent Kit. עכשיו אפשר לכתוב את ה-Notebook הראשון.

3. הוספת יומנים גולמיים באמצעות Spark Serverless

בקטע הזה, תטמיעו יומני עסקאות בפורמט JSON גולמי באגם הנתונים. ‫Managed Service for Apache Spark ‏ (Spark Serverless) מתחבר ישירות לאחסון המקורי של BigQuery. תשתמשו במחבר ל-BigQuery הרגיל כדי לנהל נתונים טבלאיים ולאפשר שליחת שאילתות וניתוח נתונים ישירים.

היכרות עם סביבת זמן הריצה של Spark Serverless שהוגדרה מראש

לפני שמריצים קוד Spark, בודקים את תבנית Serverless Runtime שהוגדרה מראש על ידי סקריפט ההגדרה. התבנית הזו מגדירה את העורף של סביבת הביצוע של היעד ומאגדת את יחסי התלות הנדרשים של המחבר.

  1. בסרגל הפעילות של IDE, פותחים את החלונית Google Cloud Data Agent Kit.
  2. מרחיבים את התפריט הנפתח Apache Spark ואז מרחיבים את Serverless.
  3. לוחצים לחיצה ימנית על fraud-pipeline-runtime ובוחרים באפשרות פרופיל כדי לפתוח את תצוגת ההגדרות שלו בכלי העריכה.
  4. בכרטיסייה פרופיל, גוללים למטה ומרחיבים את מאפיינים כדי לבדוק את התלויות המותאמות אישית שמצורפות לסביבה:
    • ‫spark.jars: מכיל את gs://spark-lib/spanner/spark-3.5-spanner-1.4.0.jar, שמשתמש במחבר Spark Spanner כדי לאפשר לעבודות Spark לכתוב תוצאות של הסקה ישירות ל-Cloud Spanner בהמשך ה-Lab. (הערה: Dataproc Serverless כולל את מחבר Spark BigQuery של Google Cloud כברירת מחדל, כך שלא נדרש שום קובץ jar נוסף כדי לקרוא ולכתוב טבלאות BigQuery).

עיון במאפייני זמן הריצה של Spark Serverless

  1. בצד ימין, לוחצים על הכרטיסייה סשנים אינטראקטיביים. היא ריקה כרגע כי עדיין לא הפעלתם קוד. ברגע שתריצו את המחברת בשלב הבא, יוקצה באופן דינמי סשן חי של מחשוב ללא שרתים, והוא יופיע כאן.

הוספת נתונים באמצעות Data Agent Kit

במקום להגדיר ידנית סשן Spark או לכתוב סקריפטים לטעינת PySpark מאפס, תוכלו לבצע תכנות בזוגות עם סוכן באמצעות Data Agent Kit.

  1. לוחצים על סמל הפעלת הסוכן בסרגל הכלים בפינה השמאלית העליונה כדי לפתוח את החלונית צ'אט עם נציג.
  2. מדביקים את ההנחיה הבאה בצ'אט (חשוב להחליף את ${PROJECT_ID} במזהה הפרויקט בפועל ב-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. אם הסוכן מבקש הרשאה להריץ פקודות אימות ברקע (למשל "האם לאפשר הרצת הפקודה הזו?"), בודקים את הפקודה המוצעת ובוחרים באפשרות כן, לאפשר הפעם (או כן, ולאפשר תמיד).
  2. כשהסוכן מסיים ליצור את הקובץ, לוחצים על הלחצן הכחול אישור הכול (או על סמל סימן הוי) בחלק התחתון של חלונית הצ'אט כדי לשמור את notebooks/01_ingestion.ipynb בסביבת העבודה.

סוכן יוצר את הנוטבוק להעלאת נתונים

בדיקה והרצה של המחברת

  1. פותחים את notebooks/01_ingestion.ipynb שנוצר ב-IDE.
  2. בודקים את קוד PySpark של לוגיקת הכתיבה של מחבר ל-BigQuery.
  3. לוחצים על Run All בסרגל הכלים של ה-notebook ב-IDE.
  4. אם זו הפעם הראשונה שאתם מריצים מחברת Spark מרחוק, יכול להיות שתקבלו הנחיה להתקין תלות מקומית בסביבת הפיתוח המשולבת. אם מופיעה בקשה, לוחצים על Install dependencies for Remote Spark Kernels (התקנת תלות עבור ליבות Spark מרוחקות) ומאשרים את תיבות הדו-שיח של ההתקנה, ואז לוחצים שוב על Run All (הפעלת הכול).
  5. בתפריט הנפתח Select Kernel, בוחרים באפשרות Remote Spark Kernels -> fraud-pipeline-runtime on Serverless Spark. (טיפ: אם תבנית זמן הריצה שהוגדרה מראש לא מופיעה ברשימה, לוחצים על סמל הרענון בתפריט הנפתח של בחירת ליבת הקרנל בפינה השמאלית העליונה כדי לטעון מחדש את ליבות הקרנל המרוחקות שזמינות).
  6. בודקים את שורת הסטטוס בפינה הימנית התחתונה של הכלי לעריכה. יוצג Connecting to kernel: fraud-pipeline-runtime on Serverless Spark.... זו ההפעלה הראשונית של קצה העורף של ליבת זמן הריצה של Spark Serverless, ולכן ייקח כמה דקות להקצאה ולהפעלה.
  7. אחרי שהליבה מסיימת להתחבר, המחברת מתחילה באופן אוטומטי להריץ את כל התאים ברצף כדי לעבד את יומני העסקאות הגולמיים למערך הנתונים ב-BigQuery.

אימות

אחרי שההפעלה מסתיימת, בודקים את הקטלוג של Data Agent Kit כדי לוודא שהטבלה נוצרה:

אימות הטבלה הגולמית בכלי לבדיקת קטלוג

  1. בסרגל הפעילות של IDE, פותחים את החלונית Google Cloud Data Agent Kit.
  2. מרחיבים את הקטע CATALOG.
  3. מרחיבים את מזהה הפרויקט.
  4. מרחיבים את BigQuery.
  5. מרחיבים את קבוצת הנתונים transactions_dataset_evals.
  6. לוחצים על הטבלה raw_transactions כדי לפתוח את תצוגת הפרטים שלה בכלי העריכה הראשי.
  7. בחלונית הניווט הימנית, בוחנים את הכרטיסיות Data,‏ Schema ו-Details כדי לבדוק את הרשומות והמטא-נתונים שהועברו.

סיכום הקטע: השתמשתם בשפה טבעית בצ'אט עם הסוכן כדי ליצור עומס עבודה מלא של Spark Serverless. לאחר מכן הפעלתם אותה כדי לעבד יומני JSON לא מובְנים לטבלה (גולמית) ב-BigQuery.

4. ביטול כפילויות ונרמול באמצעות dbt

לפני אימון מודל ה-ML, תצטרכו לאכוף את איכות הנתונים על ידי הסרת יומני סטרימינג כפולים, בידוד רשומות לא תקינות (כמו מזהי עסקאות ריקים) וצירוף נתונים ממדידים (משלמים ומוטבים). התהליך הזה דורש טרנספורמציות של SQL שהן אידמפוטנטיות ואמינות, ולכן dbt (data build tool) מתאים מאוד למטרה הזו.

יצירת מבנה של צינור עיבוד נתונים של dbt

משתמשים בסוכן כדי ליצור פרויקט dbt על מערך הנתונים ב-BigQuery:

  1. חוזרים לחלונית Agent Chat.
  2. כדי ליצור את פרויקט dbt, צריך להזין את ההוראה הבאה:
Scaffold a us-central1 dbt project in dbt_project/ that maps raw_transactions
through to an enriched_transactions model in dataset transactions_dataset_evals.

Deduplicate by transaction_id in staging. Quarantine null IDs to an invalid_transactions model.
Join the valid staging records with dim_payers and dim_payees for the enriched_transactions model,
preserving the historical `is_fraud` label column, and finally add a transaction uniqueness test.

Create an implementation plan first.
  1. הסוכן יציג את ארטיפקט תוכנית ההטמעה בחלונית העריכה הראשית. בודקים את מבנה הקבצים המוצע ואת הלוגיקה של ה-SQL.
  2. לוחצים על המשך (ואז על אישור הכל) כדי לאפשר לסוכן ליצור את הקבצים בסביבת העבודה.

תוכנית הטמעה עם כפתור 'המשך'

  1. אחרי שהיצירה מסתיימת, הסוכן מציג הסבר מפורט עם סיכום של הרכיבים החדשים. אם מוצגת בקשה, מאשרים את כל השינויים.

אישור כל הקבצים שנוצרו בחלונית הצ'אט

פיתוח ובדיקה

למרות שהסוכן הפעיל אוטומטית את dbt compile כדי לוודא שקוד ה-SQL שנוצר תקף מבחינת התחביר, עכשיו תצרו את התצוגות והטבלאות האלה ב-BigQuery ותפעילו את בדיקות איכות הנתונים כדי לבצע אימות מקומי. (הערה: בהמשך ה-Lab, תהפכו את שלב ה-dbt הזה לאוטומטי כחלק מ-DAG של Airflow מקצה לקצה).

  1. בסרגל הפעילות בצד ימין, לוחצים על סמל הסייר (או לוחצים על Cmd/Ctrl+Shift+E).
  2. מרחיבים את dbt_project -> models כדי לבדוק את מודלי ה-SQL שנוצרו. לוחצים על enriched_transactions.sql כדי לפתוח את הלוגיקה של התכונה 'שינוי' ו'הונאה' בעורך ולבדוק אותה.
  3. בסייר הקבצים, לוחצים לחיצה ימנית על התיקייה dbt_project ובוחרים באפשרות פתיחה במסוף המשולב. פעולה זו תפתח באופן אוטומטי חלונית טרמינל שמוגדרת ישירות לdbt_project ספריית העבודה הנדרשת.
  4. אם עדיין לא התקנתם את dbt, יוצרים סביבה וירטואלית מחוץ ל-dbt_project/ (בספריית הבסיס של הבית או של סביבת העבודה) ומתקינים את המתאם של BigQuery:
python3 -m venv ~/.venv/dbt
source ~/.venv/dbt/bin/activate
pip install dbt-bigquery
  1. מריצים את מודלי dbt ואת בדיקות איכות הנתונים שמשויכות אליהם:
dbt build
  1. צופים בפלט של הטרמינל. ‏dbt יקמפל את ה-SQL, יממש את הטבלאות של הנתונים הזמניים והנתונים המועשרים ב-BigQuery ויבצע את בדיקות הנתונים.

יצירה ובדיקה של פרויקט dbt בטרמינל משולב

  1. אחרי שהבנייה מסתיימת, סוגרים את חלונית המסוף כדי לפנות מקום במסך לשלבים הבאים.

סיכום הקטע: יצרתם פרויקט dbt באמצעות הסוכן, הרצתם בדיקות של איכות הנתונים והמרתם את הרשומות הגולמיות לטבלאות ביניים ולטבלאות מועשרות ב-BigQuery.

5. אימון מודל מבוזר לזיהוי הונאות באמצעות Random Forest

אחרי שהעסקאות המועשרות יתממשו ב-BigQuery, תוכלו ליצור מודל של למידת מכונה כדי לסווג אירועים שהם הונאה. ‫Random Forest היא שיטת למידה משולבת שמתאימה במיוחד לנתוני סיווג טבלאיים. הפעלת RandomForestClassifier ב-Spark Serverless מפזרת את אימון המודל על צמתי עובדים בלי לדרוש מכם לנהל את התשתית.

בשלב הזה, תשתמשו בסוכן כדי ליצור את צינור העיבוד לאימון של Spark ML.

יצירת נוטבוק לאימון מודלים של למידת מכונה

  1. פותחים את החלונית Agent Chat.
  2. מזינים את ההנחיה הבאה כדי לעצב את רצף אימון המודל (חשוב להחליף את ${PROJECT_ID} במזהה הפרויקט הפעיל):
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. בודקים את התוכנית או את הקוד שנוצר על ידי הסוכן ולוחצים על המשך או על אישור הכל כדי לשמור את notebooks/02_training.ipynb בסביבת העבודה.

הסוכן יוצר את הנוטבוק לאימון

בדיקה והרצה של המחברת

  1. פותחים את notebooks/02_training.ipynb בעורך.
  2. בודקים את השלבים בצינור עיבוד הנתונים של PySpark ML לקידוד תכונות, להרכבת וקטורים וללוגיקה של סיווג Random Forest.
  3. לוחצים על Run All בסרגל הכלים של ה-notebook ב-IDE.
  4. כשנפתח התפריט הנפתח Select Kernel, בוחרים באפשרות fraud-pipeline-runtime on Serverless Spark.

בחירת ליבת Serverless Spark עבור מחברת האימון

אימות

אחרי שההפעלה מסתיימת, מוודאים שהמודל אומן ויוצא בצורה תקינה:

  1. כדי לאמת את הציון של השטח מתחת לעקומת ROC (AUC) שדווח, מעיינים בפלט של תא ההערכה בחלק התחתון של המחברת.
  2. כדי לוודא שפריטי המודל נשמרו בהצלחה ב-GCS, מרחיבים את חלונית הסייר STORAGE בסרגל הצד של Data Agent Kit.
  3. מחפשים את ה-bucket שמסתיים ב--models (שמשויך למזהה הפרויקט הפעיל), מרחיבים אותו ומעמיקים כדי לוודא שספריית fraud_model ושלבי הצינור שלה קיימים.

אימות שהמודל נשמר ב-GCS

סיכום הקטע: השתמשתם בסוכן כדי ליצור פייפליין אימון של PySpark ML, אימנתם מודל של Random Forest בטבלה ב-BigQuery המועשרת וייצאתם את המודל ל-Cloud Storage.

6. הסקת מסקנות באצווה וכתיבה ב-Cloud Spanner

בעזרת מודל חיזוי מאומן שמאוחסן ב-Cloud Storage, תריצו הסקה (inference) של קבוצות נתונים על טרנזקציות חדשות שזורמות דרך BigQuery. צריך להפנות עסקאות בסיכון גבוה למערכת תפעולית כדי שצוות התאימות יוכל לבדוק אותן. ‫Cloud Spanner מספק מסד נתונים טרנזקציוני שניתן להתאמה לתור הבדיקה הזה.

יצירת מחברת של מסקנות אצווה

שימוש בסוכן כדי ליצור מחברת הסקה שמקשרת בין BigQuery,‏ Cloud Storage ו-Cloud Spanner:

  1. פותחים את החלונית Agent Chat.
  2. מזינים את ההנחיה הבאה:
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. כדי לשמור את notebooks/03_inference.ipynb בסביבת העבודה, מאשרים את ה-notebook שנוצר.

הסוכן שיוצר את מחברת ההסקה

בדיקה והרצה של המחברת

  1. פותחים את notebooks/03_inference.ipynb שנוצר בכלי העריכה.
  2. בודקים את רצף ההסקה של PySpark:
    • יחסי תלות: תבנית Serverless Runtime מספקת את יחסי התלות הנדרשים של cloud-spanner JAR להרצת Spark.
    • עיצוב נתונים: הסקריפט משמיט עמודות וקטוריות מורכבות של Spark ML (כמו תכונות גולמיות והסתברויות) לפני הכתיבה, כדי להתאים לסכימת הטבלה של Spanner.
    • ‫Spanner Connector: הוא כותב את השורות המסומנות באמצעות .format("cloud-spanner") כדי להוסיף אותן ישירות לתור הבדיקה.
  3. לוחצים על Run All בסרגל הכלים של ה-notebook ב-IDE.
  4. כשמופיעה בקשה לבחור ליבה, בוחרים באפשרות fraud-pipeline-runtime on Serverless Spark.

אימות

אחרי ש-notebook ההסקות יסיים את העיבוד, תוכלו לשלוח שאילתות למסד הנתונים התפעולי של Spanner ישירות בתוך סביבת הפיתוח המשולבת:

  1. בסרגל הפעילות של IDE, פותחים את החלונית Google Cloud Data Agent Kit.
  2. מרחיבים את הקטע CATALOG.
  3. מרחיבים את מזהה הפרויקט ואז מרחיבים את Spanner.
  4. עוברים אל cymbal-fraud -> fraud-db -> Tables -> SparkEvalFraudReviewQueue.
  5. לוחצים לחיצה ימנית על הטבלה ובוחרים באפשרות Query Table (שאילתת הטבלה), ואז מריצים את השאילתה:
SELECT *
FROM `SparkEvalFraudReviewQueue`
LIMIT 100;
  1. בחלונית Query Results שבהמשך, אמורות להופיע שורות חדשות שנוספו ומייצגות עסקאות עם רמת סיכון גבוהה שסומנו לבדיקה ידנית.

אימות שורות ב-Cloud Spanner

סיכום הקטע: השתמשתם בסוכן כדי ליצור מחברת של הסקת מסקנות באצווה, קיבלתם ציון לרשומות לא מסומנות ב-BigQuery באמצעות המודל שאומן, וכתבתם טרנזקציות בסיכון גבוה ישירות ל-Cloud Spanner.

7. יצירת תשתית ותזמור באמצעות Managed Airflow

הצינור שלך מורכב כרגע משלבים נפרדים: מחברת קליטה, מפרויקט טרנספורמציה של dbt וממחברת הסקה של אצווה. כדי להפוך את התהליך הזה למוכן לייצור, צריך לחבר את כל הפעולות לגרף תלות מתוזמן.

‫Managed Service for Apache Airflow (שנקרא בעבר Cloud Composer) מספק מנוע תזמור מנוהל לתהליך העבודה הזה. חבילת Data Agent כוללת את התכונה Orchestration Pipelines, שמתרגמת הגדרות של צינורות YAML דקלרטיביים ישירות ל-DAG של Airflow.

הגדרת הפייפליין

משתמשים בסוכן כדי ליצור את ההגדרה של צינור עיבוד הנתונים לניהול התהליך:

  1. בצ'אט עם סוכן, מזינים את ההנחיה הבאה (לא לשכוח להחליף את ${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.

בדיקת ההגדרות של ה-DAG

הכלי Data Agent Kit Orchestrator משתמש בהגדרות YAML דקלרטיביות כדי להגדיר ולפרוס צינורות עיבוד נתונים ב-Apache Airflow, וכך מאפשר לשלוט בגרסאות של ההגדרות ולפרוס אותן באמצעות CI/CD.

בחלונית Explorer של סביבת הפיתוח המשולבת, בודקים את שני קובצי הצינורות שהסוכן יצר בשורש של סביבת העבודה:

  1. ‫deployment.yaml: פתיחת הקובץ. הוא משמש כרישום של הסביבה שלכם. הוא ממפה את צינור עיבוד הנתונים הלוגי dev לסביבה cymbal-airflow, מגדיר את אזור ההפעלה (us-central1) ומגדיר את קטגוריית artifact_storage ה-DAG והתלות שבהם מתבצעת ההכנה.
  2. ‫fraud_analysis_pipeline.yaml: פתיחת הקובץ. ההגדרה הזו מגדירה את גרף הביצוע. הוא מציין את לוח הזמנים של הטריגר (interval: '0 0 * * *') ומסדר את שלושת השלבים ברצף בקטע actions:
    • פעולת הטמעה notebook של 01_ingestion.ipynb שפועלת ב-Dataproc Serverless.
    • פעולת טרנספורמציה pipeline שמטרגטת את ספריית dbt_project, עם תלות dependsOn שמצביעה על שלב ההטמעה.
    • פעולת הסקה notebook עבור 03_inference.ipynb עם תלות dependsOn שמפנה לשלב dbt, כולל מאפיין Spanner JAR.
  3. הסוכן גם יסכם את הארטיפקטים שנוצרו בכרטיסייה Walkthrough בחלונית העריכה, ויפרט את ההגדרות והאימותים שבוצעו.

הגדרת DAG אינטראקטיבית

ערכת הכלים Data Agent Kit מציגה את הגדרות הפייפליין כתרשים חזותי אינטראקטיבי, שמאפשר לבדוק ולערוך את מאפייני ה-DAG של Airflow.

  1. בסרגל הפעילות של IDE, פותחים את החלונית Google Cloud Data Agent Kit.
  2. בקטע DATA ENGINEERING, מרחיבים את Orchestration Pipelines.
  3. לוחצים על fraud_analysis_pipeline.yaml כדי לפתוח את ה-DAG הוויזואלי ב-Canvas בעורך הראשי.

קנבס ויזואלי של תזמור DAG

  1. לוחצים על הצומת Schedule trigger בחלק העליון. תפריט נפתח של הגדרות ייפתח בצד שמאל, ויוצג בו מחרוזת Cron שנותחה (0 0 * * *). תוכלו לשנות פרמטרים כמו backfill ו-catchup.
  2. לוחצים על צומת משימה של נוטבוק (למשל, שלב ההטמעה או ההסקה). התפריט הנפתח מתעדכן ומציג את מיפויי ההרצה הספציפיים של Dataproc Serverless ואת מאפייני המחבר.
  3. שימו לב להיפר-קישור של שם קובץ ה-notebook (לדוגמה, 01_ingestion.ipynb) בתוך בלוק הצומת. לחיצה על הסמל תפתח את המחברת ישירות בעורך.
  4. בסרגל הצד הימני, מתחת ל-Orchestration Pipelines, לוחצים על Deployment configuration. בתצוגה הזו מוצגים קלאסטר הסביבה של dev היעד וארטיפקטים של קטגוריית GCS של הפלט.

סיכום הקטע: יצרתם הגדרה של צינור עיבוד נתונים באמצעות הסוכן, והגדרתם תלות בין משימות של קליטה, dbt והסקת מסקנות בלוח ציור אינטראקטיבי.

8. פריסה, הפעלה ומעקב

אחרי שמגדירים את ה-DAG באופן מקומי, מתחברים לסביבת Managed Airflow שהוקצתה במהלך ההגדרה ומפעילים את הצינור.

הגדרת Managed Service for Apache Airflow

לפני הפריסה, צריך להגדיר את החיבור של Scheduler בהגדרות של Data Agent Kit כדי שהתוסף יכוון לסביבת Managed Airflow שלכם:

  1. בסרגל הפעילות של IDE, פותחים את החלונית Google Cloud Data Agent Kit.
  2. בקטע SETTINGS, לוחצים על הגדרות.
  3. בתפריט הימני, לוחצים על Scheduler (מתזמן).
  4. קובעים את ההגדרות:
    • מזהה הפרויקט: בוחרים את מזהה הפרויקט הפעיל.
    • אזור: בוחרים באפשרות us-central1.
    • סביבה: בוחרים באפשרות cymbal-airflow.
  5. לוחצים על שמירה.

הגדרות של Managed Service for Apache Airflow

פריסת ה-DAG

עכשיו אפשר לפרוס את הפייפליין שהגדרתם ישירות בסביבת Managed Airflow שלכם מלוח הציור החזותי:

  1. בסרגל הצד Google Cloud Data Agent Kit, מרחיבים את DATA ENGINEERING > Orchestration Pipelines ולוחצים על fraud_analysis_pipeline.yaml כדי לפתוח את בד הציור של ה-DAG החזותי.
  2. בפינה השמאלית העליונה של סרגל הכלים של אזור העריכה, לוחצים על הלחצן הכחול הפעלת צינור.
  3. בתפריט הנפתח לבחירת סביבה, בוחרים באפשרות dev.
  4. אפשר לעקוב אחרי ההתקדמות בהודעה שמופיעה באזור הסטטוס התחתון (Running pipeline: Building pipeline locally...). התוסף יקמפל אוטומטית את ה-DAG, יארוז את המחברת ואת נכסי ה-dbt ויעלה אותם לדלי GCS של סביבת Managed Airflow (התהליך הזה נמשך כ-3-4 דקות).

פריסת צינור הנתונים מלוח הציור החזותי

מעקב אחרי ההרצה

אחרי שהקומפילציה המקומית מסתיימת וההתראה הקופצת מאשרת את הפעולה Triggered a new run for pipeline... successfully, עוקבים אחרי ההרצה הפעילה:

  1. בסרגל הצד Google Cloud Data Agent Kit, מרחיבים את האפשרויות DATA ENGINEERING > Orchestration Pipelines.
  2. לוחצים על ניהול צינורות.
  3. בטבלה 'ניהול צינורות', לוחצים על fraud_analysis_pipeline כדי לפתוח את היסטוריית ההרצה של הצינור.

סקירה כללית על ניהול צינורות

  1. בתצוגה Execution History (היסטוריית ההרצה), בוחרים את ההרצה הפעילה מהיומן.
  2. ככל שהביצוע מתקדם בכל משימה בצינור (הטמעה, טרנספורמציה של dbt והסקת מסקנות), מתעדכנים אינדיקטורים של הסטטוס ומתמלאים משכי הזמן של המשימות. לוחצים על משימה כלשהי כדי לבדוק את הפלט של ההפעלה שלה בשידור חי ואת היומנים של Airflow DAG.

היסטוריית ההרצה של צינורות בזמן אמת ופרטי המשימות

סיכום הקטע: הגדרתם את החיבור של Airflow Scheduler, פרסתם את צינור הניתוח מקצה לקצה ב-Managed Airflow ועקבתם אחרי ביצוע בזמן אמת, כדי לוודא שהמערכת פועלת בצורה תקינה מיומנים גולמיים ועד לתחזיות הסופיות ב-Cloud Spanner.

9. הסרת המשאבים

כדי להימנע מחיובים שוטפים בפרויקט בענן של Google על המשאבים שבהם השתמשתם ב-Codelab הזה, צריך להשבית את הסביבה באמצעות הסקריפט האוטומטי.

  1. בחלונית Terminal (או ב-Cloud Shell), עוברים לספריית הסקריפטים ומריצים את הפקודה:
cd ~/devrel-demos/codelabs/agentic-data-labs/data-science/scripts
chmod +x teardown.sh
./teardown.sh
  1. הסקריפט יציג רשימה של כל המשאבים שהוא מתכנן למחוק ויבקש אישור:
    • סביבת Managed Airflow‏ (cymbal-airflow)
    • Cloud Spanner Instance (cymbal-fraud)
    • מערך נתונים ב-BigQuery (transactions_dataset_evals)
    • קטגוריות של Cloud Storage (gs://${PROJECT_ID}-fin-clearing-raw ו-gs://${PROJECT_ID}-models)
    • חשבון שירות של Worker (composer-worker-sa)
  2. כדי לאשר, מקלידים y. סקריפט ההסרה יסיר את כל שירותי ה-GCP שהוקצו וינקה את הקבצים המקומיים.

10. מעולה!

יצרתם צינור עיבוד נתונים מקיף לזיהוי הונאות שכולל את Cloud Storage,‏ BigQuery,‏ Managed Service for Apache Spark (Spark Serverless),‏ dbt,‏ Cloud Spanner ו-Managed Service for Apache Airflow, באמצעות תכנות בזוגות עם Google Cloud Data Agent Kit בתוך Antigravity IDE.

מה השגתם

  1. ‫📥 הטמעה של יומני טרנזקציות גולמיים בטבלה ב-BigQuery באמצעות Managed Service for Apache Spark ו-Data Agent Kit.
  2. ‫🧹 הסרת כפילויות מנתונים וביצוע נורמליזציה שלהם באמצעות יצירת פרויקט dbt עם בדיקות של איכות הנתונים.
  3. 🤖 מודל יער אקראי מבוזר שאומן באמצעות RandomForestClassifier ויוצא ל-Cloud Storage.
  4. ⚡ הפעלת הסקה (inference) של נתונים בכמות גדולה על עסקאות נכנסות והעברת רשומות בסיכון גבוה ל-Cloud Spanner לצורך ביקורת.
  5. ‫🔄 תזמור, פריסה וניטור של תהליך העבודה כ-DAG מתוזמן של Airflow באמצעות Managed Service for Apache Airflow וכלי ה-DAG לניהול חזותי של סביבת הפיתוח המשולבת (IDE).

מושגים מרכזיים

קונספט

מה למדתם

Data Agent Kit

תכנות בזוגות בתוך סביבת הפיתוח המשולבת (IDE) באמצעות שפה טבעית כדי ליצור מחברות PySpark, להגדיר מודלים של dbt ולהגדיר DAG של Airflow

BigQuery

אחסון טבלאי ניתן להרחבה ל-SQL אנליטי, לטרנספורמציות של dbt ולאימון של למידת מכונה

‫Spark Serverless

ביצוע ללא שרת (serverless) לטעינת נתונים מבוזרת של PySpark והדרכה של למידת מכונה (ML) מסוג Random Forest

Cloud Spanner Connector

כתיבת תחזיות של הסקת מסקנות באצווה של Spark ישירות לתורים של ביקורות במסד נתונים תפעולי

הצהרות YAML DAG

הגדרות הצהרה של צינורות עיבוד נתונים מוצגות כתרשימים חזותיים אינטראקטיביים של Airflow בסביבת הפיתוח המשולבת

ניהול חזותי של DAG

בדיקת התלויות של צינור עיבוד הנתונים, פריסה אל Managed Airflow ומעקב אחר היסטוריית הביצוע של משימות בזמן אמת בתוך סביבת הפיתוח המשולבת

השלבים הבאים