ניפוי באגים בבעיות בתזמון משימות

Managed Airflow (דור 3) | Managed Airflow (דור 2) | Managed Airflow (דור 1 מדור קודם)

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

מבוא

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

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

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

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

מטרות

בקטע הזה מפורטים יעדים לדוגמאות במדריך הזה.

דוגמה: תזמון לא תקין וזמן אחזור שנגרם בגלל ריבוי משימות בו-זמניות

  • מעלים את ה-DAG לדוגמה שפועל כמה פעמים בו-זמנית ומאבחנים את התקלות בתזמון ואת בעיות ההשהיה באמצעות Cloud Monitoring.

  • מבצעים אופטימיזציה של קוד ה-DAG על ידי איחוד המשימות ומעריכים את ההשפעה על הביצועים.

  • כדאי לפזר את המשימות בצורה אחידה יותר לאורך זמן ולהעריך את ההשפעה על הביצועים.

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

דוגמה: שגיאות בניתוח DAG וחביון שנגרמו בגלל קוד מורכב

  • העלאה של DAG לדוגמה עם משתני Airflow ואבחון בעיות בניתוח באמצעות Cloud Monitoring.

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

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

עלויות

במדריך הזה השתמשנו ברכיבים הבאים של Google Cloud, והשימוש בהם כרוך בתשלום:

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

לפני שמתחילים

בקטע הזה מתוארות פעולות שצריך לבצע לפני שמתחילים את המדריך.

יצירה והגדרה של פרויקט

במדריך הזה תצטרכו Google Cloud פרויקט. מגדירים את הפרויקט באופן הבא:

  1. במסוף Google Cloud , בוחרים פרויקט או יוצרים פרויקט:

    מעבר לדף לבחירת הפרויקט

  2. מוודאים שהחיוב מופעל בפרויקט. איך בודקים אם החיוב מופעל בפרויקט

  3. כדי ליצור את המשאבים הנדרשים, צריך לוודא שלמשתמש בפרויקט יש את התפקידים הבאים: Google Cloud

    • אדמין של סביבה ואובייקטים באחסון (roles/composer.environmentAndStorageObjectAdmin)
    • אדמין ב-Compute (roles/compute.admin)

הפעלת ממשקי API בפרויקט

מפעילים את Managed Airflow API.

תפקידים שנדרשים להפעלת ממשקי API

כדי להפעיל ממשקי API, נדרשת ההרשאה serviceusage.services.enable. אם יצרתם את הפרויקט, סביר להניח שכבר יש לכם את ההרשאה הזו דרך התפקיד 'בעלים' (roles/owner). אחרת, תוכלו לקבל את ההרשאה הזו דרך התפקיד 'אדמין בממשק 'שימוש בשירות'' (roles/serviceusage.serviceUsageAdmin). איך מקצים תפקידים

להפעלת ה-API

יצירת סביבת Managed Airflow

איך יוצרים סביבת Managed Airflow (דור 2)

.

במסגרת יצירת הסביבה, מקצים את התפקיד Cloud Composer v2 API Service Agent Extension ‏(roles/composer.ServiceAgentV2Ext) לחשבון של סוכן השירות של Composer. ‫Managed Airflow משתמש בחשבון הזה כדי לבצע פעולות בפרויקט Google Cloud שלכם.

דוגמה: תקלה בכלי לתזמון ומשימה שנכשלה בגלל בעיות בתזמון המשימה

בדוגמה הזו מוסבר איך לנפות באגים בתזמן של משימות שגורם לזמן טעינה ארוך בגלל ריבוי משימות בו-זמניות.

העלאת ה-DAG לדוגמה לסביבה שלכם

מעלים את ה-DAG לדוגמה הבא לסביבה שיצרתם בשלבים הקודמים. במדריך הזה, ה-DAG נקרא dag_10_tasks_200_seconds_1.

ב-DAG הזה יש 200 משימות. כל משימה ממתינה שנייה אחת ומדפיסה Complete!‎ (הושלם!). ה-DAG מופעל אוטומטית אחרי ההעלאה. ‫Managed Airflow מריץ את ה-DAG הזה 10 פעמים, וכל ההרצות של ה-DAG מתבצעות במקביל.

import time
from datetime import datetime, timedelta

from airflow import DAG
from airflow.decorators import task


tasks_amount = 200
seconds = 1
minutes = 5

with DAG(
    dag_id=f"dag_10_tasks_{tasks_amount}_sec_{seconds}",
    start_date=datetime(2023, 11, 22, 20, 0),
    end_date=datetime(2023, 11, 22, 20, 49),
    schedule=timedelta(minutes=minutes),
    catchup=True,
) as dag:

    @task
    def create_subtasks(seconds: int) -> None:
        time.sleep(seconds)

    for i in range(tasks_amount):
        create_subtasks(seconds)

אבחון בעיות שקשורות לתקלה בתזמן ולכשל במשימה

אחרי שהרצת ה-DAG מסתיימת, פותחים את ממשק המשתמש של Airflow ולוחצים על dag_10_tasks_200_seconds_1 DAG. אפשר לראות ש-10 הפעלות של DAG הסתיימו בהצלחה, ובכל אחת מהן בוצעו 200 משימות בהצלחה.

בודקים את יומני המשימות של Airflow:

  1. נכנסים לדף Environments במסוף Google Cloud .

    מעבר אל Environments

  2. ברשימת הסביבות, לוחצים על שם הסביבה. הדף Environment details ייפתח.

  3. עוברים לכרטיסייה Logs (יומנים), ואז ל-All logs (כל היומנים) > Airflow logs (יומני Airflow) > Workers (עובדים) > View in Logs Explorer (הצגה ב-Logs Explorer).

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

ההיסטוגרמה של יומני העובדים של Airflow עם שגיאות ואזהרות
    שמסומנות בצבעים אדום וכתום
איור 1. היסטוגרמה של יומני עובדים ב-Airflow (לחצו כדי להגדיל)

ה-DAG לדוגמה הניב כ-130 אזהרות ו-60 שגיאות. לוחצים על עמודה שמכילה פסים צהובים ואדומים. יופיעו ביומנים חלק מהאזהרות והשגיאות הבאות:

State of this instance has been externally set to success. Terminating
instance.

Received SIGTERM. Terminating subprocesses.

worker: Warm shutdown (MainProcess).

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

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

בודקים את פרטי המעקב:

  1. עוברים לכרטיסייה מעקב ובוחרים באפשרות סקירה כללית.

  2. בודקים את הגרף של משימות Airflow.

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

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

  3. בודקים את הגרף עובדים פעילים:

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

    בתרשים Active workers (עובדים פעילים) אפשר לראות שה-DAG הפעיל את התכונה 'שינוי גודל אוטומטי' עד למגבלה המקסימלית המותרת של שלושה עובדים במהלך ההרצה של ה-DAG.

  4. גרפים של השימוש במשאבים יכולים להצביע על חוסר קיבולת בעובדי Airflow להרצת משימות בתור. בכרטיסייה Monitoring (מעקב), בוחרים באפשרות Workers (תהליכי עבודה) ומעיינים בתרשימים Total worker CPU usage (סה"כ שימוש במעבד של תהליכי העבודה) ו-Total worker memory usage (סה"כ שימוש בזיכרון של תהליכי העבודה).

    בתרשים של השימוש במעבד על ידי עובדי Airflow, השימוש במעבד עולה עד למגבלה המקסימלית.
    איור 4. תרשים של סך השימוש במעבד (CPU) של העובדים (לחצו כדי להגדיל)
    בתרשים של השימוש בזיכרון על ידי עובדי Airflow, השימוש בזיכרון עולה, אבל לא מגיע למגבלה המקסימלית
    איור 5. תרשים של השימוש הכולל בזיכרון של העובדים (אפשר ללחוץ כדי להגדיל)

    התרשימים מצביעים על כך שההפעלה של יותר מדי משימות בו-זמנית גרמה לחריגה ממגבלת המעבד. המשאבים היו בשימוש במשך יותר מ-30 דקות, שזה אפילו יותר מהמשך הכולל של 200 משימות ב-10 הפעלות של DAG שפועלות אחת אחרי השנייה.

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

איחוד המשימות

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

הדוגמה הבאה של DAG משנה את מספר המשימות בדוגמה הראשונית מ-200 ל-20, ומגדילה את זמן ההמתנה מ-1 ל-10 שניות כדי לחקות משימות מאוחדות יותר שמבצעות את אותה כמות עבודה.

מעלים את ה-DAG לדוגמה הבא לסביבה שיצרתם. במדריך הזה, ה-DAG נקרא dag_10_tasks_20_seconds_10.

import time
from datetime import datetime, timedelta

from airflow import DAG
from airflow.decorators import task


tasks_amount = 20
seconds = 10
minutes = 5

with DAG(
    dag_id=f"dag_10_tasks_{tasks_amount}_sec_{seconds}",
    start_date=datetime(2021, 12, 22, 20, 0),
    end_date=datetime(2021, 12, 22, 20, 49),
    schedule=timedelta(minutes=minutes),
    catchup=True,
) as dag:

    @task
    def create_subtasks(seconds: int) -> None:
        time.sleep(seconds)

    for i in range(tasks_amount):
        create_subtasks(seconds)

הערכת ההשפעה של משימות מאוחדות יותר על תהליכי התזמון:

  1. מחכים עד שהרצת ה-DAG תסתיים.

  2. בממשק המשתמש של Airflow, בדף DAGs, לוחצים על ה-DAG‏ dag_10_tasks_20_seconds_10. יוצגו 10 הפעלות של DAG, ולכל אחת מהן יהיו 20 משימות שהושלמו.

  3. נכנסים לדף Environments במסוף Google Cloud .

    מעבר אל Environments

  4. ברשימת הסביבות, לוחצים על שם הסביבה. הדף Environment details ייפתח.

  5. עוברים לכרטיסייה Logs (יומנים), ואז ל-All logs (כל היומנים) > Airflow logs (יומני Airflow) > Workers (עובדים) > View in Logs Explorer (הצגה ב-Logs Explorer).

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

    ההיסטוגרמה של יומני העובדים של Airflow עם שגיאות ואזהרות
    מראה את הירידה במספר השגיאות והאזהרות אחרי איחוד המשימות
    איור 6. היסטוגרמה של יומני עובדים ב-Airflow אחרי שהמשימות אוחדו (אפשר ללחוץ כדי להגדיל)

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

  6. בכרטיסייה Monitoring, בוחרים באפשרות Workers ובודקים את הגרפים.

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

    הגרף של משימות Airflow לאורך זמן מראה שהעלייה החדה במשימות Airflow נמשכה פרק זמן קצר יותר מבעבר.
    איור 7. גרף המשימות של Airflow אחרי איחוד המשימות (אפשר ללחוץ כדי להגדיל)

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

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

    בודקים את הגרפים של צריכת משאבי העובדים. למרות שההבדל בין המשאבים שנעשה בהם שימוש בדוגמה עם המשימות המאוחדות לבין הדוגמה הראשונית הוא די משמעותי, השימוש ב-CPU עדיין מגיע ל-70% מהמגבלה.

    בתרשים של השימוש במעבד (CPU) על ידי עובדי Airflow, השימוש במעבד עולה עד 70% מהמגבלה המקסימלית.
    איור 9. תרשים של סך השימוש במעבד של העובדים אחרי איחוד המשימות (לחצו כדי להגדיל)
    בתרשים של השימוש בזיכרון על ידי עובדי Airflow, השימוש בזיכרון עולה אבל לא מגיע למגבלה המקסימלית
    איור 10. גרף השימוש הכולל בזיכרון של העובדים אחרי שהמשימות אוחדו (אפשר ללחוץ כדי להגדיל)

חלוקת המשימות באופן שווה יותר לאורך זמן

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

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

במדריך הזה, תוכלו להוסיף פרמטרים ברמת ה-DAG ל-DAG dag_10_tasks_20_seconds_10 כדי לחלק את המשימות בצורה שווה יותר לאורך זמן:

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

  2. מוסיפים את הארגומנט max_active_tasks=5 למנהל ההקשר של DAG. הארגומנט הזה קובע את המספר המקסימלי של מופעי משימות שיכולים לפעול בו-זמנית בכל DAG.

מעלים את ה-DAG לדוגמה הבא לסביבה שיצרתם. במדריך הזה, ה-DAG נקרא dag_10_tasks_20_seconds_10_scheduled.py.

import time
from datetime import datetime, timedelta

from airflow import DAG
from airflow.decorators import task


tasks_amount = 20
seconds = 10
minutes = 5
active_runs = 1
active_tasks = 5


with DAG(
    dag_id=f"dag_10_tasks_{tasks_amount}_sec_{seconds}_runs_{active_runs}_tasks_{active_tasks}",
    start_date=datetime(2021, 12, 22, 20, 0),
    end_date=datetime(2021, 12, 22, 20, 49),
    schedule=timedelta(minutes=minutes),
    max_active_runs=active_runs,
    max_active_tasks=active_tasks,
    catchup=True,
) as dag:

    @task
    def create_subtasks(seconds: int) -> None:
        time.sleep(seconds)

    for i in range(tasks_amount):
        create_subtasks(seconds)

הערכת ההשפעה של חלוקת המשימות לאורך זמן על תהליכי התזמון:

  1. מחכים עד שהרצת ה-DAG תסתיים.

  2. נכנסים לדף Environments במסוף Google Cloud .

    מעבר אל Environments

  3. ברשימת הסביבות, לוחצים על שם הסביבה. הדף Environment details ייפתח.

  4. עוברים לכרטיסייה Logs (יומנים), ואז ל-All logs (כל היומנים) > Airflow logs (יומני Airflow) > Workers (עובדים) > View in Logs Explorer (הצגה ב-Logs Explorer).

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

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

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

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

אופטימיזציה של הגדרות הסביבה

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

מספר העובדים והמקבילות של העובדים

אתם יכולים לשנות את המספר המקסימלי של העובדים כדי ש-Managed Airflow ישנה באופן אוטומטי את גודל הסביבה שלכם במסגרת המגבלות שהוגדרו.

הפרמטר [celery]worker_concurrency מגדיר את המספר המקסימלי של משימות שעובד יחיד יכול לקחת מתור המשימות. שינוי הפרמטר הזה משנה את מספר המשימות שעובד יחיד יכול לבצע בו-זמנית. אפשר לשנות את אפשרות ההגדרה הזו של Airflow על ידי החלפתה. כברירת מחדל, בו-זמניות (concurrency) של ה-worker מוגדרת על סמך מספר המופעים של משימות מקבילות קלות משקל שה-worker יכול להכיל. כלומר, הערך שלו תלוי במגבלות של משאבי העובדים. ערך המקביליות של העובדים לא תלוי במספר העובדים בסביבה שלכם.

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

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

    במקרים כאלה, מומלץ להגדיל את המספר המינימלי של העובדים ואת מספר העובדים המקסימלי שפועלים בו-זמנית כדי להימנע מהגדלה מהירה מדי של מספר העובדים.

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

החשיבות של מקביליות

מתזמני Airflow שולטים בתזמון של הפעלות DAG ומשימות ספציפיות מ-DAG. אפשרות ההגדרה [core]parallelism Airflow קובעת כמה משימות מתזמן Airflow יכול להוסיף לתור של המבצע אחרי שכל התנאים המוקדמים של המשימות האלה מתקיימים.

מקביליות היא מנגנון הגנה של Airflow שקובע כמה משימות יכולות לפעול בו-זמנית לכל מתזמן, ללא קשר למספר העובדים. ערך המקביליות, כפול מספר המתזמנים באשכול, הוא המספר המקסימלי של מופעי משימות שהסביבה שלכם יכולה להוסיף לתור.

בדרך כלל, הערך של [core]parallelism הוא מכפלה של מספר מקסימלי של עובדים ושל [celery]worker_concurrency. הוא מושפע גם ממאגר. אפשר לשנות את אפשרות ההגדרה הזו של Airflow על ידי החלפתה. מידע נוסף על שינוי הגדרות Airflow שקשורות לשינוי גודל זמין במאמר שינוי גודל ההגדרה של Airflow.

מציאת הגדרות אופטימליות של סביבות

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

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

אתם יכולים לשנות את הפרמטרים של מקביליות, מספר ה-workers וה-concurrency של ה-workers כדי להריץ את dag_10_tasks_20_seconds_10 DAG בלי להגביל את המשימות הפעילות. בדוגמה הזו, ה-DAG מופעל 10 פעמים וכל הפעלה מכילה 20 משימות קטנות. אם רוצים להפעיל את כולם בו-זמנית:

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

  • ה-workers של Airflow צריכים להיות מסוגלים להריץ 20 משימות בו-זמנית, ולכן צריך להגדיר את worker concurrency ל-20.

  • ל-workers צריכים להיות מספיק מעבד וזיכרון כדי לטפל בכל המשימות. ההשפעה של מקביליות העובדים תלויה במעבד ובזיכרון של העובדים, ולכן תצטרכו לפחות worker_concurrency / 12 במעבד ו-least worker_concurrency / 8 בזיכרון.

  • תצטרכו להגדיל את רמת המקביליות כדי להתאים אותה לרמת המקביליות הגבוהה יותר של העובדים. כדי שעובדים יוכלו לבחור 20 משימות מהתור, הכלי לתזמון צריך לתזמן קודם את 20 המשימות האלה.

משנים את הגדרות הסביבה באופן הבא:

  1. נכנסים לדף Environments במסוף Google Cloud .

    מעבר אל Environments

  2. ברשימת הסביבות, לוחצים על שם הסביבה. הדף Environment details ייפתח.

  3. עוברים לכרטיסייה Environment configuration (הגדרת הסביבה).

  4. מחפשים את ההגדרה Resources (משאבים) > Workloads (עומסי עבודה) ולוחצים על Edit (עריכה).

  5. בקטע Worker, בשדה Memory, מציינים את מגבלת הזיכרון החדשה של עובדי Airflow. במדריך הזה, משתמשים ב-4GB.

  6. בשדה CPU, מציינים את מגבלת המעבד החדשה לעובדי Airflow. במדריך הזה, משתמשים ב-2 vCPU.

  7. שומרים את השינויים וממתינים כמה דקות עד שהעובדים של Airflow יופעלו מחדש.

לאחר מכן, מבטלים את ההגדרה של מקביליות ואת אפשרויות ההגדרה של Airflow לגבי מספר העובדים:

  1. עוברים לכרטיסייה Airflow Configuration Overrides (שינויים בהגדרות של Airflow).

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

  3. עוקפים את הגדרת ההקביליות:

    קטע מפתח ערך
    core parallelism 20
  4. לוחצים על Add Airflow Configuration Override (הוספת שינוי בהגדרות של Airflow) ומשנים את הגדרת ה-concurrency של העובד:

    קטע מפתח ערך
    celery worker_concurrency 20
  5. לוחצים על שמירה ומחכים עד שההגדרה של הסביבה תתעדכן.

מפעילים שוב את אותו DAG לדוגמה עם ההגדרות המותאמות:

  1. בממשק המשתמש של Airflow, עוברים לדף DAGs.

  2. מאתרים את dag_10_tasks_20_seconds_10 DAG ומוחקים אותו.

    אחרי שמוחקים את ה-DAG, ‏ Airflow בודק את תיקיית ה-DAG בדלי של הסביבה ומריץ שוב את ה-DAG באופן אוטומטי.

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

ההיסטוגרמה של יומני העובדים ב-Airflow עם שגיאות ואזהרות
        לא מציגה שגיאות ואזהרות אחרי שהוגדרה הסביבה
איור 12. היסטוגרמה של יומני עובדים ב-Airflow אחרי התאמת הגדרות הסביבה (לחיצה להגדלה)

הגדרות הסביבה והגדרות Airflow ממלאות תפקיד חשוב בתזמון המשימות, אבל אי אפשר להגדיל את ההגדרות מעבר לגבולות מסוימים.

מומלץ לבצע אופטימיזציה של קוד ה-DAG, לאחד משימות ולהשתמש בתזמון כדי לשפר את הביצועים והיעילות.

דוגמה: שגיאות בניתוח DAG וחביון בגלל קוד DAG מורכב

בדוגמה הזו, בודקים את זמן האחזור של ניתוח DAG לדוגמה שמדמה עודף של משתני Airflow.

יצירת משתנה Airflow חדש

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

  1. נכנסים לדף Environments במסוף Google Cloud .

    מעבר אל Environments

  2. בעמודה Airflow webserver, לוחצים על הקישור Airflow של הסביבה.

  3. עוברים אל ניהול > משתנים > הוספת רשומה חדשה.

  4. מגדירים את הערכים הבאים:

    • מפתח: example_var
    • ערך: test_airflow_variable

העלאת ה-DAG לדוגמה לסביבה שלכם

מעלים את ה-DAG לדוגמה הבא לסביבה שיצרתם בשלבים הקודמים. במדריך הזה, ה-DAG נקרא dag_for_loop_airflow_variable.

ה-DAG הזה מכיל לולאת for שפועלת 1,000 פעמים ומדמה עודף של משתני Airflow. בכל איטרציה, המשתנה example_var נקרא ונוצרת משימה. כל משימה מכילה פקודה אחת שמדפיסה את הערך של המשתנה.

from datetime import datetime
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from airflow.models import Variable

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2023, 2, 17),
    'retries': 0
}

dag = DAG(
    'dag_for_loop_airflow_variable',
    default_args=default_args,
    catchup=False,
    schedule="@daily"
)

for i in range(1000):
    a = Variable.get('example_var', 'N/A')
    task = BashOperator(
        task_id=f'task_{i}',
        bash_command='echo variable foo=${foo_env}',
        dag=dag,
        env={'foo_env': a}
    )

אבחון בעיות הניתוח

זמן הניתוח של DAG הוא משך הזמן שנדרש לתזמן של Airflow לקרוא קובץ DAG ולנתח אותו. כדי שהמתזמן של Airflow יוכל לתזמן משימות מ-DAG, הוא צריך לנתח את קובץ ה-DAG כדי לגלות את המבנה של ה-DAG והמשימות שהוגדרו.

אם ניתוח של DAG נמשך זמן רב, הוא צורך את הקיבולת של מתזמן המשימות ועשוי להפחית את הביצועים של הפעלות DAG.

כדי לעקוב אחרי זמן הניתוח של ה-DAG:

  1. מריצים את dags report הפקודה Airflow CLI ב-CLI של gcloud כדי לראות את זמן הניתוח של כל ה-DAG:

    gcloud composer environments run ENVIRONMENT_NAME \
        --location LOCATION \
        dags report
    

    מחליפים את מה שכתוב בשדות הבאים:

    • ENVIRONMENT_NAME: השם של הסביבה.
    • LOCATION: האזור שבו נמצאת הסביבה.
  2. בפלט של הפקודה, מחפשים את ערך משך הזמן של ה-DAG‏ dag_for_loop_airflow_variables. ערך גדול עשוי להצביע על כך ש-DAG הזה לא מיושם בצורה אופטימלית. אם יש לכם כמה DAG, תוכלו לזהות מהטבלה של הפלט אילו DAG כוללים זמן ניתוח ארוך.

    דוגמה:

    file                | duration       | dag_num | task_num | dags
    ====================+================+=========+==========+=====================
    /dag_for_loop_airfl | 0:00:14.773594 | 1       | 1000     | dag_for_loop_airflow
    ow_variable.py      |                |         |          | _variable
    /airflow_monitoring | 0:00:00.003035 | 1       | 1        | airflow_monitoring
    .py
    
    
  3. בודקים את זמני הניתוח של DAG במסוף Google Cloud :

    1. נכנסים לדף Environments במסוף Google Cloud .

    מעבר אל Environments

  4. ברשימת הסביבות, לוחצים על שם הסביבה. הדף Environment details ייפתח.

  5. עוברים לכרטיסייה Logs (יומנים), ואז ל-All logs (כל היומנים) > DAG processor manager (מנהל מעבדי DAG).

  6. בודקים את היומנים של dag-processor-manager ומזהים בעיות אפשריות.

    רשומה ביומן של DAG לדוגמה מראה שזמן הניתוח של ה-DAG הוא 46.3 שניות
    איור 13. יומני DAG processor manager מציגים את זמני הניתוח של DAG (לחצו כדי להגדיל)

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

אופטימיזציה של קוד ה-DAG

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

שדות התבניות של Airflow מאפשרים לכם לשלב ערכים ממשתני Airflow ומתבניות Jinja ב-DAG. כך נמנעת הפעלה מיותרת של פונקציות במהלך פעימות הלב של מתזמן המשימות.

כדי להטמיע את דוגמת ה-DAG בצורה טובה יותר, מומלץ להימנע משימוש במשתני Airflow בקוד Python ברמה העליונה של DAG. במקום זאת, מעבירים משתני Airflow לאופרטורים קיימים באמצעות תבנית Jinja, מה שיגרום לעיכוב בקריאת הערך עד להפעלת המשימה.

מעלים את הגרסה החדשה של ה-DAG לדוגמה לסביבה. במדריך הזה, ה-DAG נקרא dag_for_loop_airflow_variable_optimized.

from datetime import datetime
from airflow import DAG
from airflow.operators.bash_operator import BashOperator

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2023, 2, 17),
    'retries': 0
}

dag = DAG(
    'dag_for_loop_airflow_variable_optimized',
    default_args=default_args,
    catchup=False,
    schedule='@daily'
)

for i in range(1000):
    task = BashOperator(
        task_id=f'bash_use_variable_good_{i}',
        bash_command='echo variable foo=${foo_env}',
        dag=dag,
        env={'foo_env': '{{ var.value.get("example_var") }}'},
    )

בודקים את זמן הניתוח של ה-DAG החדש:

  1. מחכים עד שהרצת ה-DAG תסתיים.

  2. מריצים שוב את הפקודה dags report כדי לראות את זמן הניתוח של כל ה-DAG:

    file                | duration       | dag_num | task_num | dags
    ====================+================+=========+==========+=====================
    /dag_for_loop_airfl | 0:00:37.000369 | 1       | 1000     | dag_for_loop_airflow
    ow_variable.py      |                |         |          | _variable
    /dag_for_loop_airfl | 0:00:01.109457 | 1       | 1000     | dag_for_loop_airflow
    ow_variable_optimiz |                |         |          | _variable_optimized
    ed.py               |                |         |          |
    /airflow_monitoring | 0:00:00.040510 | 1       | 1        | airflow_monitoring
    .py                 |                |         |          |
    
  3. בודקים שוב את היומנים של dag-processor-manager ומנתחים את משך הניתוח.

    רשומה ביומן של DAG לדוגמה מראה שזמן הניתוח של ה-DAG הוא 4.21
    שניות
    איור 14. בקטעי היומן של מנהל התהליכים של DAG מוצגים זמני הניתוח של DAG אחרי שהקוד של DAG עבר אופטימיזציה (אפשר ללחוץ כדי להגדיל)

החלפת משתני הסביבה בתבניות Airflow פישטה את קוד ה-DAG והפחיתה את זמן האחזור של הניתוח פי עשרה בערך.

אופטימיזציה של הגדרות סביבת Airflow

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

במדריך הזה, ניתוח הקבצים של DAG לוקח הרבה זמן, ומחזורי הניתוח מתחילים לחפוף, מה שגורם למיצוי הקיבולת של מתזמן המשימות. בדוגמה שלנו, ניתוח ה-DAG הראשון נמשך יותר מ-5 שניות, ולכן צריך להגדיר את מתזמן המשימות כך שיפעל בתדירות נמוכה יותר כדי להשתמש במשאבים בצורה יעילה יותר. תשנו את הגדרת ברירת המחדל של Airflow‏ scheduler_heartbeat_sec. ההגדרה הזו קובעת את התדירות (בשניות) שבה מתזמן המשימות יפעל. כברירת מחדל, הערך שמוגדר הוא 5 שניות. אפשר לשנות את אפשרות ההגדרה הזו של Airflow על ידי החלפתה.

משנים את scheduler_heartbeat_sec אפשרות ההגדרה של Airflow:

  1. נכנסים לדף Environments במסוף Google Cloud .

    מעבר אל Environments

  2. ברשימת הסביבות, לוחצים על שם הסביבה. הדף Environment details ייפתח.

  3. עוברים לכרטיסייה Airflow Configuration Overrides (שינויים בהגדרות של Airflow).

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

  5. שינוי של אפשרות ההגדרה ב-Airflow:

    קטע מפתח ערך
    scheduler scheduler_heartbeat_sec 10
  6. לוחצים על שמירה ומחכים עד שההגדרה של הסביבה תתעדכן.

בודקים את מדדי מתזמן הפגישות:

  1. עוברים לכרטיסייה מעקב ובוחרים באפשרות מתזמנים.

  2. בתרשים Scheduler heartbeat, לוחצים על לחצן More options (סמל של שלוש נקודות) ואז על View in the Metrics Explorer.

בתרשים של פעילות ה-Scheduler מוצגות פעימות לב בתדירות נמוכה יותר
איור 15. גרף פעימות הלב של מתזמן הבקשות (לחיצה להגדלה)

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

הקצאת משאבים נוספים לכלי לתזמון

ב-Managed Airflow (דור 2), אפשר להקצות יותר משאבי CPU וזיכרון לתזמן. כך תוכלו לשפר את הביצועים של הכלי לתזמון ולהאיץ את זמן הניתוח של ה-DAG.

הקצאת מעבד (CPU) וזיכרון נוספים למתזמן:

  1. נכנסים לדף Environments במסוף Google Cloud .

    מעבר אל Environments

  2. ברשימת הסביבות, לוחצים על שם הסביבה. הדף Environment details ייפתח.

  3. עוברים לכרטיסייה Environment configuration (הגדרת הסביבה).

  4. מחפשים את ההגדרה Resources (משאבים) > Workloads (עומסי עבודה) ולוחצים על Edit (עריכה).

  5. בקטע Scheduler, בשדה Memory, מציינים את מגבלת הזיכרון החדשה. במדריך הזה, משתמשים ב-4GB.

  6. בשדה CPU, מציינים את מגבלת המעבד החדשה. במדריך הזה, משתמשים ב-2 vCPU.

  7. שומרים את השינויים וממתינים כמה דקות עד שהמתזמנים של Airflow יופעלו מחדש.

  8. עוברים לכרטיסייה Logs (יומנים), ואז ל-All logs (כל היומנים) > DAG processor manager (מנהל מעבדי DAG).

  9. בודקים את היומנים של dag-processor-manager ומשווים את משך הניתוח של ה-DAG לדוגמה:

    רשומה ביומן של ה-DAG לדוגמה מראה שזמן הניתוח של ה-DAG האופטימלי הוא 1.5 שניות. ב-DAG לא ממוטב, זמן הניתוח הוא 28.71 שניות
    איור 16. בקטעי היומן של מנהל העיבוד של DAG מוצגים זמני הניתוח של DAG אחרי הקצאה של עוד משאבים למתזמן (אפשר ללחוץ כדי להגדיל)

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

מומלץ להקצות משאבים רק אחרי שמטמיעים את האופטימיזציות האפשריות של קוד ה-DAG וההגדרה של Airflow.

הסרת המשאבים

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

מחיקת הפרויקט

  1. במסוף Google Cloud , נכנסים לדף Manage resources.

    כניסה לדף Manage resources

  2. ברשימת הפרויקטים, בוחרים את הפרויקט שרוצים למחוק ולוחצים על Delete.
  3. כדי למחוק את הפרויקט, כותבים את מזהה הפרויקט בתיבת הדו-שיח ולוחצים על Shut down.

מחיקת משאבים בודדים

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

מחיקת סביבת Managed Airflow במהלך התהליך הזה, נמחק גם הדלי של הסביבה.

המאמרים הבאים