פתרון בעיות בתזמון של Airflow

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

בדף הזה מפורטים שלבים לפתרון בעיות ומידע על בעיות נפוצות שקשורות לתזמון של Airflow ולמעבדי DAG.

זיהוי המקור של הבעיה

כדי להתחיל לפתור את הבעיה, צריך לזהות אם היא מתרחשת:

  • בזמן הניתוח של ה-DAG, בזמן שה-DAG מנותח על ידי מעבד DAG של Airflow
  • בזמן ההרצה, בזמן שה-DAG מעובד על ידי מתזמן Airflow

מידע נוסף על זמן הניתוח וזמן ההפעלה זמין במאמר ההבדל בין זמן הניתוח של DAG לבין זמן ההפעלה של DAG.

בדיקת בעיות בעיבוד DAG

  1. בדיקת היומנים של מעבד ה-DAG.
  2. בדיקת זמני הניתוח של DAG.

מעקב אחרי משימות שפועלות ומשימות בתור

כדי לבדוק אם יש משימות תקועות בתור, פועלים לפי השלבים הבאים.

  1. במסוף Google Cloud , עוברים לדף Environments.

    מעבר אל Environments

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

  3. עוברים לכרטיסייה מעקב.

  4. בכרטיסייה Monitoring, בודקים את התרשים Airflow tasks בקטע DAG runs ומזהים בעיות אפשריות. משימות Airflow הן משימות שנמצאות במצב של המתנה בתור ב-Airflow. הן יכולות לעבור לתור של ברוקר Celery או Kubernetes Executor. משימות בתור של Celery הן מופעים של משימות שמוכנסים לתור של ברוקר Celery.

פתרון בעיות בזמן הניתוח של DAG

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

ניתוח ותזמון של DAG ב-Managed Airflow (דור קודם 1) וב-Airflow 1

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

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

התפלגות של מספר המשימות והזמן שנדרש לביצוען

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

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

הגדרה של Airflow לצורך התאמה להיקף השימוש

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

  • מקבילות של עובדים

    הפרמטר [celery]worker_concurrency קובע את המספר המקסימלי של משימות ש-worker של Airflow יכול לבצע בו-זמנית. אם מכפילים את הערך של הפרמטר הזה במספר העובדים של Airflow בסביבת Managed Airflow, מקבלים את המספר המקסימלי של משימות שאפשר להריץ ברגע נתון בסביבה. המספר הזה מוגבל על ידי אפשרות ההגדרה [core]parallelism Airflow, שמתוארת בהמשך.

  • מספר מקסימלי של הפעלות DAG פעילות

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

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

    אפשר גם להגדיר את הערך הזה ברמת ה-DAG באמצעות הפרמטר max_active_runs.

  • מספר המשימות הפעילות המקסימלי לכל DAG

    אפשרות ההגדרה [core]max_active_tasks_per_dag Airflow קובעת את המספר המקסימלי של מופעי משימות שיכולים לפעול בו-זמנית בכל DAG.

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

    אפשר גם להגדיר את הערך הזה ברמת ה-DAG באמצעות הפרמטר max_active_tasks.

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

  • מקביליות וגודל המאגר

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

    זהו פרמטר גלובלי לכל הגדרת Airflow.

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

    אפשר להגדיר את גודל המאגר בממשק המשתמש של Airflow (Admin > Pools). משנים את גודל המאגר לרמת המקביליות שצפויה בסביבה שלכם.

    בדרך כלל, [core]parallelism מוגדר כמכפלה של מספר העובדים המקסימלי ושל [celery]worker_concurrency.

פתרון בעיות בהרצת משימות ובהוספת משימות לתור

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

הפעלות של DAG לא מתבצעות

תיאור הבעיה:

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

  • ההרצה של ה-DAG תמיד תהיה בעתיד, וה-DAG אף פעם לא יורץ.

  • הפעלות קודמות של DAG מסומנות כהפעלות שבוצעו בהצלחה למרות שהן לא בוצעו.

מידע נוסף זמין במסמכי התיעוד של Apache Airflow.

פתרונות אפשריים:

  • פועלים לפי ההמלצות בתיעוד של Apache Airflow.

  • הגדרת start_date סטטי ל-DAG. אפשר גם להשתמש ב-catchup=False כדי להשבית את ההרצה של ה-DAG לתאריכים קודמים.

  • מומלץ להימנע משימוש ב-datetime.now() או ב-days_ago(<number of days>), אלא אם אתם מודעים לתופעות הלוואי של הגישה הזו.

שימוש בתכונה TimeTable של Airflow scheduler

לוחות זמנים זמינים החל מ-Airflow 2.2.

אפשר להגדיר ל-DAG טבלת זמנים באחת מהשיטות הבאות:

אפשר גם להשתמש בלוחות זמנים מובנים.

משאבי אשכול מוגבלים

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

  • יוצרים סביבה חדשה עם סוג מכונה שמספק ביצועים טובים יותר ומעבירים אליה את ה-DAG.
  • יוצרים עוד סביבות Managed Airflow ומפצלים את ה-DAG ביניהן.
  • משנים את סוג המכונה לצמתי GKE, כמו שמתואר במאמר שדרוג סוג המכונה לצמתי GKE. ההליך הזה עלול לגרום לשגיאות, ולכן הוא האפשרות הכי פחות מומלצת.
  • משדרגים את סוג המכונה של מופע Cloud SQL שמריץ את מסד הנתונים של Airflow בסביבה שלכם, למשל באמצעות הפקודות gcloud composer environments update. יכול להיות שהביצועים הנמוכים של מסד הנתונים של Airflow הם הסיבה לכך שהמתזמן איטי.

איך להימנע מתזמון משימות במהלך חלונות תחזוקה

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

שימוש ב-wait_for_downstream ב-DAGs

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

משימות שנמצאות בתור זמן רב מדי יבוטלו וישובצו מחדש

אם משימת Airflow נשארת בתור יותר מדי זמן, המתזמן יתזמן אותה מחדש להפעלה אחרי שיעבור פרק הזמן שמוגדר באפשרות ההגדרה של [scheduler]task_queued_timeoutAirflow. ערך ברירת המחדל הוא 2400. בגרסאות Airflow שקודמות לגרסה 2.3.1, המשימה מסומנת גם כנכשלה ומבוצע ניסיון חוזר אם היא עומדת בדרישות לניסיון חוזר.

דרך אחת לבחון את הסימפטומים של המצב הזה היא להסתכל על התרשים עם מספר המשימות בתור ("כרטיסיית מעקב" בממשק המשתמש של Managed Airflow). אם העליות בתרשים הזה לא יורדות תוך שעתיים בערך, סביר להניח שהמשימות יתוזמנו מחדש (ללא יומנים), ואחריהן יופיעו רשומות ביומן של המתזמן עם הכיתוב "Adopted tasks were still pending ...". במקרים כאלה, יכול להיות שתופיע ההודעה 'לא נמצא קובץ יומן…' ביומני המשימות של Airflow כי המשימה לא בוצעה.

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

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

גישה לפרמטר min_file_process_interval ב-Managed Airflow

ב-Managed Airflow, משתנה האופן שבו מתזמן Airflow משתמש ב-[scheduler]min_file_process_interval.

Airflow 1

במקרה של Managed Airflow באמצעות Airflow 1, המשתמשים יכולים להגדיר את הערך של [scheduler]min_file_process_interval בין 0 ל-600 שניות. ערכים שגבוהים מ-600 שניות יניבו את אותן תוצאות כמו אם [scheduler]min_file_process_interval מוגדר ל-600 שניות.

Airflow 2

בגרסאות של Managed Airflow שקודמות לגרסה 1.19.9, [scheduler]min_file_process_interval המערכת מתעלמת מ-

גרסאות Managed Airflow חדשות יותר מ-1.19.9:

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

אפשר להשתמש ב-[scheduler]min_file_process_interval כדי להגדיר את התדירות שבה מתבצע ניתוח של DAG, אבל הפרמטר הזה לא יכול להיות ארוך יותר מהזמן שנדרש למתזמן כדי לבצע לולאות [scheduler]num_runs כשמתזמנים את ה-DAG.

סימון משימות כנכשלות אחרי שמגיעים ל-dagrun_timeout

אם הפעלת DAG לא מסתיימת תוך dagrun_timeout (פרמטר של DAG), מתזמן המשימות מסמן את המשימות שלא הסתיימו (פועלות, מתוזמנות וממתינות בתור) ככאלה שנכשלו.

פתרון:

  • הארכת dagrun_timeout כדי לעמוד בדרישות הזמן הקצוב לתפוגה.

תסמינים של עומס כבד על מסד הנתונים של Airflow

לפעמים ביומני התזמון של Airflow מופיעה רשומת אזהרה כזו:

Scheduler heartbeat got an exception: (_mysql_exceptions.OperationalError) (2006, "Lost connection to MySQL server at 'reading initial communication packet', system error: 0")"

יכול להיות שסימפטומים דומים יופיעו גם ביומני העובדים של Airflow:

ל-MySQL:

(_mysql_exceptions.OperationalError) (2006, "Lost connection to MySQL server at
'reading initial communication packet', system error: 0")"

ב-PostgreSQL:

psycopg2.OperationalError: connection to server at ... failed

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

פתרונות אפשריים:

שרת האינטרנט מציג את האזהרה 'נראה שהמתזמן לא פועל'

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

לפעמים, אם יש עומס כבד על מתזמן המשימות, יכול להיות שהוא לא יוכל לדווח על הפעימה שלו כל [scheduler]scheduler_heartbeat_sec.

במצב כזה, יכול להיות שיוצג בשרת האינטרנט של Airflow האזהרה הבאה:

The scheduler does not appear to be running. Last heartbeat was received <X>
seconds ago.

פתרונות אפשריים:

  • להגדיל את משאבי המעבד (CPU) והזיכרון של הכלי לתזמון.

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

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

  • מגדילים את הערך של אפשרות ההגדרה [scheduler]scheduler_health_check_threshold ב-Airflow, כדי ששרת האינטרנט ימתין זמן רב יותר לפני שידווח על חוסר הזמינות של המתזמן.

פתרונות עקיפים לבעיות שנתקלים בהן במהלך מילוי חוסרים ב-DAG

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

Airflow 2

gcloud composer environments run \
  ENVIRONMENT_NAME \
  --location LOCATION \
   dags backfill -- -B \
   -s START_DATE \
   -e END_DATE \
   DAG_NAME

כדי להריץ מחדש רק משימות שנכשלו ב-DAG ספציפי, משתמשים גם בארגומנט --rerun-failed-tasks.

Airflow 1

gcloud composer environments run \
  ENVIRONMENT_NAME \
  --location LOCATION \
  backfill -- -B \
  -s START_DATE \
  -e END_DATE \
  DAG_NAME

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

מחליפים את:

  • ENVIRONMENT_NAME בשם הסביבה.
  • LOCATION עם האזור שבו הסביבה ממוקמת.
  • START_DATE עם ערך לפרמטר start_date של DAG, בפורמט YYYY-MM-DD.
  • END_DATE עם ערך לפרמטר end_date של DAG, בפורמט YYYY-MM-DD.
  • DAG_NAME בשם של ה-DAG.

לפעמים פעולת מילוי החוסר (backfill) יוצרת מצב של קיפאון (deadlock), שבו אי אפשר לבצע מילוי חוסר כי יש נעילה על משימה. לדוגמה:

2022-11-08 21:24:18.198 CET DAG ID Task ID Run ID Try number
2022-11-08 21:24:18.201 CET -------- --------- -------- ------------
2022-11-08 21:24:18.202 CET 2022-11-08 21:24:18.203 CET These tasks are deadlocked:
2022-11-08 21:24:18.203 CET DAG ID Task ID Run ID Try number
2022-11-08 21:24:18.204 CET ----------------------- ----------- ----------------------------------- ------------
2022-11-08 21:24:18.204 CET <DAG name> <Task name> backfill__2022-10-27T00:00:00+00:00 1
2022-11-08 21:24:19.249 CET Command exited with return code 1
...
2022-11-08 21:24:19.348 CET Failed to execute job 627927 for task backfill

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

  • כדי להשבית את לוח הזמנים המינימלי, צריך לשנות את הערך של [core]schedule_after_task_execution ל-False.

  • מריצים מילוי חוסרים לטווח תאריכים מצומצם יותר. לדוגמה, אפשר להגדיר את START_DATE ואת END_DATE כדי לציין תקופה של יום אחד בלבד.

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