תיאום של צינורות עיבוד נתונים

בדף הזה מוסבר על תזמור צינורות נתונים באמצעות Managed Service for Apache Airflow וטריגרים. מומלץ להשתמש ב-Managed Airflow כדי לתזמן צינורות עיבוד נתונים ב-Cloud Data Fusion. אם אתם צריכים דרך פשוטה יותר לנהל תזמור, אתם יכולים להשתמש בטריגרים.

מלחין/ה

תזמור צינורות נתונים באמצעות Managed Airflow

היתרונות של תזמור הפעלת צינורות ב-Cloud Data Fusion באמצעות Managed Airflow:

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

כדי לתזמן צינורות עיבוד נתונים של Cloud Data Fusion באמצעות Managed Airflow, פועלים לפי התהליך הבא:

  1. מגדירים את סביבת Managed Airflow.

    • יוצרים סביבת Managed Airflow. אם אין לכם פרויקט, אתם יכולים להקצות את הסביבה ב Google Cloud פרויקט. הסביבה הזו היא סביבת העבודה שלכם לניהול תהליכים.
    • מתן הרשאות. מוודאים שלחשבון השירות של Managed Airflow יש את ההרשאות הנדרשות לגישה ל-Cloud Data Fusion (למשל, הרשאה להפעלה, לעצירה ולרישום של צינורות).
  2. הגדרת גרפים אציקליים מכוונים (DAG) לתיאום.

    • יוצרים DAG: ב-Managed Airflow, יוצרים DAG שמגדיר את תהליך העבודה של תזמור צינורות Cloud Data Fusion.
    • אופרטורים של Cloud Data Fusion: אפשר להשתמש באופרטורים של Cloud Data Fusion ב-Managed Airflow בתוך DAG. האופרטורים האלה מאפשרים לכם ליצור אינטראקציה עם Cloud Data Fusion באופן פרוגרמטי.

אופרטורים של Cloud Data Fusion

לתיזמור של צינורות נתונים ב-Cloud Data Fusion יש את האופרטורים הבאים:

CloudDataFusionStartPipelineOperator

מפעיל את הביצוע של פייפליין Cloud Data Fusion לפי המזהה שלו. הוא כולל את הפרמטרים הבאים:

  • מזהה הפייפליין
  • מיקום (אזורGoogle Cloud )
  • מרחב השמות של הפייפליין
  • ארגומנטים בזמן ריצה (אופציונלי)
  • המתנה להשלמה (אופציונלי)
  • זמן קצוב לתפוגה (אופציונלי)
CloudDataFusionStopPipelineOperator

מאפשרת להפסיק פייפליין פעיל של Cloud Data Fusion.

CloudDataFusionDeletePipelineOperator

מחיקת צינור ב-Cloud Data Fusion.

יצירת תהליך העבודה של תרשים ה-DAG

כשיוצרים את תהליך העבודה של DAG, חשוב להתייחס לנקודות הבאות:

  • הגדרת יחסי תלות: משתמשים במבנה DAG כדי להגדיר יחסי תלות בין משימות. לדוגמה, יכול להיות שיש לכם משימה שמחכה לסיום מוצלח של צינור עיבוד ב-namespace אחד לפני הפעלה של צינור עיבוד אחר ב-namespace אחר.
  • תזמון: אפשר לתזמן את ה-DAG כך שיפעל במרווחי זמן ספציפיים, למשל מדי יום או מדי שעה, או להגדיר שהוא יופעל באופן ידני.

מידע נוסף מופיע במאמר סקירה כללית על Managed Airflow.

טריגרים

ניהול צינורות עיבוד נתונים באמצעות טריגרים

טריגרים של Cloud Data Fusion מאפשרים להריץ באופן אוטומטי צינורות (pipeline) במורד הזרם אחרי השלמה (הצלחה, כשל או כל תנאי שצוין) של צינור אחד או יותר במעלה הזרם.

טריגרים שימושיים למשימות הבאות:

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

לדוגמה, יש לכם מערך נתונים שמכיל את כל המידע על המשלוחים של החברה. על סמך הנתונים האלה, אתם רוצים לענות על כמה שאלות עסקיות. לשם כך, יוצרים צינור אחד שמנקה את הנתונים הגולמיים על המשלוחים, שנקרא Shipments Data Cleaning. לאחר מכן יוצרים צינור שני, Delayed Shipments USA, שקורא את הנתונים המנוקים ומאתר את המשלוחים בארה"ב שהתעכבו מעבר לסף שצוין. אפשר להפעיל את הפייפליין Delayed Shipments USA ברגע שהפייפליין Shipments Data Cleaning שמוגדר במעלה הזרם מסתיים בהצלחה.

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

כדי לתזמן צינורות באמצעות טריגרים, פועלים לפי התהליך הבא:

  1. יצירת צינורות עיבוד נתונים במעלה הזרם ובמורד הזרם.

    • ב-Cloud Data Fusion Studio, מתכננים ומפעילים את צינורות עיבוד הנתונים שמרכיבים את שרשרת התיאום.
    • כדאי לחשוב איזו השלמה של צינור עיבוד נתונים תפעיל את צינור עיבוד הנתונים הבא (במורד הזרם) בתהליך העבודה.
  2. אופציונלי: העברת ארגומנטים של זמן ריצה לצינורות עיבוד נתונים במעלה הזרם.

  3. יוצרים טריגר לדואר נכנס בפייפליין במורד הזרם.

    • ב-Cloud Data Fusion Studio, עוברים לדף List. בכרטיסייה Deployed (פריסה), לוחצים על השם של צינור הנתונים במורד הזרם. מוצגת תצוגת הפריסה של צינור הנתונים.
    • בצד ימין של הדף, לוחצים על מפעילים של הודעות נכנסות. מוצגת רשימה של צינורות זמינים.
    • לוחצים על צינור הנתונים במעלה הזרם. בוחרים אחד או יותר ממצבי השלמת הצינור (Succeeds,‏ Fails או Stops) כשהתנאי הוא מתי צינור הנתונים במורד הזרם צריך לפעול.
    • אם רוצים שצינור עיבוד הנתונים במעלה הזרם ישתף מידע (שנקרא הגדרת מטען ייעודי) עם צינור עיבוד הנתונים במורד הזרם, לוחצים על Trigger config ופועלים לפי השלבים להעברת הגדרת מטען ייעודי כארגומנטים של זמן ריצה. אם לא, לוחצים על הפעלת הטריגר.
  4. בודקים את הטריגר.

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

העברת הגדרות של מטען ייעודי (payload) כארגומנטים של זמן ריצה

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

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

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

כדי להעביר את הגדרות המטען הייעודי (payload) כארגומנטים של זמן ריצה:

  1. אחרי שלוחצים על Trigger config (הגדרת הטריגר), יופיעו כל הארגומנטים של זמן הריצה שהגדרתם קודם עבור צינור הנתונים במעלה הזרם. בוחרים את הארגומנטים של זמן הריצה להעברה מצינור הנתונים במעלה הזרם לצינור הנתונים במורד הזרם כשהטריגר הזה מופעל.
  2. לוחצים על הכרטיסייה Plugin config (הגדרת התוסף) כדי לראות רשימה של מה שיועבר מצינור עיבוד הנתונים במעלה הזרם לצינור עיבוד הנתונים במורד הזרם כשהוא מופעל.
  3. לוחצים על Configure and Enable Trigger (הגדרת הטריגר והפעלתו).