Managed Airflow (דור 3) | Managed Airflow (דור 2) | Managed Airflow (דור 1 מדור קודם)
בדף הזה מפורטים שלבים לפתרון בעיות ומידע על בעיות נפוצות בתהליכי עבודה.
הרבה בעיות בהרצת DAG נובעות מביצועים לא אופטימליים של הסביבה. כדי לבצע אופטימיזציה של הסביבה, אפשר להיעזר במדריך אופטימיזציה של הביצועים והעלויות של הסביבה.
חלק מהבעיות בהפעלות של DAG יכולות להיגרם בגלל שתזמן Airflow לא פועל בצורה תקינה או אופטימלית. כדי לפתור את הבעיות האלה, צריך לפעול לפי ההוראות לפתרון בעיות בתכנון פגישות.
פתרון בעיות בתהליך העבודה
כדי להתחיל לפתור את הבעיה:
בודקים את היומנים של Airflow.
כדי להגדיל את רמת הרישום ביומן של Airflow, צריך לבטל את ההגדרה של אפשרות התצורה הבאה של Airflow.
קטע מפתח ערך logginglogging_levelערך ברירת המחדל הוא INFO. מגדירים את הערך ל-DEBUGכדי לקבל יותר פירוט בהודעות היומן.בודקים את לוח הבקרה של Monitoring.
בודקים את Cloud Monitoring.
במסוף Google Cloud , בודקים אם יש שגיאות בדפים של רכיבי הסביבה.
בממשק האינטרנט של Airflow, בודקים בתצוגת הגרף של ה-DAG אם יש מופעים של משימות שנכשלו.
קטע מפתח ערך webserverdag_orientationLR,TB,RLאוBT
חקירת כשלים במשימות Airflow באמצעות Gemini Cloud Assist
Gemini Cloud Assist investigations הוא כלי לניתוח שורש הבעיה (RCA) לפתרון בעיות בתשתית ובאפליקציות בסביבות ענן מורכבות ומבוזרות. החקירות יכולות לעזור לכם להבין, לאבחן ולפתור בעיות ב- Google Cloud. בעזרת חקירות, אתם יכולים לייעל את התגובה לאירועים על ידי קיצור הזמן לפתרון ושיפור הזמינות הכוללת, והכול במאמץ מופחת.
ב-Managed Airflow, אתם יכולים להתחיל חקירות של Gemini Cloud Assist למשימות Airflow שנכשלו מממשק המשתמש של DAG. Managed Airflow מאכלס באופן אוטומטי פרטים כמו תיאור הבעיה וטווח הזמן, וכולל את הסביבה שלכם כמשאב רלוונטי.
התחלת חקירות והצגתן
כדי להתחיל חקירה חדשה ב-Gemini Cloud Assist למשימת Airflow שנכשלה, או כדי להציג חקירה קיימת:
נכנסים לדף Environments במסוף Google Cloud .
בוחרים סביבה כדי לראות את הפרטים שלה.
בדף פרטי הסביבה, עוברים לכרטיסייה DAGs.
לוחצים על השם של DAG.
בדף פרטי DAG, עוברים לכרטיסייה היסטוריית הפעלות ולוחצים על הפעלת DAG עם משימות שנכשלו.
בעמודה State (מצב) של משימת Airflow שנכשלה, לוחצים על Investigate (בדיקה):
- כדי לבדוק משימה חדשה, לוחצים על בדיקה.
- אם כבר מתבצעת חקירה לגבי משימה מסוימת, לוחצים על הצגת החקירה כדי לבדוק את החקירה הקיימת. לחלופין, אפשר להתחיל בדיקה נוספת בלחיצה על בדיקה חדשה.
אפשר להמשיך אל יצירה, הפעלה ובדיקה של החקירה באמצעות Gemini Cloud Assist.
דוגמה לחקירה
בדוגמה הזו מוצג תהליך של בדיקת משימה שנכשלה.
במרכז הבקרה DAG Statistics (נתונים סטטיסטיים של DAG) בדף Monitoring (מעקב) >, בודקים את הנתונים של הפעלות DAG שנכשלו:
איור 1. גרף של הפעלות DAG שהושלמו (לחיצה להגדלה) עוברים אל DAGs. בעמודה Failed runs (1h) (הפעלות שנכשלו (שעה)) אפשר לראות של-DAG
create_large_txt_file_print_logsהיו כמה הפעלות שנכשלו בשעה האחרונה. לוחצים על השם של ה-DAG.
איור 2. רשימה של תרשימי DAG עם נתונים סטטיסטיים של הפעלות DAG (אפשר ללחוץ כדי להגדיל) לוחצים על אחת מהרצות ה-DAG שנכשלו ואז לוחצים על Investigate (חקירה) לצד הרשומה של משימת Airflow שנכשלה ומתחילים את החקירה.
איור 3. רשימת המשימות בתרשימי ה-DAG שנכשלו (אפשר ללחוץ כדי להגדיל) צריך להמתין עד שהבדיקה תסתיים.
ברשימה relevant observations (תצפיות רלוונטיות) מפורט תהליך החקירה. בדוגמה הזו, המשימה נכשלה בלי ליצור יומנים, אבל Gemini Cloud Assist הצליח למצוא את הסיבה לכישלון ביומני מתזמן המשימות של Airflow, שבהם המשימה הופסקה כזומבי.
איור 4. תצפית על זיהוי משימת זומבי (לחצו כדי להגדיל) לבסוף, Gemini Cloud Assist מסכם את הממצאים ומספק השערה יחד עם המלצות לפתרון הבעיה. בדוגמה הזו, המשימה נכשלה כי לעובד של Airflow לא היו מספיק משאבים כדי לעבד אותה. התמיכה בכך מבוססת על תצפיות של הפעלה מחדש של workerPod כמה פעמים עם שגיאת OOM, ועל סיום המשימה על ידי המתזמן כמשימת זומבי.
איור 5. השערה לגבי מיצוי משאבים של Airflow Worker Pod (לחצו כדי להגדיל)
ניפוי באגים בכשלים של אופרטורים
כדי לנפות באגים בכשל של אופרטור:
- בודקים אם יש שגיאות שקשורות למשימה.
- בודקים את היומנים של Airflow.
- בודקים את Cloud Monitoring.
- בודקים יומנים ספציפיים של אופרטורים.
- מתקנים את השגיאות.
- מעלים את ה-DAG לתיקייה
/dags. - בממשק האינטרנט של Airflow, מנקים את המצבים הקודמים של ה-DAG.
- ממשיכים או מפעילים את ה-DAG.
פתרון בעיות בביצוע משימות
Airflow היא מערכת מבוזרת עם הרבה ישויות כמו מתזמן, מפעיל, עובדים שתקשרים ביניהם דרך תור משימות ומסד הנתונים של Airflow, ושולחים אותות (כמו SIGTERM). בתרשים הבא מוצג סקירה כללית של הקישורים בין רכיבי Airflow.
במערכת מבוזרת כמו Airflow, יכולות להיות בעיות בקישוריות לרשת, או בעיות לסירוגין בתשתית הבסיסית. זה יכול להוביל למצבים שבהם משימות נכשלות ומתוזמנות מחדש לביצוע, או שמשימות לא מושלמות בהצלחה (לדוגמה, משימות זומבי או משימות שנתקעו במהלך הביצוע). ל-Airflow יש מנגנונים להתמודדות עם מצבים כאלה, והוא מחדש את הפעולה הרגילה באופן אוטומטי. בקטעים הבאים מוסברות בעיות נפוצות שמתרחשות במהלך ביצוע משימות על ידי Airflow.
פתרון בעיות במשימות של KubernetesExecutor
CeleryKubernetesExecutor הוא סוג של מפעיל ב-Managed Airflow (דור 3) שיכול להשתמש ב-CeleryExecutor וב-KubernetesExecutor בו-זמנית.
מידע נוסף על פתרון בעיות במשימות שמופעלות באמצעות KubernetesExecutor זמין בדף Use CeleryKubernetesExecutor.
משימות נכשלות בלי ליצור יומנים
יכולות להיות כמה סיבות לכך שמופע של משימה ייכשל בלי שיופקו יומנים. לדוגמה, זה יכול לקרות בגלל שגיאות בניתוח של DAG, עיכובים בסנכרון של DAG או אם פוד של Airflow worker מוצא להסרה במהלך ביצוע משימה (ראו משימה נכשלת בגלל הסרת פוד).
שגיאות או פסק זמן בניתוח DAG
אם יש שגיאות תכנות בקובץ DAG או אם ניתוח קובץ DAG נמשך יותר מדי זמן, יכול להיות שתזמן המשימות של Airflow יוכל לתזמן משימות, אבל עובדי Airflow לא יוכלו להריץ אותן. במקרה כזה, אפשר לסמן משימה כFailed בלי שיהיה יומן רישום של ההפעלה שלה.
תסמינים
יומני העובדים של Airflow ב-Cloud Logging מכילים הודעות כמו:
airflow.exceptions.AirflowException: Dag "example-dag" could not be found; either it does not exist or it failed to parse.-
airflow.exceptions.AirflowTaskTimeout: Timeout, PID: 12345(יכול להיות שההודעה הזו לא תכיל את שם ה-DAG או את נתיב הקובץ). ERROR - Failed to import: /home/airflow/gcs/dags/example-dag.py(ללא tracebacks מפורטים).
אם יש שגיאות בייבוא DAG, יכול להיות שהן יופיעו בממשק המשתמש של Airflow, או בהודעות
Broken DAGבדף פרטי הסביבה במסוףGoogle Cloud .
פתרון
בודקים ביומני העובדים של Airflow אם יש שגיאות שקשורות לניתוח של DAG. אם מוצגות שגיאות מסוג
AirflowTaskTimeout, יכול להיות שחלף הזמן הקצוב לניתוח של DAG. הזמן הקצוב לתפוגה של ניתוח (parsing) בעובדי Airflow נשלט על ידיdagbag_import_timeout.אם זמני הניתוח של DAG ארוכים, כדאי לבדוק אם יש תחרות על המעבד באשכול של הסביבה. אם אין מספיק CPU לעובדים: הגדלת ה-CPU של עובדי Airflow. אפשרות אחרת היא להקטין את worker_concurrency, כמו שמתואר במאמר בנושא אופטימיזציה של הסביבה.
אם השימוש במעבד נמוך, כדאי לבצע אופטימיזציה של הגדרת ה-DAG כדי לקצר את זמן הניתוח, למשל על ידי הימנעות מקוד ברמה העליונה.
אם מופיעות שגיאות
Failed to importאו שגיאות בממשק המשתמש של Airflow או במסוףGoogle Cloud , כדאי לבדוק את היומנים של מעבד ה-DAG כדי לראות את פרטי ה-traceback, כמו שמתואר במאמר פתרון בעיות במעבד ה-DAG. אפשר גם להריץ את הפקודה הבאה ב-CLI של gcloud כדי לראות את השגיאות בייבוא ה-DAG:gcloud composer environments run ENVIRONMENT_NAME \ --location LOCATION \ dags list-import-errorsאם ה-DAG שלכם מבצע קריאות לשירותים חיצוניים במהלך הניתוח, כדאי להוסיף בלוקים של
try...exceptסביב הקריאות האלה כדי לטפל בשגיאות זמניות.אם אי אפשר לבצע אופטימיזציה של ניתוח DAG, צריך להגדיל את
dagbag_import_timeout. מבטלים את הגדרת ברירת המחדל של Airflow, שהיא 30 שניות, ומגדירים ערך גבוה יותר, למשל 120 שניות.
עיכובים בסנכרון קובצי DAG
כשמעלים או מעדכנים קובצי DAG בדלי של הסביבה, לוקח זמן עד שהקבצים האלה מסונכרנים עם העובדים והמתזמנים של Airflow.
הסנכרון הזה מתבצע באופן עצמאי בכל המתזמנים והעובדים. אם מפעילים ריצת DAG זמן קצר אחרי העלאה או עדכון של קובץ DAG, והקובץ עדיין לא מסונכרן עם worker שמבצע משימה, המשימה תיכשל בלי יומנים, ויכול להיות שתראו את השגיאה airflow.exceptions.AirflowException: Dag "example-dag" could not be
found... ביומני ה-worker.
הסנכרון הזה בדרך כלל נמשך דקה או שתיים, אבל הוא יכול להימשך יותר זמן אם יש לכם הרבה קבצים או קבצים גדולים בתיקיות dags/ או plugins/ בדלי.
פתרון
צריך להמתין לפחות 2 דקות אחרי העלאה או עדכון של DAG או פלאגינים לפני שמפעילים DAG או מאפשרים אותם.
משימות תקועות במצב 'בתור'
בגרסאות Airflow קודמות לגרסה 2.6.3, לפעמים משימות נתקעות באופן קבוע במצב queued. מצב כזה יכול לקרות אם משימה מסומנת כממתינה בתור במסד הנתונים של Airflow, אבל היא לא קיימת בפועל ב-Celery.
במצב כזה, יכול להיות שעובדי Airflow ייכשלו בבדיקות הפעילות ויופעלו מחדש, מה שעלול לגרום לכך שמשימות ייכשלו עם השגיאות 'לא נמצא קובץ יומן'.
הבעיה הזו נפתרה בגרסה Airflow 2.6.3 ואילך. אם אתם משתמשים בגרסה קודמת של Airflow, אתם יכולים לשדרג את הסביבה לגרסת תמונה שמשתמשת ב-Airflow 2.6.3 ואילך.
כפתרון עקיף, אפשר לנקות ידנית משימות שנתקעו במצב 'בהמתנה'. בממשק המשתמש של Airflow, עוברים אל Browse > Task Instances, מוצאים מופעים של משימות שנתקעו במצב queued ומגדירים את המצב שלהם ל-failed.
המשימות מופסקות באופן פתאומי
במהלך ביצוע המשימה, יכול להיות שהתהליכים של Airflow יסתיימו באופן פתאומי בגלל בעיות שלא קשורות ספציפית למשימה עצמה. במאמר סיבות נפוצות לבעיות מפורטת רשימה של תרחישים כאלה ופתרונות אפשריים. בקטעים הבאים מפורטים כמה תסמינים נוספים שיכולים לנבוע מסיבות השורש האלה:
משימות זומבי
מערכת Airflow מזהה שני סוגים של חוסר התאמה בין משימה לבין תהליך שמבצע את המשימה:
משימות זומבי הן משימות שאמורות לפעול אבל לא פועלות. זה יכול לקרות אם התהליך של המשימה הסתיים או לא מגיב, אם העובד של Airflow לא דיווח על סטטוס המשימה בזמן כי הוא עמוס מדי, או אם מכונת ה-VM שבה המשימה מבוצעת כובתה. מערכת Airflow מאתרת משימות כאלה מעת לעת, ומבצעת את המשימה מחדש או נכשלת בה, בהתאם להגדרות המשימה.
איך לגלות משימות זומבי
resource.type="cloud_composer_environment" resource.labels.environment_name="ENVIRONMENT_NAME" log_id("airflow-scheduler") textPayload:"Detected zombie job"משימות לא פעילות הן משימות שלא אמורות לפעול. Airflow מוצא משימות כאלה מעת לעת ומסיים אותן.
מידע נוסף על פתרון בעיות שקשורות למשימות זומבי זמין במאמר בנושא סיבות נפוצות לבעיות.
אותות SIGTERM
אותות SIGTERM משמשים את Linux, Kubernetes, Airflow scheduler ו-Celery כדי להפסיק תהליכים שאחראים להפעלת עובדי Airflow או משימות Airflow.
יכולות להיות כמה סיבות לשליחת אותות SIGTERM בסביבה:
משימה הפכה למשימת זומבי וצריך להפסיק אותה.
מתזמן המשימות זיהה כפילות של משימה ושולח למשימה אותות של סיום המופע ו-SIGTERM כדי להפסיק אותה.
בהתאמה אופקית של קבוצות Pod לעומס, מישור הבקרה של GKE שולח אותות SIGTERM כדי להסיר Pods שכבר לא נחוצים.
מתזמן המשימות יכול לשלוח אותות SIGTERM לתהליך DagFileProcessorManager. האותות האלה של SIGTERM משמשים את Scheduler לניהול מחזור החיים של התהליך DagFileProcessorManager, ואפשר להתעלם מהם בבטחה.
דוגמה:
Launched DagFileProcessorManager with pid: 353002 Sending Signals.SIGTERM to group 353002. PIDs of all processes in the group: [] Sending the signal Signals.SIGTERM to group 353002 Sending the signal Signals.SIGTERM to process 353002 as process group is missing.מרוץ תהליכים בין קריאה חוזרת (callback) של פעימת הלב לבין קריאות חוזרות של יציאה ב-local_task_job, שעוקב אחרי ההרצה של המשימה. אם בדיקת הפעימות מזהה שמשימה סומנה כהצלחה, היא לא יכולה להבחין בין מצב שבו המשימה עצמה הצליחה לבין מצב שבו נאמר ל-Airflow להתייחס למשימה כהצלחה. עם זאת, הוא יסיים את הפעולה של מפעיל המשימות בלי לחכות לסיום שלו.
אפשר להתעלם בבטחה מאותות SIGTERM כאלה. המשימה כבר במצב מוצלח, והביצוע של הפעלת ה-DAG כולה לא יושפע.
רשומת היומן
Received SIGTERM.היא ההבדל היחיד בין יציאה רגילה לבין סיום המשימה במצב של הצלחה.איור 7. מרוץ תהליכים בין אותות פעימת הלב לבין קריאות חוזרות (callback) ליציאה (אפשר ללחוץ כדי להגדיל) רכיב Airflow משתמש ביותר משאבים (CPU, זיכרון) מהמותר בצומת האשכול.
שירות GKE מבצע פעולות תחזוקה ושולח אותות SIGTERM לקבוצות Pod שפועלות בצומת שעומד לעבור שדרוג.
כשמופסקת פעילות של מופע של משימה באמצעות SIGTERM, אפשר לראות את רשומות היומן הבאות ביומנים של Airflow worker שהריץ את המשימה:
{local_task_job.py:211} WARNING - State of this instance has been externally set to queued. Terminating instance. {taskinstance.py:1411} ERROR - Received SIGTERM. Terminating subprocesses. {taskinstance.py:1703} ERROR - Task failed with exception
פתרונות אפשריים:
הבעיה הזו מתרחשת כשמכונה וירטואלית שמריצה את המשימה לא מקבלת מספיק זיכרון. ההגדרה הזו לא קשורה להגדרות של Airflow, אלא לכמות הזיכרון שזמינה למכונה הווירטואלית.
ב-Managed Airflow (דור 3), אפשר להקצות יותר משאבי CPU וזיכרון לעובדי Airflow.
אפשר להקטין את הערך של אפשרות ההגדרה
[celery]worker_concurrencyconcurrency ב-Airflow. האפשרות הזו קובעת כמה משימות יבוצעו בו-זמנית על ידי worker נתון של Airflow.
מידע נוסף על אופטימיזציה של הסביבה זמין במאמר אופטימיזציה של הביצועים והעלויות של הסביבה.
המשימה ב-Airflow הופסקה על ידי Negsignal.SIGKILL
לפעמים המשימה שלכם עשויה להשתמש ביותר זיכרון ממה שהוקצה לעובד של Airflow.
במקרה כזה, יכול להיות שהטעינה תופסק על ידי Negsignal.SIGKILL. המערכת שולחת את האות הזה כדי למנוע צריכת זיכרון נוספת, שעשויה להשפיע על ההפעלה של משימות אחרות ב-Airflow. ביומן של Airflow worker יכול להיות שתראו את רשומת היומן הבאה:
{local_task_job.py:102} INFO - Task exited with return code Negsignal.SIGKILL
הפונקציה Negsignal.SIGKILL יכולה להופיע גם כקוד -9.
פתרונות אפשריים:
הפחתת
worker_concurrencyשל עובדי Airflow.הגדלת כמות הזיכרון שזמינה לעובדי Airflow.
כדי לנהל משימות שדורשות הרבה משאבים ב-Managed Airflow, אפשר להשתמש ב-KubernetesPodOperator או ב-GKEStartPodOperator כדי לבודד את המשימות ולהקצות משאבים בהתאמה אישית.
אופטימיזציה של המשימות כדי להשתמש בפחות זיכרון.
המשימה נכשלת בגלל עומס על המשאבים
תסמין: במהלך ביצוע משימה, תהליך המשנה של Airflow Worker שאחראי לביצוע משימת Airflow מופסק באופן פתאומי. השגיאה שמופיעה ביומן של Airflow worker עשויה להיראות כמו השגיאה שבהמשך:
...
File "/opt/python3.8/lib/python3.8/site-packages/celery/app/trace.py", line 412, in trace_task R = retval = fun(*args, **kwargs) File "/opt/python3.8/lib/python3.8/site-packages/celery/app/trace.py", line 704, in __protected_call__ return self.run(*args, **kwargs) File "/opt/python3.8/lib/python3.8/site-packages/airflow/executors/celery_executor.py", line 88, in execute_command _execute_in_fork(command_to_exec) File "/opt/python3.8/lib/python3.8/site-packages/airflow/executors/celery_executor.py", line 99, in _execute_in_fork
raise AirflowException('Celery command failed on host: ' + get_hostname())airflow.exceptions.AirflowException: Celery command failed on host: airflow-worker-9qg9x
...
פתרון:
ב-Managed Airflow (דור 3), מגדילים את מגבלות הזיכרון של עובדי Airflow.
אם הסביבה שלכם יוצרת גם משימות זומבי, כדאי לעיין במאמר בנושא פתרון בעיות שקשורות למשימות זומבי.
מדריך לניפוי באגים בבעיות שקשורות לזיכרון או לאחסון מלאים זמין במאמר ניפוי באגים בבעיות שקשורות לזיכרון או לאחסון מלאים ב-DAG.
המשימה נכשלת בגלל הוצאה של Pod
ה-Pods של Google Kubernetes Engine כפופים למחזור החיים של Kubernetes Pod ולפינוי Pod. הסיבה הנפוצה ביותר להוצאת Pod משימוש ב-Managed Airflow היא עליות פתאומיות בעומס המשימות.
פינוי של Pod יכול להתרחש כש-Pod מסוים משתמש יתר על המידה במשאבים של צומת, ביחס לציפיות לגבי צריכת המשאבים שהוגדרו לצומת. לדוגמה, יכול להיות שפינוי יקרה כשכמה משימות שצורכות הרבה זיכרון פועלות ב-Pod, והעומס המשולב שלהן גורם לחריגה ממגבלת צריכת הזיכרון של הצומת שבו ה-Pod הזה פועל.
אם מתבצעת הוצאה של Pod של Airflow worker, כל מופעי המשימות שפועלים ב-Pod הזה מופסקים, ומאוחר יותר מסומנים ככאלה שנכשלו על ידי Airflow.
היומנים נשמרים בזיכרון הזמני. אם פוד של Worker מוצא מהמערכת לפני שמאגר הנתונים הזמני מתרוקן, היומנים לא מופקים. אם המשימה נכשלת בלי יומנים, זה סימן שתהליכי העבודה של Airflow מופעלים מחדש בגלל חוסר זיכרון (OOM). יכול להיות שחלק מהיומנים יופיעו ב-Cloud Logging גם אם יומני Airflow לא הופקו.
כדי לראות את היומנים:
במסוף Google Cloud , עוברים לדף Environments.
ברשימת הסביבות, לוחצים על שם הסביבה. הדף Environment details ייפתח.
עוברים לכרטיסייה יומנים.
כדי לראות את היומנים של כל אחד מהעובדים ב-Airflow, עוברים אל All logs > Airflow logs > Workers.
פתרון:
הגדלת מגבלות הזיכרון לעובדי Airflow.
מוודאים שהמשימות ב-DAG הן אידמפוטנטיות וניתנות לניסיון חוזר.
מומלץ להימנע מהורדה של קבצים מיותרים למערכת הקבצים המקומית של עובדי Airflow.
ל-Airflow workers יש קיבולת מוגבלת במערכת הקבצים המקומית. נפח האחסון של Airflow worker יכול להיות מ-1GB עד 10GB. כשנגמר נפח האחסון, מישור הבקרה של GKE מפנה את ה-Pod של Airflow worker. כל המשימות שהעובד שהוצא מהמערכת ביצע ייכשלו.
דוגמאות לפעולות בעייתיות:
- הורדת קבצים או אובייקטים ואחסונם באופן מקומי בתהליך עבודה (worker) של Airflow. במקום זאת, כדאי לאחסן את האובייקטים האלה ישירות בשירות מתאים, כמו קטגוריה של Cloud Storage.
- גישה לאובייקטים גדולים בתיקייה
/dataמעובד Airflow. תהליך העבודה של Airflow מוריד את האובייקט למערכת הקבצים המקומית שלו. במקום זאת, כדאי להטמיע את ה-DAG כך שקבצים גדולים יעובדו מחוץ ל-Pod של Airflow worker.
סיבות נפוצות לבעיה
Airflow worker ran out of memory
כל תהליך עבודה של Airflow יכול להריץ עד [celery]worker_concurrency מופעים של משימות בו-זמנית. אם צריכת הזיכרון המצטברת של מופעי המשימות האלה חורגת ממגבלת הזיכרון של עובד Airflow, תהליך אקראי יופסק כדי לפנות משאבים.
לפעמים, מחסור בזיכרון בתהליך עובד של Airflow עלול לגרום לשליחה של מנות פגומות במהלך סשן של SQL Alchemy למסד הנתונים, לשרת DNS או לכל שירות אחר שמופעל על ידי DAG. במקרה כזה, יכול להיות שהצד השני של החיבור ידחה או ינתק חיבורים מה-worker של Airflow. לדוגמה:
"UNKNOWN:Error received from peer
{created_time:"2024-11-31T10:09:52.217738071+00:00", grpc_status:14,
grpc_message:"failed to connect to all addresses; last error: UNKNOWN:
ipv4:<ip address>:443: handshaker shutdown"}"
פתרונות:
מבצעים אופטימיזציה של משימות כדי להשתמש בפחות זיכרון, למשל על ידי הימנעות מקוד ברמה העליונה.
הפחתה
[celery]worker_concurrency.Increase memory for Airflow workers כדי להתאים לשינויים ב-
[celery]worker_concurrency.
העובד של Airflow הוצא
הוצאת pods היא חלק רגיל מהפעלת עומסי עבודה ב-Kubernetes. מערכת GKE מפנה תאי Pod אם נגמר להם האחסון או כדי לפנות משאבים לעומסי עבודה עם עדיפות גבוהה יותר.
פתרונות:
- אם הפינוי נגרם בגלל חוסר מקום באחסון, אפשר לצמצם את השימוש באחסון או להסיר קבצים זמניים ברגע שהם לא נחוצים יותר.
לחלופין, אפשר להגדיל את נפח האחסון הזמין או להריץ עומסי עבודה בתא ייעודי עם
KubernetesPodOperator.
תהליך העבודה של Airflow הופסק
יכול להיות שעובדי Airflow יוסרו באופן חיצוני. אם משימות שפועלות כרגע לא מסתיימות במהלך תקופת סיום מבוקר, הן מופסקות ויכול להיות שהן יזוהו כזומבים.
תרחישים ופתרונות אפשריים:
הפעלת העובדים של Airflow מחדש במהלך שינויים בסביבה, כמו שדרוגים או התקנת חבילות:
גילוי שינויים בסביבת Composer
resource.type="cloud_composer_environment" resource.labels.environment_name="ENVIRONMENT_NAME" log_id("cloudaudit.googleapis.com%2Factivity")אפשר לבצע פעולות כאלה כשלא מופעלות משימות קריטיות, או להפעיל ניסיונות חוזרים של משימות.
יכול להיות שרכיבים שונים לא יהיו זמינים באופן זמני במהלך פעולות תחזוקה.
אתם יכולים לציין חלונות זמן לתחזוקה כדי לצמצםחפיפה עם הביצוע של המשימות הקריטיות.
תהליך העבודה של זרימת האוויר היה תחת עומס כבד
כמות המשאבים של המעבד (CPU) והזיכרון שזמינים לעובד של Airflow מוגבלת על ידי הגדרות הסביבה. אם ניצול המשאבים מתקרב למגבלות, יכול להיות שזה יגרום לתחרות על משאבים ולעיכובים מיותרים במהלך ביצוע המשימה. במקרים קיצוניים, כשאין מספיק משאבים לאורך זמן, יכול להיות שייווצרו משימות זומבי.
פתרונות:
- עוקבים אחרי השימוש במעבד ובשימוש בזיכרון של ה-worker-ים ומשנים אותו כדי שלא יעלה על 80%.
מסד הנתונים של Airflow היה תחת עומס כבד
מסד נתונים משמש רכיבים שונים של Airflow כדי לתקשר ביניהם, ובמיוחד כדי לאחסן את אותות החיים של מופעי משימות. מחסור במשאבים במסד הנתונים גורם לזמני שאילתה ארוכים יותר, ויכול להשפיע על ביצוע המשימות.
לפעמים השגיאות הבאות מופיעות ביומנים של Airflow worker:
(psycopg2.OperationalError) connection to server at <IP address>,
port 3306 failed: server closed the connection unexpectedly
This probably means the server terminated abnormally before or while
processing the request.
פתרונות:
- מומלץ להימנע משימוש בהרבה הוראות
Variables.getבקוד DAG ברמה העליונה. במקום זאת, משתמשים בתבניות Jinja כדי לאחזר ערכים של משתני Airflow. לא מומלץ להשתמש ב-CloudLoggingHandler לטיפול ברישום ביומן. ה-handler הזה של הרישום גורם לעומס יתר במסד הנתונים של הסביבה ולכשלים הבאים ב-DAG עם השגיאות 'השרת סגר את החיבור באופן לא צפוי'.
לבצע אופטימיזציה (הפחתה) של השימוש בהוראות xcom_push ו-xcom_pull ב-Jinja Templates בקוד DAG ברמה העליונה.
כדאי לשקול לשדרג לסביבה גדולה יותר (בינונית או גדולה).
מסד הנתונים של Airflow לא היה זמין באופן זמני
יכול להיות שיעבור זמן עד שעובד של Airflow יזהה שגיאות לסירוגין, כמו בעיות זמניות בקישוריות, ויטפל בהן בצורה תקינה. יכול להיות שהיא תעבור את סף ברירת המחדל לזיהוי מכשירי זומבי.
הסבר על פסק זמן (timeout) של אותות דופק ב-Airflow
resource.type="cloud_composer_environment"
resource.labels.environment_name="ENVIRONMENT_NAME"
log_id("airflow-worker")
textPayload:"Heartbeat time limit exceeded"פתרונות:
מגדילים את הזמן הקצוב לתפוגה של משימות זומבי ומבטלים את הערך של אפשרות התצורה
[scheduler]scheduler_zombie_task_thresholdב-Airflow:קטע מפתח ערך הערות schedulerscheduler_zombie_task_thresholdזמן קצוב לתפוגה חדש (בשניות) ערך ברירת המחדל הוא 300.
המשימות נכשלות כי קרתה שגיאה במהלך הביצוע
סיום של מופע
Airflow משתמש במנגנון terminating instance כדי להפסיק משימות של Airflow. המנגנון הזה משמש במצבים הבאים:
- כשתזמן מסיים משימה שלא הושלמה בזמן.
- כשהזמן שמוקצב למשימה מסתיים או שהיא מתבצעת יותר מדי זמן.
כש-Airflow מסיים מופעים של משימות, אפשר לראות את רשומות היומן הבאות ביומנים של Airflow worker שהריץ את המשימה:
INFO - Subtask ... WARNING - State of this instance has been externally set
to success. Terminating instance.
INFO - Subtask ... INFO - Sending Signals.SIGTERM to GPID <X>
INFO - Subtask ... ERROR - Received SIGTERM. Terminating subprocesses.
פתרונות אפשריים:
בודקים את קוד המשימה כדי לראות אם יש שגיאות שעלולות לגרום להרצה שלה למשך זמן ארוך מדי.
הגדלת המעבד והזיכרון של עובדי Airflow, כדי שהמשימות יתבצעו מהר יותר.
מגדילים את הערך של האפשרות
[celery_broker_transport_options]visibility_timeoutבהגדרות של Airflow.כתוצאה מכך, מתזמן המשימות מחכה זמן רב יותר עד שמשימה מסתיימת, לפני שהוא מחשיב את המשימה כמשימת זומבי. האפשרות הזו שימושית במיוחד למשימות שלוקחות הרבה זמן, שנמשכות שעות רבות. אם הערך נמוך מדי (לדוגמה, 3 שעות), מתזמן המשימות יתייחס למשימות שפועלות במשך 5 או 6 שעות כאל משימות תקועות (משימות זומבי).
מגדילים את הערך של אפשרות ההגדרה
[core]killed_task_cleanup_timeAirflow.ערך ארוך יותר מאפשר לעובדי Airflow יותר זמן לסיים את המשימות שלהם בצורה חלקה. אם הערך נמוך מדי, יכול להיות שמשימות Airflow יופסקו בפתאומיות, בלי מספיק זמן לסיים את העבודה בצורה תקינה.
הפעלת DAG לא מסתיימת בתוך הזמן הצפוי
תיאור הבעיה:
לפעמים הפעלת DAG לא מסתיימת כי משימות Airflow נתקעות והפעלת DAG נמשכת יותר מהצפוי. בתנאים רגילים, משימות של Airflow לא נשארות במצב 'בהמתנה בתור' או 'פועל' ללא הגבלת זמן, כי ל-Airflow יש פסק זמן ונהלי ניקוי שעוזרים להימנע ממצב כזה.
פתרון:
משתמשים בפרמטר
dagrun_timeoutעבור DAG. לדוגמה:dagrun_timeout=timedelta(minutes=120). לכן, כל הרצה של DAG צריכה להסתיים לפני שתפוג ההמתנה של הרצת ה-DAG. מידע נוסף על מצבי משימות ב-Airflow זמין במאמרי העזרה של Apache Airflow.משתמשים בפרמטר task execution timeout כדי להגדיר זמן קצוב לתפוגה כברירת מחדל למשימות שמופעלות על סמך אופרטורים של Apache Airflow.
החיבור לשרת Postgres נותק במהלך ביצוע השאילתה או מיד אחריו
חריגים ל-Lost connection to Postgres server during query קורים בדרך כלל כשהתנאים הבאים מתקיימים:
- ה-DAG שלכם משתמש ב-
PythonOperatorאו באופרטור בהתאמה אישית. - ה-DAG שולח שאילתות למסד הנתונים של Airflow.
אם מבוצעות כמה שאילתות מפונקציה שאפשר להפעיל, יכול להיות שה-tracebacks יצביעו באופן שגוי על שורה self.refresh_from_db(lock_for_update=True) בקוד של Airflow. זו השאילתה הראשונה במסד הנתונים אחרי הפעלת המשימה. הסיבה האמיתית לחריגה מתרחשת לפני כן, כשסשן של SQLAlchemy לא נסגר כמו שצריך.
הסשנים של SQLAlchemy מוגבלים לשרשור ונוצרים בפונקציה שאפשר לקרוא לה. אפשר להמשיך את הסשן מאוחר יותר בתוך קוד Airflow. אם יש עיכובים משמעותיים בין שאילתות בסשן אחד, יכול להיות שהחיבור כבר נסגר על ידי שרת Postgres. הזמן הקצוב לתפוגת החיבור בסביבות Managed Airflow מוגדר ל-10 דקות בערך.
פתרון:
- משתמשים ב-decorator
airflow.utils.db.provide_session. הדקורטור הזה מספק סשן תקין למסד הנתונים של Airflow בפרמטרsessionוסוגר את הסשן בצורה נכונה בסוף הפונקציה. - אל תשתמשו בפונקציה אחת שפועלת לאורך זמן. במקום זאת, מעבירים את כל השאילתות במסד הנתונים לפונקציות נפרדות, כך שיהיו כמה פונקציות עם העיטור
airflow.utils.db.provide_session. במקרה כזה, הסשנים נסגרים אוטומטית אחרי שליפת תוצאות השאילתה.
הפרעות זמניות בחיבור למסד הנתונים של מטא-נתונים של Airflow
Managed Airflow פועל על גבי תשתית מבוזרת. המשמעות היא שמעת לעת עשויות להופיע בעיות זמניות, והן עלולות להפריע להרצת המשימות שלכם ב-Airflow.
במקרים כאלה, יכול להיות שיופיעו הודעות השגיאה הבאות ביומני העובדים של Airflow:
"Can't connect to Postgres server on 'airflow-sqlproxy-service.default.svc.cluster.local' (111)"
או
"Can't connect to Postgres server on 'airflow-sqlproxy-service.default.svc.cluster.local' (104)"
בעיות לסירוגין כאלה יכולות להיגרם גם מפעולות תחזוקה שמתבצעות בסביבות Managed Airflow שלכם.
בדרך כלל שגיאות כאלה הן לסירוגין, ואם המשימות שלכם ב-Airflow הן אידמפוטנטיות והגדרתם ניסיונות חוזרים, הן לא ישפיעו עליכם. אפשר גם להגדיר חלונות זמנים לתחזוקה.
סיבה נוספת לשגיאות כאלה יכולה להיות מחסור במשאבים באשכול של הסביבה. במקרים כאלה, אפשר להגדיל את סביבת העבודה או לבצע בה אופטימיזציה לפי ההוראות במאמרים הגדלת סביבות עבודה או אופטימיזציה של סביבת העבודה.
הרצת DAG מסומנת כהצלחה אבל לא בוצעו משימות
אם הפעלת DAG execution_date מוקדמת יותר מ-DAG start_date, יכול להיות שתראו הפעלות DAG שלא כוללות הפעלות משימות, אבל עדיין מסומנות כהפעלות מוצלחות.
הסיבה
המצב הזה יכול לקרות באחד מהמקרים הבאים:
אי התאמה נגרמת בגלל הבדל באזור הזמן בין
execution_dateלביןstart_dateשל DAG. לדוגמה, זה יכול לקרות כשמשתמשים ב-pendulum.parse(...)כדי להגדיר אתstart_date.הערך של
start_dateב-DAG מוגדר כערך דינמי, למשלairflow.utils.dates.days_ago(1)
הפתרון
מוודאים ש
execution_dateו-start_dateמשתמשים באותו אזור זמן.כדי למנוע הפעלה של DAG עם תאריכי התחלה קודמים, צריך לציין
start_dateסטטי ולשלב אותו עםcatchup=False.
שיטות מומלצות
ההשפעה של פעולות עדכון או שדרוג על ביצוע משימות ב-Airflow
פעולות עדכון או שדרוג משבשות משימות Airflow שמופעלות כרגע, אלא אם משימה מופעלת במצב ניתן להשהיה.
מומלץ לבצע את הפעולות האלה כשאתם מצפים להשפעה מינימלית על ההרצות של משימות Airflow, ולהגדיר מנגנוני ניסיון חוזר מתאימים ב-DAG ובמשימות.
לא כדאי לתזמן DAG שנוצר באופן פרוגרמטי באותו זמן
יצירה של אובייקטים של DAG באופן פרוגרמטי מקובץ DAG היא שיטה יעילה ליצירה של הרבה DAG דומים עם הבדלים קטנים בלבד.
חשוב לא לתזמן את כל ה-DAGs האלה להפעלה מיידית. יש סיכוי גבוה שלעובדי Airflow אין מספיק משאבי CPU וזיכרון כדי להריץ את כל המשימות שתוזמנו באותו הזמן.
כדי למנוע בעיות בתזמון של DAGs פרוגרמטיים:
- להגדיל את מספר העובדים בו-זמנית ולהרחיב את הסביבה, כדי שתוכלו לבצע יותר משימות בו-זמנית.
- כדאי ליצור DAGs באופן שיפיץ את התזמונים שלהם באופן שווה לאורך זמן, כדי להימנע מתזמון של מאות משימות בו-זמנית, וכך לתת לעובדי Airflow זמן להריץ את כל המשימות המתוזמנות.
שליטה בזמן הביצוע של DAG, משימות וביצועים מקבילים של אותו DAG
אם רוצים לקבוע כמה זמן יימשך ביצוע של DAG מסוים, אפשר להשתמש בפרמטר dagrun_timeout של ה-DAG. לדוגמה, אם אתם מצפים שריצת DAG אחת (לא משנה אם הביצוע יסתיים בהצלחה או בכישלון) לא תימשך יותר משעה, אתם יכולים להגדיר את הפרמטר הזה ל-3, 600 שניות.
אפשר גם לקבוע את משך הזמן המקסימלי של משימה יחידה ב-Airflow. כדי לעשות זאת, אפשר להשתמש בפקודה execution_timeout.
אם רוצים לקבוע כמה הפעלות של DAG יהיו פעילות עבור DAG מסוים, אפשר להשתמש ב[core]max-active-runs-per-dag
אפשרות ההגדרה של Airflow.
אם רוצים שרק מופע אחד של DAG יפעל בכל רגע נתון, מגדירים את הפרמטר max-active-runs-per-dag לערך 1.
איך נמנעים מעלייה בתנועת הרשת אל מסד הנתונים של Airflow וממנו
כמות התנועה ברשת בין אשכול GKE בסביבה שלכם לבין מסד הנתונים של Airflow תלויה במספר ה-DAG, במספר המשימות ב-DAG ובאופן שבו ה-DAG ניגש לנתונים במסד הנתונים של Airflow. הגורמים הבאים עשויים להשפיע על השימוש ברשת:
שאילתות למסד הנתונים של Airflow. אם קובצי ה-DAG מבצעים הרבה שאילתות, הם יוצרים נפח תנועה גדול. דוגמאות: בדיקת הסטטוס של משימות לפני שממשיכים למשימות אחרות, שליחת שאילתה לטבלת XCom, יצירת dump של תוכן מסד הנתונים של Airflow.
מספר גדול של משימות. ככל שיש יותר משימות לתזמון, כך נוצר יותר נפח תנועה ברשת. השיקול הזה רלוונטי גם למספר הכולל של המשימות ב-DAG וגם לתדירות התזמון. כשמתזמן Airflow מתזמן הפעלות של DAG, הוא שולח שאילתות למסד הנתונים של Airflow ויוצר תנועה.
ממשק האינטרנט של Airflow יוצר תעבורת נתונים ברשת כי הוא שולח שאילתות למסד הנתונים של Airflow. שימוש אינטנסיבי בדפים עם גרפים, משימות ודיאגרמות עלול ליצור נפח גדול של תעבורת רשת.