Managed Airflow (דור 3) | Managed Airflow (דור 2) | Managed Airflow (דור 1 מדור קודם)
בדף הזה מפורטים שלבים לפתרון בעיות ומידע על בעיות נפוצות בתהליכי עבודה.
הרבה בעיות בהרצת DAG נובעות מביצועים לא אופטימליים של הסביבה. כדי לבצע אופטימיזציה של הסביבה, אפשר להיעזר במדריך אופטימיזציה של הביצועים והעלויות של הסביבה.
חלק מהבעיות בהפעלות של DAG יכולות להיגרם בגלל שתזמן Airflow לא פועל בצורה תקינה או אופטימלית. כדי לפתור את הבעיות האלה, צריך לפעול לפי ההוראות לפתרון בעיות בתכנון פגישות.
פתרון בעיות באמצעות Managed Airflow Agent
סוכן Managed Airflow יכול לעזור לכם להבין, לאבחן ולפתור בעיות במשימות Airflow שנכשלו ובריצות DAG, על ידי ניתוח יומני ביצוע של משימות, קוד מקור של DAG ומדדי סביבה. מידע נוסף זמין במאמר בנושא פתרון בעיות במשימות ובריצות של DAG ב-Managed Airflow Agent.
תהליך עבודה לפתרון בעיות
כדי להתחיל לפתור את הבעיה:
בודקים את היומנים של Airflow.
כדי להגדיל את רמת הרישום ביומן של Airflow, צריך לבטל את ההגדרה של אפשרות התצורה הבאה של Airflow.
קטע מפתח ערך logginglogging_levelערך ברירת המחדל הוא INFO. מגדירים את הערך ל-DEBUGכדי לקבל רמת פירוט גבוהה יותר בהודעות היומן.בודקים את לוח הבקרה של Monitoring.
בודקים את Cloud Monitoring.
במסוף Google Cloud , בודקים אם יש שגיאות בדפים של רכיבי הסביבה.
בממשק האינטרנט של Airflow, בודקים בתצוגת הגרף של ה-DAG אם יש מופעים של משימות שנכשלו.
קטע מפתח ערך webserverdag_orientationLR,TB,RLאוBT
ניפוי באגים בכשלים של אופרטורים
כדי לנפות באגים בכשל של אופרטור:
- בודקים אם יש שגיאות שקשורות למשימה.
- בודקים את היומנים של Airflow.
- בודקים את Cloud Monitoring.
- בודקים יומנים ספציפיים למפעיל.
- מתקנים את השגיאות.
- מעלים את ה-DAG לתיקייה
/dags. - בממשק האינטרנט של Airflow, מנקים את המצבים הקודמים של ה-DAG.
- ממשיכים או מפעילים את ה-DAG.
פתרון בעיות בביצוע משימות
Airflow היא מערכת מבוזרת עם הרבה ישויות כמו מתזמן, מפעיל, עובדים שתקשרים ביניהם דרך תור משימות ומסד הנתונים של Airflow, ושולחים אותות (כמו SIGTERM). התרשים הבא מציג סקירה כללית של הקישורים בין רכיבי Airflow.
במערכת מבוזרת כמו Airflow, יכולות להיות בעיות בקישוריות לרשת, או בעיות לסירוגין בתשתית הבסיסית. זה יכול להוביל למצבים שבהם משימות נכשלות ומתוזמנות מחדש לביצוע, או שמשימות לא מושלמות בהצלחה (לדוגמה, משימות זומבי או משימות שנתקעו במהלך הביצוע). ל-Airflow יש מנגנונים להתמודדות עם מצבים כאלה, והוא מחדש באופן אוטומטי את הפעולה הרגילה. בקטעים הבאים מוסברות בעיות נפוצות שמתרחשות במהלך ביצוע משימות על ידי Airflow.
משימות נכשלות בלי ליצור יומנים
יכולות להיות כמה סיבות לכך שמופע של משימה ייכשל בלי שיופקו יומנים. לדוגמה, זה יכול לקרות בגלל שגיאות בניתוח של 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, והקובץ עדיין לא מסונכרן עם עובד שמבצע משימה, המשימה תיכשל ללא יומנים, ויכול להיות שיופיע airflow.exceptions.AirflowException: Dag "example-dag" could not be
found... ביומני העובדים.
הסנכרון הזה בדרך כלל נמשך דקה או שתיים, אבל הוא יכול להימשך יותר זמן אם יש לכם הרבה קבצים או קבצים גדולים בתיקיות 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.מרוץ תהליכים בין התקשרות חזרה של פעימת הלב לבין התקשרויות חזרה של יציאה ב-local_task_job, שעוקב אחרי הרצת המשימה. אם בדיקת הפעילות מזהה שמשימה סומנה כהצלחה, היא לא יכולה להבחין בין מצב שבו המשימה עצמה הצליחה לבין מצב שבו Airflow קיבל הוראה להתייחס למשימה כאל משימה שהצליחה. עם זאת, הוא יסיים את הפעולה של מפעיל המשימות בלי לחכות לסיום שלו.
אפשר להתעלם בבטחה מאותות SIGTERM כאלה. המשימה כבר במצב מוצלח, והביצוע של הרצת ה-DAG כולה לא יושפע.
רשומת היומן
Received SIGTERM.היא ההבדל היחיד בין יציאה רגילה לבין סיום המשימה במצב מוצלח.איור 2. מרוץ תהליכים בין אותות פעימת הלב לבין קריאות חוזרות (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 (דור 2), אפשר להקצות יותר משאבי 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 (דור 2), מגדילים את מגבלות הזיכרון של עובדי Airflow.
אם הסביבה שלכם יוצרת גם משימות זומבי, כדאי לעיין במאמר בנושא פתרון בעיות שקשורות למשימות זומבי.
מדריך לניפוי באגים בבעיות שקשורות לזיכרון או לאחסון מלאים זמין במאמר ניפוי באגים בבעיות שקשורות לזיכרון או לאחסון מלאים ב-DAG.
המשימה נכשלת בגלל הוצאה של Pod
ה-Pods של Google Kubernetes Engine כפופים למחזור החיים של Kubernetes Pod ולפינוי Pod. הסיבה הנפוצה ביותר להוצאת Pod משימוש ב-Managed Airflow היא קפיצות בשימוש במשאבים של משימות.
פינוי של Pod יכול להתרחש כש-Pod מסוים משתמש יתר על המידה במשאבים של צומת, ביחס לציפיות לגבי צריכת המשאבים שהוגדרו לצומת. לדוגמה, יכול להיות שפינוי יקרה כשכמה משימות שצורכות הרבה זיכרון פועלות ב-Pod, והעומס המשולב שלהן גורם לחריגה ממגבלת צריכת הזיכרון של הצומת שבו ה-Pod הזה פועל.
אם פוד של Airflow worker מפונה, כל מופעי המשימות שפועלים בפוד הזה מופסקים, ומאוחר יותר מסומנים ככאלה שנכשלו על ידי Airflow.
היומנים נשמרים בזיכרון הזמני. אם פוד של Worker מוצא מהמערכת לפני שמאגר הנתונים הזמני מתרוקן, היומנים לא מופקים. אם המשימה נכשלת בלי יומנים, זה מצביע על כך שתהליכי העבודה של Airflow מופעלים מחדש בגלל חוסר זיכרון (OOM). יכול להיות שחלק מהיומנים יופיעו ב-Cloud Logging גם אם יומני Airflow לא הופקו.
כדי לראות את היומנים:
במסוף Google Cloud , עוברים לדף Environments.
ברשימת הסביבות, לוחצים על שם הסביבה. הדף Environment details ייפתח.
עוברים לכרטיסייה יומנים.
אפשר לראות את היומנים של כל אחד מה-workers של Airflow בקטע All logs > Airflow logs > Workers.
תיאור הבעיה:
במסוף Google Cloud , עוברים לדף Workloads.
אם יש
airflow-workerפודים שמוצג בהםEvicted, לוחצים על כל פוד שהוצא ומחפשים את ההודעהThe node was low on resource: memoryבחלק העליון של החלון.
פתרון:
הגדלת מגבלות הזיכרון לעובדי Airflow.
בודקים את היומנים של תרמילי
airflow-workerכדי לראות מהן הסיבות האפשריות להוצאה. מידע נוסף על אחזור יומנים מ-Pods ספציפיים זמין במאמר פתרון בעיות בעומסי עבודה שנפרסו.מוודאים שהמשימות ב-DAG הן אידמפוטנטיות וניתנות לניסיון חוזר.
מומלץ להימנע מהורדה של קבצים מיותרים למערכת הקבצים המקומית של עובדי Airflow.
ל-Airflow workers יש קיבולת מוגבלת במערכת הקבצים המקומית. לכל עובד ב-Airflow יכול להיות נפח אחסון של 1GB עד 10GB. כשנגמר נפח האחסון, מישור הבקרה של GKE מפנה את ה-Pod של Airflow worker. כל המשימות שהעובד שהוצא ביצע ייכשלו.
דוגמאות לפעולות בעייתיות:
- הורדת קבצים או אובייקטים ואחסונם באופן מקומי ב-worker של Airflow. במקום זאת, כדאי לאחסן את האובייקטים האלה ישירות בשירות מתאים, כמו קטגוריה של Cloud Storage.
- גישה לאובייקטים גדולים בתיקייה
/dataמעובד של Airflow. תהליך העבודה של Airflow מוריד את האובייקט למערכת הקבצים המקומית שלו. במקום זאת, כדאי להטמיע את ה-DAG כך שקבצים גדולים יעובדו מחוץ ל-Pod של Airflow worker.
סיבות נפוצות לבעיה
הזיכרון של Airflow worker אזל
כל תהליך עבודה של Airflow יכול להריץ עד [celery]worker_concurrency מופעים של משימות בו-זמנית. אם צריכת הזיכרון המצטברת של מופעי המשימות האלה חורגת ממגבלת הזיכרון של עובד Airflow, תהליך אקראי יופסק כדי לפנות משאבים.
גילוי אירועים של חוסר זיכרון (OOM) בתהליך העבודה של Airflow
resource.type="k8s_node"
resource.labels.cluster_name="GKE_CLUSTER_NAME"
log_id("events")
jsonPayload.message:"Killed process"
jsonPayload.message:("airflow task" OR "celeryd")לפעמים, מחסור בזיכרון בתהליך עובד של 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.בגרסאות של Managed Airflow (דור 2) שקודמות לגרסה 2.6.0, צריך לעדכן את
[celery]worker_concurrencyבאמצעות הנוסחה הנוכחית אם הערך הזה נמוך יותר.
הוצאה של Airflow worker
הוצאת pods היא חלק רגיל מהפעלת עומסי עבודה ב-Kubernetes. GKE מוציא משימוש פודים אם נגמר להם האחסון או כדי לפנות משאבים לעומסי עבודה עם עדיפות גבוהה יותר.
הסרת עובדים מ-Airflow
resource.type="k8s_pod"
resource.labels.cluster_name="GKE_CLUSTER_NAME"
resource.labels.pod_name:"airflow-worker"
log_id("events")
jsonPayload.reason="Evicted"פתרונות:
- אם הפינוי נגרם בגלל חוסר מקום באחסון, אפשר לצמצם את השימוש באחסון או להסיר קבצים זמניים ברגע שהם לא נחוצים יותר.
לחלופין, אפשר להגדיל את נפח האחסון הזמין או להריץ עומסי עבודה בתא ייעודי עם
KubernetesPodOperator.
תהליך העבודה של Airflow הופסק
יכול להיות שעובדי Airflow יוסרו באופן חיצוני. אם משימות שפועלות כרגע לא מסתיימות במהלך תקופת סיום מבוקר, הן מופסקות ויכול להיות שהן יזוהו כזומבים.
גילוי סיום של פודים של עובדים ב-Airflow
resource.type="k8s_cluster" resource.labels.cluster_name="GKE_CLUSTER_NAME" protoPayload.methodName:"pods.delete" protoPayload.response.metadata.name:"airflow-worker"
תרחישים ופתרונות אפשריים:
הפעלת העובדים של Airflow מחדש במהלך שינויים בסביבה, כמו שדרוגים או התקנת חבילות:
גילוי שינויים בסביבת ה-Composer
resource.type="cloud_composer_environment" resource.labels.environment_name="ENVIRONMENT_NAME" log_id("cloudaudit.googleapis.com%2Factivity")אפשר לבצע פעולות כאלה כשלא מופעלות משימות קריטיות, או להפעיל ניסיונות חוזרים של משימות.
יכול להיות שרכיבים שונים לא יהיו זמינים באופן זמני במהלך פעולות תחזוקה.
הכרת פעולות התחזוקה של GKE
resource.type="gke_nodepool" resource.labels.cluster_name="GKE_CLUSTER_NAME" protoPayload.metadata.operationType="UPGRADE_NODES"
אתם יכולים לציין חלונות זמן לתחזוקה כדי לצמצם
חפיפה עם הביצוע של המשימות הקריטיות.
בגרסאות של Managed Airflow (דור 2) שקודמות לגרסה 2.4.5, יכול להיות שעובד של Airflow שמסיים את הפעולה שלו יתעלם מאות SIGTERM וימשיך להריץ משימות:
התאמה אוטומטית לעומס (autoscaling) ב-Composer
resource.type="cloud_composer_environment" resource.labels.environment_name="ENVIRONMENT_NAME" log_id("airflow-worker-set") textPayload:"Workers deleted"אפשר לשדרג לגרסה מאוחרת יותר של Managed Airflow שבה הבעיה הזו נפתרה.
תהליך העבודה של זרימת האוויר היה תחת עומס כבד
כמות משאבי המעבד והזיכרון שזמינים לעובד של Airflow מוגבלת על ידי הגדרות הסביבה. אם ניצול המשאבים מתקרב למגבלות, יכול להיות שזה יגרום למאבק על משאבים ולעיכובים מיותרים במהלך ביצוע המשימה. במקרים קיצוניים, כשאין מספיק משאבים לאורך זמן, יכול להיות שייווצרו משימות זומבי.
פתרונות:
- עוקבים אחרי השימוש במעבד ובשימוש בזיכרון של ה-worker-ים ומשנים אותו כדי שלא יעלה על 80%.
שאילתות ב-Cloud Logging כדי לגלות את הסיבות להפעלה מחדש או להוצאה של Pod
הסביבות של Managed Airflow משתמשות באשכולות GKE כשכבת תשתית מחשוב. בקטע הזה אפשר למצוא שאילתות שימושיות שיעזרו לכם למצוא סיבות להפעלה מחדש או להוצאה של עובד או מתזמן של Airflow.
אפשר לשנות את השאילתות שמוצגות בהמשך באופן הבא:
אתם יכולים לציין את ציר הזמן הנדרש ב-Cloud Logging. לדוגמה, 6 השעות האחרונות, 3 הימים האחרונים, או שאתם יכולים להגדיר טווח זמן מותאם אישית.
צריך לציין את שם האשכול של הסביבה ב-CLUSTER_NAME.
אפשר להגביל את החיפוש לפוד ספציפי על ידי הוספת התו POD_NAME.
גילוי מאגרי תגים שהופעלו מחדש
resource.type="k8s_node"
log_id("kubelet")
jsonPayload.MESSAGE:"will be restarted"
resource.labels.cluster_name="CLUSTER_NAME"
שאילתה חלופית להגבלת התוצאות ל-Pod ספציפי:
resource.type="k8s_node"
log_id("kubelet")
jsonPayload.MESSAGE:"will be restarted"
resource.labels.cluster_name="CLUSTER_NAME"
"POD_NAME"
גילוי השבתה של קונטיינרים כתוצאה מאירוע של חוסר זיכרון
resource.type="k8s_node"
log_id("events")
(jsonPayload.reason:("OOMKilling" OR "SystemOOM")
OR jsonPayload.message:("OOM encountered" OR "out of memory"))
severity=WARNING
resource.labels.cluster_name="CLUSTER_NAME"
שאילתה חלופית להגבלת התוצאות ל-Pod ספציפי:
resource.type="k8s_node"
log_id("events")
(jsonPayload.reason:("OOMKilling" OR "SystemOOM")
OR jsonPayload.message:("OOM encountered" OR "out of memory"))
severity=WARNING
resource.labels.cluster_name="CLUSTER_NAME"
"POD_NAME"
איך מוצאים קונטיינרים שהפסיקו לפעול
resource.type="k8s_node"
log_id("kubelet")
jsonPayload.MESSAGE:"ContainerDied"
severity=DEFAULT
resource.labels.cluster_name="CLUSTER_NAME"
שאילתה חלופית להגבלת התוצאות ל-Pod ספציפי:
resource.type="k8s_node"
log_id("kubelet")
jsonPayload.MESSAGE:"ContainerDied"
severity=DEFAULT
resource.labels.cluster_name="CLUSTER_NAME"
"POD_NAME"
מסד הנתונים של 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 workers, כדי שהמשימות יתבצעו מהר יותר.
מגדילים את הערך של האפשרות
[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.
אם מבוצעות כמה שאילתות מפונקציה שאפשר להפעיל, יכול להיות שפרטי השגיאה יצביעו באופן שגוי על שורה self.refresh_from_db(lock_for_update=True) מסוימת בקוד של Airflow. זו השאילתה הראשונה במסד הנתונים אחרי הפעלת המשימה. הסיבה האמיתית לחריגה מתרחשת לפני כן, כשסשן של SQLAlchemy לא נסגר כמו שצריך.
הסשנים של SQLAlchemy מוגבלים לשרשור ונוצרים בפונקציה שאפשר לקרוא לה. אפשר להמשיך את הסשן מאוחר יותר בתוך קוד Airflow. אם יש עיכובים משמעותיים בין שאילתות בסשן אחד, יכול להיות שהחיבור כבר נסגר על ידי שרת Postgres. התוקף של החיבור בסביבות מנוהלות של 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 היא שיטה יעילה ליצירת DAGs דומים רבים שיש ביניהם רק הבדלים קטנים.
חשוב לא לתזמן את כל ה-DAG האלה להפעלה מיידית. יש סיכוי גבוה שלעובדי 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. שימוש אינטנסיבי בדפים עם גרפים, משימות ודיאגרמות עלול ליצור נפח גדול של תעבורת רשת.