مسار رصد عمليات الاحتيال باستخدام "حزمة وكيل البيانات" وبيئة التطوير المتكاملة Antigravity

1. مقدمة

تخيّل أنّك عالم بيانات في Cymbal Financial، وهي شركة لمعالجة الدفعات بكميات كبيرة. حدثت موجة من التأخيرات في التسوية، ويشتبه فريق الامتثال في حدوث عمليات احتيال منسَّقة. عليك إنشاء مسار لنقل سجلّات معاملات غرفة المقاصة الأولية، وتنظيف البيانات، وتدريب نموذج تعلُّم آلي، وتشغيل الاستدلال المجمّع، ونقل المعاملات عالية الخطورة إلى قائمة مراجعة في Cloud Spanner لإجراء التدقيق اليدوي.

يتطلّب ذلك عادةً كتابة رمز إعداد متكرّر (دفاتر ملاحظات Spark، وإعدادات dbt، وبرامج التدريب النصية، ومخططات DAG في Airflow) والتبديل المستمر بين واجهات وحدة التحكّم والمحرّرات.

في هذا الدرس التطبيقي حول الترميز، ستعمل مع وكيل باستخدام حزمة أدوات وكيل البيانات (DAK) من Google Cloud داخل بيئة تطوير Antigravity المتكاملة. باستخدام اللغة الطبيعية الحوارية، سيساعدك الوكيل في إنشاء دفاتر ملاحظات Spark، وتجميع مشروع dbt، وإنشاء حلقة استدلال، وتنظيم سير العمل باستخدام الخدمة المُدارة لـ Apache Airflow.

الإجراءات التي ستنفذّها

  • استيعاب سجلات غرفة المقاصة من Cloud Storage باستخدام الخدمة المُدارة لـ Apache Spark (الحوسبة بدون خادم من Spark) في جدول BigQuery
  • إزالة البيانات المكرّرة وتسوية المعاملات باستخدام dbt لإنشاء طبقات بيانات نظيفة (البيانات الأولية، وبيانات الإعداد، والبيانات المحسّنة)
  • تدريب نموذج تصنيف Random Forest موزّع (RandomForestClassifier) على Spark Serverless
  • تنفيذ استنتاج مجمّع على المعاملات الجديدة وكتابة تنبيهات بشأن المخاطر العالية مباشرةً إلى Cloud Spanner
  • تنظيم مسار التعلّم بأكمله وإعداده بشكل مرئي وتفعيله باستخدام Managed Service for Apache Airflow ورصد الرسوم البيانية الموجّهة غير الدورية التفاعلية داخل IDE

المتطلبات

  • متصفّح ويب، مثل Chrome
  • مشروع على Google Cloud تم تفعيل الفوترة فيه (ننصح باستخدام مشروع جديد ومخصّص للمختبرات العملية).
  • الإلمام بأساسيات SQL وPython وPySpark
  • بيئة التطوير المتكاملة Antigravity مع اشتراك في Google AI Pro (يُنصح به)

يجب أن تكون تكلفة الموارد التي تم إنشاؤها في هذا الدرس التطبيقي حول الترميز أقل من 5 دولارات أمريكية. احرص على اتّباع تعليمات التنظيف في نهاية الدرس العملي لحذف الموارد التي تم توفيرها.

2. إعداد البيئة

لبدء المعمل، عليك تشغيل نص برمجي للتمهيد. يعمل هذا النص البرمجي تلقائيًا على تفعيل واجهات برمجة التطبيقات المطلوبة في Google Cloud Platform، وإنشاء حزمة Cloud Storage لنقل البيانات، وإنشاء مجموعات بيانات وهمية للمعاملات والأدلة، وتحميل الأدلة المرجعية إلى BigQuery، وبدء عملية توفير Cloud Spanner و"الخدمة المُدارة لـ Apache Airflow" (المعروفة سابقًا باسم Cloud Composer) في الخلفية.

اختيار مشروع أو إنشاؤه

اختَر مشروعًا حاليًا أو أنشِئ مشروعًا جديدًا في Google Cloud Console.

تأكيد الفوترة

تأكَّد من تفعيل الفوترة لمشروعك على Google Cloud. يمكنك الاطّلاع على مزيد من المعلومات حول كيفية تنفيذ هذا الإجراء من خلال اتّباع هذا الدليل.

تشغيل نص التهيئة البرمجي

ستستخدم Google Cloud Shell (أو shell المحلي الذي تم إعداده باستخدام Google Cloud CLI) لبدء عملية إعداد البيئة.

  1. افتح Google Cloud Console.
  2. انقر على تفعيل 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. نفِّذ نص التهيئة البرمجي الأوّلي لنشر جميع الموارد إلى us-central1:
chmod +x setup.sh setup_spanner.sh setup_composer.sh
export REGION=us-central1
./setup.sh
  1. عندما ينتهي النص البرمجي، سيظهر لك ناتج ملخّص يشير إلى أنّ مجموعة بياناتك في BigQuery وحزمة Cloud Storage جاهزتان. في الخلفية، سيستمر توفير Cloud Spanner (يستغرق ذلك حوالي دقيقتَين) و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) من خلال اختيار فتح مجلد. سيكون هذا بمثابة مساحة عملك المحلية في هذا الدرس التطبيقي حول الترميز.

ضبط مجلد مشروع Antigravity IDE

تثبيت إضافة "حزمة وكيل البيانات"

يوفّر إضافة "مجموعة أدوات Google Cloud Data Agent" تكاملاً عميقًا مع خدمات بيانات Google Cloud مباشرةً في المحرّر، ما يتيح لك التفاعل مع BigQuery وCloud SQL وCloud Storage وغير ذلك بدون تبديل السياقات.

  1. في بيئة التطوير المتكاملة Antigravity، انقر على رمز الإضافات في "شريط الأنشطة" في أقصى يمين الشاشة (يبدو الرمز على شكل أربعة مربعات).
  2. في شريط البحث أعلى لوحة "الإضافات"، اكتب Google Cloud Data Agent Kit.
  3. ابحث عن الإضافة المسماة Google Cloud Data Agent Kit التي نشرها googlecloudtools.
  4. انقر على الزر تثبيت.
  5. قد تظهر رسالة تسألك: "هل تثق في الناشر googlecloudtools وإضافاته؟". انقر على الوثوق بالناشرين والتثبيت للمتابعة.

تثبيت إضافة &quot;حزمة أدوات وكيل البيانات&quot;

بعد التثبيت، سيظهر رمز جديد Google Cloud Data Agent Kit في شريط الأنشطة في أقصى يمين بيئة التطوير المتكاملة Antigravity.

  1. من المفترض أن تفتح تلقائيًا صفحة إعداد بعنوان "مرحبًا بك في Google Cloud Data Agent Kit". إذا لم تكن مسجِّلاً الدخول إلى حسابك على Cloud، اتّبِع أي طلبات للسماح بالوصول.
  2. في قسم ملخّص الإعداد، ابحث عن حقل المشروع. انقر على القائمة المنسدلة واختَر مشروعك على Google Cloud. اضبط منطقتك على us-central1. بعد ذلك، انقر على ضبط خوادم MCP.

الإعداد الأوّلي لإضافة Data Agent Kit

  1. انقر على ضبط خوادم MCP. ضمن لوحة إعدادات MCP، تأكَّد من تفعيل خوادم MCP البعيدة التالية:
    • BigQuery
    • Spanner
    • دفاتر الملاحظات

بعد ذلك، انقر على البدء.

إعداد خوادم MCP

استكشاف خيارات الإعداد

بعد اكتمال عملية الإعداد، ستنتقل إلى صفحة "البدء باستخدام Google Cloud Data Agent Kit".

  1. ضمن "الإعداد والتكوين"، انقر على البدء.
  2. سيؤدي ذلك إلى فتح لوحة إعدادات حزمة Data Agent. استكشاف علامات التبويب:
    • المشروع والمنطقة: تحقَّق من رقم تعريف المشروع الذي اخترته وتأكَّد من أنّ نص التهيئة البرمجي فعّل جميع واجهات برمجة التطبيقات المطلوبة (Compute Engine وCloud Storage وBigQuery وSpanner وما إلى ذلك).
    • ‫BigQuery: اضبط الموقع الجغرافي التلقائي لطلبات البحث في BigQuery. استخدِم المنطقة us-central1.
    • إعداد خوادم MCP: يمكنك الاطّلاع على خوادم MCP المفعّلة (مثل BigQuery وNotebooks وSpanner وما إلى ذلك) التي تتيح لوكلاء الذكاء الاصطناعي التفاعل بأمان مع بياناتك.
    • المهارات: استكشِف المهارات المُعدّة مسبقًا التي تزوّد الوكلاء بقدرات متخصّصة لإنجاز مهام البيانات المعقّدة.

لوحة إعدادات &quot;حزمة أدوات وكيل البيانات&quot;

ملخّص القسم: شغّلت نص الإعداد الأوّلي لإنشاء مواد عرض GCS وBigQuery أثناء إنشاء Spanner وAirflow في الخلفية. بعد ذلك، فتحت المشروع في Antigravity IDE وفعّلت إضافة "حزمة أدوات Google Cloud Data Agent". أنت الآن جاهز لكتابة دفتر ملاحظاتك الأول.

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، ثم وسِّع بلا خادم.
  3. انقر بزر الماوس الأيمن على fraud-pipeline-runtime واختَر الملف الشخصي لفتح طريقة عرض الإعدادات في المحرّر.
  4. في علامة التبويب الملف الشخصي، انتقِل للأسفل ووسِّع الخصائص لفحص التبعيات المخصّصة المرفقة بالبيئة:
    • spark.jars: يحتوي على gs://spark-lib/spanner/spark-3.5-spanner-1.4.0.jar الذي يستخدم أداة ربط Spark Spanner للسماح لمهام Spark بكتابة نتائج الاستنتاج مباشرةً إلى Cloud Spanner في وقت لاحق من المختبر. (ملاحظة: تتضمّن خدمة Dataproc Serverless موصّل Spark BigQuery من Google Cloud تلقائيًا، بدون الحاجة إلى أي إعدادات إضافية لملف jar من أجل قراءة جداول BigQuery وكتابتها).

استكشاف خصائص وقت التشغيل غير الخادم في Spark

  1. لاحِظوا علامة التبويب الجلسات التفاعلية على يمين الصفحة. وهي فارغة حاليًا لأنّك لم تنفّذ أي رمز برمجي بعد. بعد تشغيل دفتر الملاحظات في الخطوة التالية، سيتم توفير جلسة حوسبة مباشرة بدون خادم بشكل ديناميكي وستظهر هنا.

استيعاب البيانات باستخدام "حزمة Data Agent"

بدلاً من إعداد جلسة 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 الذي تم إنشاؤه حديثًا في بيئة التطوير المتكاملة.
  2. راجِع رمز PySpark لمنطق الكتابة في أداة ربط BigQuery.
  3. انقر على تنفيذ الكل في شريط أدوات ورقة الملاحظات في بيئة التطوير المتكاملة.
  4. إذا كانت هذه هي المرة الأولى التي تشغّل فيها دفتر ملاحظات Spark عن بُعد، قد تطلب منك بيئة التطوير المتكاملة تثبيت التبعيات المحلية. إذا طُلب منك ذلك، انقر على تثبيت التبعيات لـ Remote Spark Kernels وأكِّد مربّعات حوار التثبيت، ثم انقر على تشغيل الكل مرة أخرى.
  5. في القائمة المنسدلة اختيار النواة، اختَر نواة Spark عن بُعد -> fraud-pipeline-runtime على Serverless Spark. (ملاحظة: إذا لم يظهر نموذج وقت التشغيل الذي تم ضبطه مسبقًا في القائمة، انقر على رمز إعادة التحميل في أعلى يسار القائمة المنسدلة لاختيار النواة لإعادة تحميل النواة البعيدة المتاحة).
  6. انظر إلى شريط الحالة في أسفل يمين المحرّر. سيظهر لك Connecting to kernel: fraud-pipeline-runtime on Serverless Spark.... بما أنّ هذه هي عملية الإطلاق الأوّلية لخادم Spark Serverless، سيستغرق توفير الخادم وتشغيله بضع دقائق.
  7. بعد أن تنتهي النواة من عملية الربط، سيبدأ دفتر الملاحظات تلقائيًا في تنفيذ جميع الخلايا بالتسلسل لمعالجة سجلات المعاملات الأولية في مجموعة بيانات BigQuery.

التحقق

بعد اكتمال التنفيذ، راجِع فهرس "حزمة أدوات وكيل البيانات" للتأكّد من إنشاء الجدول:

التحقّق من جدول Raw في &quot;مستكشف الفهرس&quot;

  1. في شريط أنشطة بيئة التطوير المتكاملة (IDE)، افتح لوحة Google Cloud Data Agent Kit.
  2. وسِّع قسم الكتالوج.
  3. وسِّع رقم تعريف مشروعك.
  4. وسِّع BigQuery.
  5. وسِّع مجموعة بيانات transactions_dataset_evals.
  6. انقر على الجدول raw_transactions لفتح عرض التفاصيل في المحرِّر الرئيسي.
  7. في شريط التنقّل الأيمن، استكشِف علامات التبويب البيانات والمخطط والتفاصيل لفحص السجلات والبيانات الوصفية التي تمّت إضافتها.

ملخّص القسم: استخدمت اللغة الطبيعية في Agent Chat لإنشاء عبء عمل كامل في Spark Serverless. بعد ذلك، نفّذتَها لمعالجة سجلّات JSON غير المنظَّمة في جدول BigQuery (الأولي).

4. إزالة التكرار وتوحيد التنسيق باستخدام dbt

قبل تدريب نموذج تعلُّم الآلة، عليك فرض جودة البيانات من خلال إزالة سجلّات البث المكرّرة وعزل السجلّات غير الصالحة (مثل أرقام تعريف المعاملات الفارغة) ودمج البيانات المتعدّدة الأبعاد (الدافعين والمدفوع لهم). تتطلّب هذه العملية تحويلات SQL متكرّرة وموثوقة، ما يجعل dbt (أداة إنشاء البيانات) خيارًا مناسبًا.

إنشاء بنية أساسية لمسار dbt

استخدِم الوكيل لإنشاء مشروع dbt على مجموعة بيانات BigQuery:

  1. ارجع إلى جزء محادثة مع الوكيل.
  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. انقر على متابعة (ثم على قبول الكل) للسماح للوكيل بإنشاء الملفات في مساحة عملك.

خطة التنفيذ مع زر &quot;متابعة&quot;

  1. بعد اكتمال عملية الإنشاء، يعرض البرنامج التعليمي جولة إرشادية تلخّص المكوّنات الجديدة. اقبل كل التغييرات إذا طُلب منك ذلك.

قبول جميع الملفات التي تم إنشاؤها في &quot;لوحة Chat&quot;

إنشاء الاختبار وتنفيذه

على الرغم من أنّ الوكيل نفّذ dbt compile تلقائيًا لضمان صحة بنية SQL التي تم إنشاؤها، عليك الآن تحويل طرق العرض والجداول هذه إلى BigQuery وإجراء اختبارات جودة البيانات للتحقّق منها محليًا. (ملاحظة: في وقت لاحق من الدرس التطبيقي، ستتمكّن من إنجاز خطوة dbt هذه بشكل آلي كجزء من مخطّط موجه غير دوري (DAG) شامل في Airflow).

  1. في شريط الأنشطة في أقصى اليمين، انقر على رمز المستكشف (أو اضغط على Cmd/Ctrl+Shift+E).
  2. وسِّع dbt_project -> models لفحص نماذج SQL التي تم إنشاؤها. انقر على enriched_transactions.sql لفتح ومراجعة منطق ميزة التحويل ومكافحة الاحتيال في المحرِّر.
  3. في "مستكشف الملفات"، انقر بزر الماوس الأيمن على المجلد dbt_project واختَر الفتح في "وحدة طرفية مدمجة" (Open in Integrated Terminal). يؤدي ذلك إلى فتح لوحة طرفية تلقائيًا يتم ضبطها مباشرةً على 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 واختباره في &quot;وحدة التحكّم المدمجة&quot;

  1. بعد انتهاء عملية الإنشاء، أغلِق لوحة الوحدة الطرفية لإتاحة مساحة على الشاشة للخطوات المتبقية.

ملخّص القسم: أنشأت مشروع dbt باستخدام الوكيل، وأجريت اختبارات جودة البيانات، وحوّلت السجلات الأولية إلى جداول BigQuery مرحلية ومحسّنة.

5- تدريب نموذج موزّع لرصد الاحتيال باستخدام "الغابة العشوائية"

بعد إعداد المعاملات المحسّنة في BigQuery، ستنشئ نموذج تعلُّم آلة لتصنيف الأحداث الاحتيالية. Random Forest هي طريقة تعلُّم جماعي مناسبة لبيانات التصنيف الجدولية. يؤدي تشغيل RandomForestClassifier على Spark Serverless إلى توزيع تدريب النموذج على عُقد عاملة بدون الحاجة إلى إدارة البنية الأساسية.

في هذه الخطوة، ستستخدم الوكيل لإنشاء مسار تدريب Spark ML.

إنشاء دفتر ملاحظات لتدريب تعلُّم الآلة

  1. افتح لوحة محادثة مع وكيل الدعم.
  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 لترميز الميزات وتجميع المتجهات ومنطق تصنيف "الغابة العشوائية".
  3. انقر على تنفيذ الكل في شريط أدوات ورقة الملاحظات في بيئة التطوير المتكاملة.
  4. عند فتح أداة اختيار القائمة المنسدلة اختيار النواة، اختَر fraud-pipeline-runtime على Serverless Spark.

اختيار نواة Serverless Spark لدفتر الملاحظات التدريبي

التحقق

بعد اكتمال التنفيذ، تأكَّد من تدريب النموذج وتصديره بشكلٍ صحيح:

  1. راجِع نواتج خلية التقييم بالقرب من أسفل دفتر الملاحظات للتحقّق من نتيجة "المساحة تحت منحنى ROC" (AUC) المُبلَغ عنها.
  2. لضمان حفظ عناصر النموذج بنجاح في GCS، وسِّع لوحة مستكشف التخزين في الشريط الجانبي لـ "حزمة أدوات وكيل البيانات".
  3. ابحث عن الحزمة التي تنتهي بـ -models (مرتبطة بمعرّف المشروع النشط)، ووسِّعها، وانتقِل إلى التفاصيل للتحقّق من وجود الدليل fraud_model ومراحل خط الأنابيب.

التأكّد من حفظ النموذج في GCS

ملخّص القسم: استخدمتَ الوكيل لإنشاء مسار تدريب نموذج تعلُّم الآلة في PySpark، ودربتَ نموذج "الغابة العشوائية" على جدول BigQuery المحسّن، وصدّرتَ النموذج إلى Cloud Storage.

6. الاستنتاج المجمّع والكتابة إلى Cloud Spanner

باستخدام نموذج تنبؤي مدرَّب ومخزَّن في Cloud Storage، ستنفّذ استنتاجًا مجمّعًا على المعاملات الجديدة التي يتم إجراؤها من خلال BigQuery. يجب توجيه المعاملات عالية الخطورة إلى نظام تشغيلي ليتمكّن فريق الامتثال من مراجعتها. توفّر Cloud Spanner قاعدة بيانات معاملات قابلة للتوسّع لقائمة المراجعة هذه.

إنشاء دفتر ملاحظات للاستدلال المجمّع

استخدِم الوكيل لإنشاء دفتر ملاحظات للاستدلال يربط بين BigQuery وCloud Storage وCloud Spanner:

  1. افتح لوحة محادثة مع وكيل الدعم.
  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 في مساحة عملك.

الوكيل الذي ينشئ دفتر الملاحظات الخاص بالاستدلال

مراجعة دفتر الملاحظات وتنفيذه

  1. افتح notebooks/03_inference.ipynb الذي تم إنشاؤه حديثًا في المحرِّر.
  2. راجِع تسلسل الاستنتاج في PySpark:
    • الملفات التابعة: يوفّر نموذج Serverless Runtime ملفات cloud-spanner JAR التابعة المطلوبة لتنفيذ Spark.
    • تنسيق البيانات: يسقط النص البرمجي أعمدة متّجهة معقّدة في Spark ML (مثل الميزات الأولية والاحتمالات) قبل الكتابة لتتطابق مع مخطط جدول Spanner.
    • Spanner Connector: يكتب هذا الموصّل الصفوف التي تم الإبلاغ عنها باستخدام .format("cloud-spanner") لإضافتها مباشرةً إلى قائمة انتظار المراجعة.
  3. انقر على تنفيذ الكل في شريط أدوات ورقة الملاحظات في بيئة التطوير المتكاملة.
  4. عندما يُطلب منك اختيار نواة، اختَر fraud-pipeline-runtime على Serverless Spark.

التحقق

بعد انتهاء معالجة دفتر ملاحظات الاستدلال، يمكنك طلب البحث في قاعدة بيانات Spanner التشغيلية مباشرةً داخل بيئة التطوير المتكاملة (IDE) باتّباع الخطوات التالية:

  1. في شريط أنشطة بيئة التطوير المتكاملة (IDE)، افتح لوحة Google Cloud Data Agent Kit.
  2. وسِّع قسم الكتالوج.
  3. وسِّع رقم تعريف مشروعك، ثم وسِّع Spanner.
  4. انتقِل إلى cymbal-fraud -> fraud-db -> Tables -> SparkEvalFraudReviewQueue.
  5. انقر بزر الماوس الأيمن على الجدول واختَر جدول الاستعلام، ثم نفِّذ الاستعلام:
SELECT *
FROM `SparkEvalFraudReviewQueue`
LIMIT 100;
  1. في لوحة نتائج الاستعلام أدناه، من المفترض أن ترى صفوفًا تم إدراجها حديثًا وتم الإبلاغ عنها على أنّها معاملات عالية الخطورة وتتطلّب مراجعة يدوية.

التحقّق من الصفوف في Cloud Spanner

ملخّص القسم: استخدمتَ الوكيل لإنشاء دفتر ملاحظات للاستدلال المجمّع، وسجّلتَ نتائج سجلّات BigQuery غير المصنّفة باستخدام النموذج المدرَّب، وكتبتَ المعاملات عالية الخطورة مباشرةً إلى Cloud Spanner.

7. إنشاء بنية أساسية وتنظيمها باستخدام Managed Airflow

يتألف مسار البيانات حاليًا من خطوات منفصلة: دفتر ملاحظات خاص باستيعاب البيانات، ومشروع تحويل dbt، ودفتر ملاحظات خاص بالاستدلال المجمَّع. لجعل هذا الإنتاج جاهزًا، عليك ربطها معًا في مخطط تبعية مجدوَل.

توفّر خدمة Managed Service for Apache Airflow (المعروفة سابقًا باسم Cloud Composer) محرّك تنسيق مُدارًا لسير العمل هذا. تتضمّن "حزمة وكيل البيانات" ميزة خطوط أنابيب التنسيق التي تحوّل تعريفات خطوط أنابيب 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.

مراجعة إعدادات مجموعة توجيه الأسيليكليك

يستخدم Data Agent Kit Orchestrator إعدادات YAML تعريفية لتحديد خطوط الأنابيب ونشرها في Apache Airflow، ما يسمح بالتحكّم في إصدارات التعريفات ونشرها من خلال CI/CD.

في لوحة IDE 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، وتضمين سمة ملف JAR الخاص بـ Spanner.
  3. سيقدّم الوكيل أيضًا ملخّصًا لهذه العناصر التي تم إنشاؤها في علامة التبويب جولة إرشادية في لوحة المحرّر، مع توضيح الإعدادات وعمليات التحقّق التي تم إجراؤها.

إعدادات DAG التفاعلية

تعرض "حزمة وكيل البيانات" إعدادات مسار البيانات على شكل رسم بياني مرئي تفاعلي لفحص خصائص DAG في Airflow وتعديلها.

  1. في شريط أنشطة بيئة التطوير المتكاملة (IDE)، افتح لوحة Google Cloud Data Agent Kit.
  2. ضمن DATA ENGINEERING، وسِّع Orchestration Pipelines.
  3. انقر على fraud_analysis_pipeline.yaml لفتح لوحة DAG المرئية في المحرِّر الرئيسي.

لوحة العرض المرئية لـ DAG الخاصة بالتنظيم

  1. انقر على العقدة Schedule trigger في أعلى الصفحة. يتم فتح قائمة منبثقة خاصة بالإعدادات على يسار الصفحة، تعرض سلسلة Cron التي تم تحليلها (0 0 * * *) وتتيح لك تعديل مَعلمات مثل "تعبئة البيانات السابقة" و"تعبئة البيانات المتأخرة".
  2. انقر على أي من عُقد مهام دفتر الملاحظات (مثل خطوة الاستيعاب أو الاستنتاج). يتم تعديل النافذة المنبثقة لعرض عمليات الربط المحدّدة لتنفيذ Dataproc Serverless وخصائص الموصل.
  3. لاحظ الرابط التشعّبي لاسم دفتر الملاحظات (مثل 01_ingestion.ipynb) داخل كتلة العقدة. سيؤدي النقر عليه إلى فتح دفتر الملاحظات مباشرةً في المحرّر.
  4. في الشريط الجانبي الأيمن ضمن "خطوط سير العمل"، انقر على Deployment configuration. تعرض طريقة العرض هذه مجموعة بيئة dev المستهدَفة ونتائج حزمة GCS.

ملخّص القسم: أنشأت إعدادات خط أنابيب التنسيق باستخدام الوكيل، وحدّدت التبعيات بين مهام الاستيعاب وdbt والاستدلال في لوحة عرض مرئية تفاعلية.

8. التفعيل والتنفيذ والمراقبة

بعد تحديد DAG محليًا، ستتصل ببيئة Managed Airflow التي تم توفيرها أثناء عملية الإعداد وستنفّذ خط الأنابيب.

إعداد Managed Service for Apache Airflow

قبل النشر، اضبط عملية ربط "نظام جدولة المهام" في إعدادات "حزمة وكيل البيانات" لكي يستهدف الامتداد بيئة Airflow المُدارة:

  1. في شريط أنشطة بيئة التطوير المتكاملة (IDE)، افتح لوحة Google Cloud Data Agent Kit.
  2. ضمن SETTINGS، انقر على الإعدادات.
  3. انقروا على المجدول من القائمة اليمنى.
  4. اضبط الإعدادات:
    • رقم تعريف المشروع: اختَر رقم تعريف مشروعك النشط.
    • المنطقة: انقر على us-central1.
    • البيئة: اختَر cymbal-airflow.
  5. انقر على حفظ.

إعدادات Managed Service for Apache Airflow

تفعيل الرسم البياني الدوراني الموجّه

ستنشر الآن مسار البيانات الذي تم إعداده مباشرةً في بيئة Managed Airflow من لوحة العرض المرئية:

  1. في الشريط الجانبي مجموعة أدوات Google Cloud Data Agent، وسِّع 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 ENGINEERING > Orchestration Pipelines.
  2. انقر على إدارة خطوط الإنتاج.
  3. في جدول "إدارة خطوط الإنتاج"، انقر على fraud_analysis_pipeline لفتح سجلّ التنفيذ.

نظرة عامة على إدارة خطوط الإنتاج

  1. في طريقة عرض سجلّ التنفيذ، اختَر عملية التشغيل النشطة من التقويم.
  2. أثناء تقدّم التنفيذ في كل مهمة من مهام خط الأنابيب (النقل والتحويل باستخدام dbt والاستنتاج)، يتم تعديل مؤشرات الحالة وتعبئة مدة المهمة. انقر على أي مهمة لفحص نتائج تنفيذها المباشر وسجلات DAG في Airflow.

سجلّ تنفيذ البنية الأساسية لبرنامج العرض المباشر وتفاصيل المهام

ملخّص القسم: لقد أعددت اتصال Airflow Scheduler، ونشرت مسار التحليل المتكامل إلى Managed Airflow، وراقبت عملية تنفيذ مباشرة، وتحقّقت من النظام من السجلات الأولية إلى توقّعات Cloud Spanner النهائية.

9. تَنظيم

لتجنُّب تحمّل رسوم مستمرة على مشروعك على السحابة الإلكترونية في Google Cloud مقابل الموارد المستخدَمة في هذا الدرس التطبيقي حول الترميز، عليك إيقاف البيئة باستخدام النص البرمجي التلقائي.

  1. في لوحة الوحدة الطرفية (أو في Cloud Shell)، انتقِل إلى دليل النصوص البرمجية ونفِّذ ما يلي:
cd ~/devrel-demos/codelabs/agentic-data-labs/data-science/scripts
chmod +x teardown.sh
./teardown.sh
  1. سيعرض النص البرمجي قائمة بجميع الموارد التي يخطّط لحذفها ويطلب تأكيد ذلك:
    • Managed Airflow Environment (cymbal-airflow)
    • مثيل Cloud Spanner (cymbal-fraud)
    • مجموعة بيانات BigQuery (transactions_dataset_evals)
    • حِزم Cloud Storage (gs://${PROJECT_ID}-fin-clearing-raw وgs://${PROJECT_ID}-models)
    • حساب خدمة المنفِّذ (composer-worker-sa)
  2. يُرجى كتابة y للتأكيد. سيزيل نص البرمجة الخاص بإيقاف التشغيل جميع خدمات Google Cloud Platform المتوفّرة وينظّف الملفات المحلية.

10. تهانينا!

لقد أنشأت مسارًا متكاملاً لرصد الاحتيال يشمل Cloud Storage وBigQuery و"الخدمة المُدارة لـ Apache Spark" (Spark Serverless) وdbt وCloud Spanner و"الخدمة المُدارة لـ Apache Airflow"، بالإضافة إلى برمجة ثنائية مع "حزمة أدوات وكيل البيانات" من Google Cloud داخل بيئة التطوير المتكاملة Antigravity IDE.

إنجازاتك

  1. 📥 تمّت إضافة سجلّات المعاملات الأولية إلى جدول BigQuery باستخدام Managed Service for Apache Spark وData Agent Kit.
  2. ‫🧹 إزالة البيانات المكررة وتوحيد تنسيقها من خلال إنشاء مشروع dbt يتضمّن اختبارات جودة البيانات
  3. 🤖 تم تدريب نموذج Random Forest موزَّع باستخدام RandomForestClassifier وتصدير النموذج المُدرَّب إلى Cloud Storage.
  4. ⚡ تنفيذ استنتاج مجمّع بشأن المعاملات الواردة وتوجيه السجلات عالية الخطورة إلى Cloud Spanner لمراجعتها.
  5. 🔄 تنظيم سير العمل ونشره ومراقبته كرسم بياني موجّه غير دوري (DAG) مجدول في Airflow باستخدام Managed Service for Apache Airflow وأدوات إدارة الرسومات البيانية الموجّهة غير الدورية المرئية في بيئة التطوير المتكاملة

المفاهيم الرئيسية

المفهوم

ما تعلّمته

مجموعة أدوات وكيل البيانات

البرمجة الثنائية داخل بيئة التطوير المتكاملة باستخدام اللغة الطبيعية لإنشاء دفاتر ملاحظات PySpark وإعداد نماذج dbt وتحديد مخططات DAG في Airflow

BigQuery

تخزين جدولي قابل للتوسّع لـ SQL التحليلي وعمليات تحويل dbt وتدريب نماذج تعلُّم الآلة

Spark Serverless

التنفيذ بدون خادم لتحميل البيانات الموزّعة باستخدام PySpark والتدريب على تعلُّم الآلة باستخدام Random Forest

Cloud Spanner Connector

كتابة توقّعات استنتاج Spark المجمّعة مباشرةً في قوائم انتظار المراجعات في قاعدة البيانات التشغيلية

تعريفات YAML DAG

تعريفات مسار البيانات التوضيحية معروضة كرسومات بيانية مرئية تفاعلية في Airflow ضمن بيئة التطوير المتكاملة

إدارة DAG المرئية

فحص تبعيات مسار البيانات، وتفعيل النشر في Managed Airflow، ورصد سجلّ تنفيذ المهام المباشرة داخل بيئة التطوير المتكاملة

الخطوات التالية