בדף הזה מוסבר על מחזור החיים של צינור, החל מקוד הצינור ועד למשימת Dataflow.
בדף הזה מוסברים המושגים הבאים:
- מהו תרשים ביצוע ואיך צינור עיבוד נתונים של Apache Beam הופך לעבודת Dataflow
- איך Dataflow מטפל בשגיאות
- איך Dataflow מבצע באופן אוטומטי מקביליות ומפיץ את לוגיקת העיבוד בצינור אל העובדים שמבצעים את העבודה
- אופטימיזציות של משימות ש-Dataflow עשוי לבצע
תרשים ביצוע
כשמריצים את צינור הנתונים של Dataflow, Dataflow יוצר גרף ביצוע מהקוד שבונה את אובייקט Pipeline, כולל כל הטרנספורמציות ופונקציות העיבוד המשויכות להן, כמו אובייקטים של DoFn. זהו תרשים הביצוע של צינור העיבוד, והשלב נקרא זמן בניית התרשים.
במהלך בניית הגרף, Apache Beam מריץ באופן מקומי את הקוד מנקודת הכניסה הראשית של קוד צינור עיבוד הנתונים, עוצר בקריאות לשלב של מקור, יעד או טרנספורמציה, והופך את הקריאות האלה לצמתים בגרף.
לכן, קטע קוד בנקודת הכניסה של צינור עיבוד נתונים (שיטת Java ו-Go main
או הרמה העליונה של סקריפט Python) מורץ באופן מקומי במחשב שמריץ את צינור עיבוד הנתונים. אותו קוד שמוצהר בשיטה של אובייקט DoFn מופעל בתהליכי העבודה של Dataflow.
לדוגמה, הדוגמה WordCount שכלולה בערכות ה-SDK של Apache Beam מכילה סדרה של טרנספורמציות לקריאה, לחילוץ, לספירה, לעיצוב ולכתיבה של המילים הנפרדות באוסף של טקסט, יחד עם ספירת המופעים של כל מילה. בתרשים הבא מוצג תהליך ההרחבה של הטרנספורמציות בצינור WordCount לגרף ביצוע:

איור 1: תרשים של ביצוע לדוגמה של WordCount
תרשים הביצוע שונה לעיתים קרובות מהסדר שבו ציינתם את ההמרות כשבניתם את צינור העיבוד. ההבדל הזה נובע מכך ששירות Dataflow מבצע אופטימיזציות ומיזוגים שונים בגרף הביצוע לפני שהוא פועל על משאבי ענן מנוהלים. שירות Dataflow מכבד את התלות בנתונים כשמריצים את צינור הנתונים. עם זאת, יכול להיות ששלבים ללא תלות בנתונים יפעלו בכל סדר.
כדי לראות את תרשים הביצוע הלא אופטימלי ש-Dataflow יצר עבור צינור עיבוד הנתונים, בוחרים את העבודה בממשק המעקב של Dataflow. מידע נוסף על הצגת משימות זמין במאמר שימוש בממשק המעקב של Dataflow.
במהלך בניית הגרף, Apache Beam מאמת שכל המשאבים שאליהם מתייחס צינור עיבוד הנתונים, כמו קטגוריות של Cloud Storage, טבלאות של BigQuery ונושאים או מינויים של Pub/Sub, אכן קיימים ונגישים. האימות מתבצע באמצעות קריאות רגילות ל-API של השירותים הרלוונטיים, ולכן חשוב שלחשבון המשתמש שמשמש להרצת צינור יש קישוריות תקינה לשירותים הנדרשים והרשאה לקרוא לממשקי ה-API של השירותים. לפני ששולחים את צינור הנתונים לשירות Dataflow, Apache Beam בודק גם שגיאות אחרות, ומוודא שגרף צינור הנתונים לא מכיל פעולות לא חוקיות.
גרף הביצוע מתורגם לפורמט JSON, וגרף הביצוע בפורמט JSON מועבר לנקודת הקצה של שירות Dataflow.
לאחר מכן, שירות Dataflow מאמת את גרף ההפעלה בפורמט JSON. אחרי שהגרף מאומת, הוא הופך לעבודה בשירות Dataflow. אפשר לראות את העבודה, את תרשים הביצוע שלה, את הסטטוס ואת פרטי היומן באמצעות ממשק המעקב של Dataflow.
Java
שירות Dataflow שולח תגובה למכונה שבה מריצים את תוכנית Dataflow. התשובה הזו מוכלת באובייקט DataflowPipelineJob, שמכיל את jobId של משימת Dataflow.
אפשר להשתמש ב-jobId כדי לעקוב אחרי המשימה, לנטר אותה ולפתור בה בעיות באמצעות ממשק המעקב של Dataflow וממשק שורת הפקודה של Dataflow.
מידע נוסף מופיע במאמרי העזרה של ה-API של DataflowPipelineJob.
Python
שירות Dataflow שולח תגובה למכונה שבה מריצים את תוכנית Dataflow. התשובה הזו מוכלת באובייקט DataflowPipelineResult, שמכיל את job_id של משימת Dataflow.
אפשר להשתמש בjob_id כדי לעקוב אחרי המשימה, לנטר אותה ולפתור בה בעיות באמצעות ממשק המעקב של Dataflow וממשק שורת הפקודה של Dataflow.
המשך
שירות Dataflow שולח תגובה למכונה שבה מריצים את תוכנית Dataflow. התשובה הזו מוכלת באובייקט dataflowPipelineResult, שמכיל את jobID של משימת Dataflow.
אפשר להשתמש בjobID כדי לעקוב אחרי המשימה, לנטר אותה ולפתור בה בעיות באמצעות ממשק המעקב של Dataflow וממשק שורת הפקודה של Dataflow.
בניית הגרף מתבצעת גם כשמריצים את צינור הנתונים באופן מקומי, אבל הגרף לא מתורגם ל-JSON ולא מועבר לשירות. במקום זאת, הגרף מופעל באופן מקומי באותו מחשב שבו הפעלתם את תוכנית Dataflow. מידע נוסף זמין במאמר הגדרת PipelineOptions להרצה מקומית.
טיפול בשגיאות ובחריגים
יכול להיות שצינור הנתונים יזרוק חריגים במהלך עיבוד הנתונים. חלק מהשגיאות האלה הן זמניות, כמו קושי זמני בגישה לשירות חיצוני. שגיאות אחרות הן קבועות, כמו שגיאות שנגרמות בגלל נתוני קלט פגומים או כאלה שלא ניתן לנתח, או בגלל מצביעים ריקים במהלך החישוב.
Dataflow מעבד רכיבים בחבילות שרירותיות, ומנסה שוב לעבד את החבילה כולה אם מתקבלת שגיאה לגבי רכיב כלשהו בחבילה. כשמפעילים את התהליך במצב אצווה, המערכת מנסה שוב ארבע פעמים חבילות שכוללות פריט שנכשל. הצינור נכשל לחלוטין אם חבילה אחת נכשלת ארבע פעמים. כשמפעילים את הצינור במצב סטרימינג, המערכת מנסה שוב ושוב להפעיל חבילה שכוללת פריט שנכשל, ולכן יכול להיות שהצינור ייעצר לתמיד.
כשמעבדים נתונים במצב אצווה, יכול להיות שיופיע מספר גדול של כשלים פרטניים לפני שמשימת צינור נתונים נכשלת לחלוטין. זה קורה כשחבילה מסוימת נכשלת אחרי ארבעה ניסיונות חוזרים. לדוגמה, אם צינור העיבוד מנסה לעבד 100 חבילות, יכול להיות ש-Dataflow ייצור כמה מאות כשלים נפרדים עד שחבילה אחת תגיע לתנאי של ארבעה כשלים ליציאה.
שגיאות בהפעלת העובדים, כמו כשל בהתקנת חבילות בעובדים, הן זמניות. במקרה כזה, המערכת תנסה שוב ושוב ללא הגבלה, והצינור עלול להיתקע לתמיד.
הקבלה והפצה
שירות Dataflow מבצע אוטומטית הקבלה וחלוקה של לוגיקת העיבוד בצינור הנתונים בין תהליכי Worker ושרשורים.
Dataflow משתמש בהפשטות במודל התכנות כדי לייצג פונקציות של עיבוד מקביל. לדוגמה, טרנספורמציות של ParDo גורמות ל-Dataflow להפיץ קוד עיבוד, שמיוצג על ידי אובייקטים של DoFn, לכמה עובדים כדי שהם יפעלו במקביל.
Dataflow תומך בשני מימדים משלימים של מקביליות:
- מקביליות אופקית: נתוני הפייפליין מפוצלים ומעובדים בכמה מופעי worker בו-זמנית, ומנוהלים באופן דינמי באמצעות התאמה אופקית לעומס (autoscaling).
- מקביליות אנכית: נתוני צינורות מעובדים בכמה ליבות CPU ושרשורים בכל מכונת עובד וירטואלית, ומנוהלים באמצעות התאמה דינמית של מספר השרשורים והתאמה אנכית לעומס.
Dataflow מנהל באופן אוטומטי את המקביליות של העבודות, מטפל באיזון מחדש דינמי של העבודה ומבצע אופטימיזציה של גרף הביצוע באמצעות מיזוג ואופטימיזציה של שילובים. גורמים כמו מקורות נתונים שלא ניתן לפצל, שלבים עם fan-out רחב, הטיה של מפתחות ומגבלות של sink בהמשך ה-פייפליין יכולים להגביל את המקביליות של פייפליינים.
במאמר בנושא מקביליות ב-Dataflow מוסבר איך Dataflow מחלק את הנתונים, משנה את קנה המידה של העובדים ופותר צווארי בקבוק של מקביליות.
אופטימיזציה של מיזוג
אחרי שגרף הביצוע של צינור הנתונים בפורמט JSON עובר אימות, יכול להיות ששירות Dataflow ישנה את הגרף כדי לבצע אופטימיזציות.
האופטימיזציות יכולות לכלול מיזוג של כמה שלבים או טרנספורמציות בתרשים הביצוע של צינור הנתונים לשלבים בודדים. מיזוג השלבים מונע מהשירות Dataflow להפוך כל PCollection ביניים בצינור העיבוד למוחשי, מה שיכול להיות יקר מבחינת הזיכרון והתקורה של העיבוד.
למרות שכל השינויים שאתם מציינים בבניית צינור הנתונים מבוצעים בשירות, כדי להבטיח את הביצוע היעיל ביותר של צינור הנתונים, יכול להיות שהשינויים יבוצעו בסדר שונה או כחלק משינוי גדול יותר. שירות Dataflow מתחשב בתלות של נתונים בין השלבים בתרשים הביצוע, אבל מעבר לכך, השלבים עשויים להתבצע בכל סדר.
דוגמה ל-Fusion
בתרשים הבא אפשר לראות איך אפשר לבצע אופטימיזציה לגרף הביצוע מהדוגמה WordCount שכלולה ב-Apache Beam SDK for Java, ולמזג אותו באמצעות שירות Dataflow כדי לבצע את הפעולה בצורה יעילה:

איור 2: דוגמה לתרשים ביצועים שעבר אופטימיזציה של WordCount
מניעת מיזוג
במקרים מסוימים, יכול להיות שמערכת Dataflow תנחש בצורה לא נכונה את הדרך האופטימלית למיזוג פעולות בצינור, מה שיגביל את היכולת של Dataflow להשתמש בכל העובדים הזמינים. במקרים כאלה, אפשר להשתמש בטרנספורמציה Redistribute כדי לתת ל-Dataflow רמז לחלוקה מחדש של הנתונים.
כדי להוסיף טרנספורמציה של Redistribute, קוראים לאחת מהשיטות הבאות:
Redistribute.arbitrarily: מציין שהנתונים כנראה לא מאוזנים. Dataflow בוחר את האלגוריתם הטוב ביותר לחלוקה מחדש של הנתונים.
Redistribute.byKey: מציין שסביר להניח ש-PCollectionשל זוגות מפתח/ערך לא מאוזן, וצריך לבצע חלוקה מחדש על סמך המפתחות. בדרך כלל, Dataflow ממקם את כל הרכיבים של מפתח יחיד באותו Thread עובד. עם זאת, אין הבטחה שהמפתחות ימוקמו באותו מקום, והרכיבים מעובדים באופן עצמאי.
אם צינור הנתונים מכיל טרנספורמציה Redistribute, בדרך כלל Dataflow מונע מיזוג של השלבים לפני ואחרי הטרנספורמציה Redistribute, ומבצע ערבוב של הנתונים כדי שהשלבים בהמשך הצינור אחרי הטרנספורמציה Redistribute יפעלו במקביל בצורה אופטימלית יותר.
מעקב אחר מיזוג
אפשר לגשת לגרף המותאם ולשלבים הממוזגים במסוף Google Cloud , באמצעות ה-CLI של gcloud או באמצעות ה-API.
המסוף
כדי לראות את השלבים והפעולות הממוזגים בתרשים במסוף, פותחים את תצוגת התרשים Stage workflow בכרטיסייה Execution details של משימת Dataflow.
כדי לראות את השלבים שמוזגו לשלב מסוים, לוחצים על השלב הממוזג בגרף. בחלונית פרטי השלב, השלבים המאוחדים מוצגים בשורה Component steps. לפעמים חלקים של טרנספורמציה מורכבת אחת משולבים בכמה שלבים.
gcloud
כדי לגשת לגרף האופטימלי ולשלבים הממוזגים באמצעות ה-CLI של gcloud, מריצים את הפקודה הבאה gcloud:
gcloud dataflow jobs describe --full JOB_ID --format json
מחליפים את JOB_ID במזהה של משימת Dataflow.
כדי לחלץ את החלקים הרלוונטיים, מעבירים את הפלט של הפקודה gcloud ל-jq:
gcloud dataflow jobs describe --full JOB_ID --format json | jq '.pipelineDescription.executionPipelineStage\[\] | {"stage_id": .id, "stage_name": .name, "fused_steps": .componentTransform }'
כדי לראות את התיאור של השלבים הממוזגים בקובץ התגובה של הפלט, בתוך המערך ComponentTransform, צריך לעיין באובייקט ExecutionStageSummary.
API
כדי לגשת לתרשים האופטימלי ולשלבים המאוחדים באמצעות ה-API, צריך להפעיל את project.locations.jobs.get.
כדי לראות את התיאור של השלבים הממוזגים בקובץ התגובה של הפלט, בתוך המערך ComponentTransform, צריך לעיין באובייקט ExecutionStageSummary.
אופטימיזציה משולבת
פעולות צבירה הן מושג חשוב בעיבוד נתונים בקנה מידה גדול.
צבירה מאחדת נתונים שרחוקים זה מזה מבחינה מושגית, ולכן היא שימושית מאוד לצורך קורלציה. מודל התכנות של Dataflow מייצג פעולות צבירה כטרנספורמציות GroupByKey, CoGroupByKey ו-Combine.
פעולות הצבירה של Dataflow משלבות נתונים בכל מערך הנתונים, כולל נתונים שאולי מפוזרים בין כמה עובדים. במהלך פעולות צבירה כאלה, לרוב הכי יעיל לשלב כמה שיותר נתונים באופן מקומי לפני שמשלבים נתונים בין מופעים. כשמחילים טרנספורמציה של GroupByKey או טרנספורמציה מצטברת אחרת, שירות Dataflow מבצע באופן אוטומטי שילוב חלקי באופן מקומי לפני פעולת הקיבוץ הראשית.
כשמבצעים שילוב חלקי או שילוב בכמה רמות, שירות Dataflow מקבל החלטות שונות בהתאם לסוג הנתונים שצינור הנתונים עובד איתם – נתונים באצווה או נתונים בסטרימינג. בנתונים מוגבלים, השירות מעדיף יעילות ויבצע כמה שיותר שילובים מקומיים. בנתונים לא מוגבלים, השירות מעדיף חביון נמוך, ויכול להיות שהוא לא יבצע שילוב חלקי, כי הוא עלול להגדיל את החביון.