انتقل إلى المحتوى الرئيسي

الوحدة 9 — Vertex AI Pipelines

الوحدات من 3 إلى 8 غطّت خطوات مشروع كشف الاحتيال، لكن كلّها أُطلقت يدويًّا من الدفتر. في الإنتاج، هذا لا يصمد أسبوعًا: خطأ بشريّ يفوّت خطوة، ترتيب خاطئ يستعمل بيانات قديمة، ولا يبقى أثر يُوضِّح ما نُشِر ومن أين أتى. Vertex AI Pipelines يُنسِّق هذه الخطوات في رسم بيانيّ موجَّه لا دوريّ (DAG) قابل للتشغيل من زرّ واحد أو من مُجدوِل، مع لسانيّة كاملة على كلّ الآثار.

KFP في لمحة

Vertex Pipelines يشغّل تصنيفات Kubeflow Pipelines (KFP). المفاهيم الأربعة:

  • مكوِّن (component): وحدة تنفيذ صغيرة، دالّة Python (أو حاوية كاملة) بمُدخلات ومخرجات موصوفة النوع.
  • مسار (pipeline): رسم بيانيّ يربط عدّة مكوّنات، بتدفّق البيانات (مخرج مكوّن يصبح مدخل الآخر).
  • أثر (artifact): كائن قابل للتخزين والاستعمال بين المكوّنات (مجموعة بيانات، نموذج، مقاييس تقييم).
  • مُنسِّق التنفيذ: Vertex نفسه، يستقبل التصنيف، يُشغِّل كلّ مكوّن على آلته الخاصّة، ويُخزِّن الآثار على GCS.

كتابة مكوّن أساسيّ

from kfp import dsl
from kfp.dsl import component, Input, Output, Dataset, Model, Metrics

@component(
base_image="python:3.11-slim",
packages_to_install=["google-cloud-bigquery", "pandas", "pyarrow"],
)
def extract_features(
project: str,
date: str,
features_out: Output[Dataset],
):
from google.cloud import bigquery
import pandas as pd

client = bigquery.Client(project=project)
df = client.query(f"""
SELECT ...
FROM `{project}.transactions.raw`
WHERE DATE(event_ts) = DATE('{date}')
""").to_dataframe(create_bqstorage_client=True)

df.to_parquet(features_out.path + ".parquet", index=False)
features_out.metadata["rows"] = len(df)
features_out.metadata["date"] = date

الملاحظات:

  • @component يُحوِّل الدالّة إلى مكوّن قابل للتضمين، مع صورة أساس ومكتبات للتثبيت.
  • Output[Dataset]: مخرج مُحدَّد نوع Dataset، يحصل تلقائيًّا على مسار GCS خالد (gs://.../pipeline_root/.../features_out/) وعلى ما وراء بيانات تُقرَأ في المُستهلك.
  • features_out.metadata: مفاتيح حرّة تُضاف إلى الأثر، تُقرأ في الواجهة وتُبنى عليها الشروط.

مكوّن التدريب بنية مماثلة، يأخذ Input[Dataset] ويُنتج Output[Model]:

@component(
base_image="europe-docker.pkg.dev/vertex-ai/training/xgboost-gpu.1-6:latest",
packages_to_install=["pandas", "pyarrow", "scikit-learn"],
)
def train_xgb(
features_in: Input[Dataset],
max_depth: int,
learning_rate: float,
model_out: Output[Model],
metrics_out: Output[Metrics],
):
import pandas as pd, xgboost as xgb
from sklearn.metrics import roc_auc_score, average_precision_score
df = pd.read_parquet(features_in.path + ".parquet")
features = [c for c in df.columns if c != "y"]
# ... تقسيم، تدريب، حفظ ...
metrics_out.log_metric("val_aucpr", 0.72)
metrics_out.log_metric("val_rocauc", 0.94)

تجميع المكوّنات في مسار

@dsl.pipeline(
name="fraud-training-pipeline",
pipeline_root="gs://fraud-pipelines-lab/pipeline_root",
)
def fraud_pipeline(project: str, date: str):
features = extract_features(project=project, date=date)

trained = train_xgb(
features_in=features.outputs["features_out"],
max_depth=6,
learning_rate=0.08,
)

eval_task = evaluate_model(
model_in=trained.outputs["model_out"],
test_features=features.outputs["features_out"],
)

with dsl.If(eval_task.outputs["passes_gate"] == "true", name="deploy-if-good"):
register_and_deploy(
model_in=trained.outputs["model_out"],
endpoint_name="fraud-endpoint",
)

المسار يعرِّف تدفّق الآثار بين المكوّنات. Vertex يستنبط الاعتماديّات: التقييم لا يبدأ قبل انتهاء التدريب، والنشر لا يبدأ قبل التقييم.

التسليم بشرط (dsl.If) يحمي الإنتاج: لن يُنشَر النموذج إذا لم يعبر التقييم عتبةً. لا رفع نموذج «سيّئ ولكن مقبول لأنّ الأسبوع كان صعبًا».

التصنيف والتنفيذ

from kfp import compiler
from google.cloud import aiplatform

compiler.Compiler().compile(
pipeline_func=fraud_pipeline,
package_path="fraud_pipeline.json",
)

aiplatform.init(project="fraud-detection-lab", location="europe-west1")

run = aiplatform.PipelineJob(
display_name="fraud-training-2026-09-05",
template_path="fraud_pipeline.json",
parameter_values={"project": "fraud-detection-lab", "date": "2026-09-05"},
pipeline_root="gs://fraud-pipelines-lab/pipeline_root",
enable_caching=True,
)
run.submit(service_account="sa-pipeline@fraud-detection-lab.iam.gserviceaccount.com")

enable_caching=True هو ذهب Vertex Pipelines: إن أُعيد تشغيل المسار بنفس المدخلات، Vertex يعيد استعمال مخرجات المكوّنات السابقة عوض إعادة تنفيذها. تجربة learning_rate جديد لا تُعيد استخراج البيانات إذا لم يتغيّر التاريخ.

اللسانيّة الكاملة

بعد كلّ تشغيل، تبويبة «Lineage» في الواجهة تعرض الرسم البيانيّ كاملًا:

  • جدول BigQuery المصدر
  • features Dataset (على GCS، بمعرف أثر) ⇒
  • model Model (على GCS) ⇒
  • metrics Metrics (JSON) ⇒
  • الإصدار في Model Registry
  • النشر على Endpoint

كلّ عقدة قابلة للنقر: مقاييسها، منشئها، تاريخها، ومكوّن المسار الذي أنتجها. حين يسأل مدير التنظيم «من أين جاء الإصدار المنشور اليوم؟»، الجواب في تبويبة واحدة، بلا شفهيّات.

الآثار والفروق مع «الملفّات»

كلّ أثر (Artifact) له:

  • URI: مسار GCS الفعليّ (يستطيع Vertex فتحه في أيّ مكوّن).
  • معرِّف خالد: يبقى ما دام المسار موجودًا.
  • ما وراء بيانات: مفاتيح مضافة يدويًّا (metadata) أو تلقائيًّا (المقاييس مثلًا).
  • نوع موصوف: Dataset, Model, Metrics, HTML, Markdown، مع كلاسات مخصّصة لك.

الفارق مع «إنزال ملفّ إلى GCS» جوهريّ: الأثر يحمل هويّة نوعيّة يعرفها Vertex، فتُظهرها الواجهة، تعرضها اللسانيّة، وتربطها بالمكوّنات المستهلِكة. ملفّ عاديّ على GCS يظلّ ملفًّا صامتًا.

الجدولة

خطّ يوميّ ليلة كلّ يوم:

from google.cloud.aiplatform.pipeline_job_schedules import PipelineJobSchedule

sched = PipelineJobSchedule(
pipeline_job=run,
display_name="fraud-nightly",
)
sched.create(
cron="0 2 * * *", # 2 صباحًا يوميًّا
max_concurrent_run_count=1,
max_run_count=None, # بلا حدّ
)

Vertex Scheduler يستقبل المسار المُصنَّف، يشغِّله كلّ ليلة، ويسجّل التنفيذات في تبويبة «Runs». التنفيذ الفاشل يظهر بلون أحمر، مع إمكانيّة إعادة تشغيل خطوة واحدة (retry) بدل المسار كاملًا.

الاعتبارات العمليّة

  • قِسْ الآلات لكلّ مكوّن: استخراج البيانات يحتاج n1-standard-2؛ التدريب يحتاج n1-standard-8 مع T4. لا تُعمِّم آلة قويّة على المسار كلّه، فوترتك سترتفع للا شيء.

    train_task = train_xgb(...).set_cpu_limit("8").set_memory_limit("32G").add_node_selector_constraint("cloud.google.com/gke-accelerator", "nvidia-tesla-t4").set_gpu_limit(1)
  • حدّ زمنيّ لكلّ مكوّن: مكوّن معلَّق يُنفَق ثمنه. set_display_name، set_retry، وset_timeout أدوات ضروريّة.

  • مسار جذر واحد للمشروع: pipeline_root هو دلو GCS الذي يحمل كلّ الآثار. جعله خارج المشاريع التجريبيّة يمنع تخبّطها في نتائج بعضها.

عدم إعادة الاستعمال في وجه التغيّر

عندما يتغيّر مصدر البيانات لكن التاريخ يبقى نفسه، enable_caching قد يعيد استعمال مخرج قديم غير مقصود. الاحتياط: أضف بصمة محتوى (content hash) إلى مدخلات المكوّنات الحسّاسة، أو أطفئ التخزين المؤقّت للمكوّنات القريبة من مصدر البيانات.

المسار البسيط أفضل من المُعقَّد

لا تُدخِل ثلاثين مكوّنًا في المسار الأوّل. ابدأ بأربعة: استخراج، تدريب، تقييم، نشر. أضف الشروط والتحويلات لاحقًا حين تصير الحاجة صريحة، لا افتراضًا. مسار عشرين خطوة يُصعِّب التصحيح على من يرثه.

في الخلاصة

  • KFP يُحوِّل خطوات المشروع إلى مكوّنات، والمكوّنات إلى مسار قابل لإعادة التنفيذ.
  • الآثار كائنات نوعيّة تُخزَّن على GCS مع ما وراء بيانات ولسانيّة قابلة للاستفسار.
  • dsl.If يجعل النشر مشروطًا بعتبات جودة موصوفة، فيصير تسليم الإنتاج قرارًا آليًّا لا بشريًّا.
  • الجدولة والتخزين المؤقّت يخفضان تكلفة إعادة التشغيل ويجعلان المسار الليليّ ممكنًا اقتصاديًّا.

الوحدة التالية والأخيرة: استدعاء نموذج أساس من Model Garden لتصنيف تعليقات المتنازعين، مع احتساب الكلفة بالرمز.