Dataflow נועד להריץ צינורות עיבוד נתונים גדולים על ידי חלוקת העבודה בין מאגר מנוהל של מופעי מחשוב. הבנה של האופן שבו Dataflow מבצע עיבוד מקביל עוזרת לתכנן צינורות יעילים, להימנע מנקודות צוואר בקבוק בביצועים ולבצע אופטימיזציה של עלויות המשאבים.
בדף הזה מוסבר איך Dataflow מבצעת עיבוד נתונים במקביל, איך היא מנהלת את הביצוע ומרחיבה אותו, מהם הגורמים הנפוצים שמגבילים את המקביליות ומהן הטכניקות שבהן אפשר להשתמש כדי לייעל את קצב העברת הנתונים בצינור.
מודלים של מקביליות: אופקי לעומת אנכי
Dataflow משיג מקביליות באמצעות שתי אסטרטגיות משלימות:
מקביליות אופקית: נתוני צינור עיבוד הנתונים מחולקים ומעובדים בו-זמנית בכמה מופעי עובדים (מכונות וירטואליות). Dataflow יכול להתאים אוטומטית את הגודל של מאגר העובדים בהתאם לביקוש לעומס העבודה באמצעות התאמה אוטומטית אופקית לעומס. כברירת מחדל, ב-Dataflow מוגדרת מגבלת משאבים של 4,000 עובדים לכל משימה. אפשר לשנות את המגבלה באמצעות בקשות להגדלת מכסה.
מקביליות אנכית: מספר ליבות CPU ותהליכים בתוך מופע עובד יחיד מעבדים נתונים של צינורות במקביל. כל מכונה וירטואלית של Worker מריצה תהליכי Worker ושרשורי Harness כדי להשתמש במשאבי מחשוב זמינים. באמצעות התאמה דינמית של מספר השרשורים, Dataflow יכולה לשנות את מספר השרשורים הפעילים לכל עובד בצינורות נתונים של אצווה על סמך השימוש במעבד והזיכרון הפנוי. ב-Dataflow Prime, התכונה Vertical Autoscaling משנה באופן דינמי את גודל הזיכרון והמחשוב שמוקצים לעובדים.
יחידות עבודה והיררכיית ביצוע
כדי לחלק את העיבוד בין תהליכים ובין threads, Dataflow מפרק צינורות של Apache Beam ליחידות עבודה נפרדות:
- PCollections ומחיצות:
PCollectionמייצג מערך נתונים מבוזר. בנתונים מוגבלים (צינורות עיבוד נתונים של אצווה), Dataflow מחלק את מערך הנתונים לפי פיצולים או רסיסים. במקרה של נתונים לא מוגבלים (צינורות עיבוד נתונים של סטרימינג), הנתונים מגיעים באופן רציף ומוטמעים כהודעות או כמחיצות של סטרימינג. - חבילות: Dataflow מקבץ רכיבים לחבילות שרירותיות לעיבוד על ידי
DoFn. החבילה היא יחידת הכשל והניסיון החוזר: אם העיבוד של רכיב מסוים מעלה חריגה שלא טופלה, המערכת מנסה לעבד מחדש את החבילה כולה. פעולות שצורכות הרבה זיכרון יכולות להגביר את העומס על הזיכרון של העובד ולהוביל לשגיאות שקשורות לחוסר זיכרון. - מיזוג של שלבים: במהלך האופטימיזציה של הגרף, Dataflow משלב טרנספורמציות סמוכות לשלבי ביצוע מאוחדים כדי לבטל את התקורה של מימוש נתונים ביניים. בשלב משולב, הרכיבים מעובדים בלולאת ביצוע צפופה בשרשור יחיד לפני שהם מועברים לשלב הבא או לגבול של ערבוב.
פרטים נוספים על תרגום צינורות ויצירת גרפים זמינים במאמר בנושא מחזור החיים של צינורות.
מקביליות מנוהלת והתאמה אוטומטית לעומס (autoscaling)
כברירת מחדל, Dataflow מנהל את המקביליות של צינורות הנתונים באופן אוטומטי, בלי שצריך לבצע כוונון ידני של החלוקה למחיצות, בדרכים הבאות:
- שינוי גודל אוטומטי אופקי:
- צינורות להרצת משימות באצווה: המערכת מעריכה את העבודה הכוללת שנותרה, את העומס על המקור ואת השימוש במעבד כדי להגדיל או להקטין את מספר העובדים בבריכת העובדים, וכך להשלים את העבודה במהירות וביעילות מבחינת עלויות.
- צינורות להעברת נתונים בזמן אמת: ניתוח של זמן האחזור של המערכת, גודל ה-backlog וניצול המעבד כדי להגדיל את מספר העובדים בזמן שיאי התפוקה ולהקטין אותו בתקופות של תנועה נמוכה. פרטים נוספים זמינים במאמר בנושא התאמה של שינוי גודל אוטומטי אופקי של סטרימינג.
- איזון דינמי של עומסי עבודה (DWR): בצינורות של עיבוד באצווה, Dataflow עוקב אחרי ההתקדמות של משימות נפרדות של עובדים. אם worker מסיים את העבודה מוקדם או אם worker אחר מפגר בגלל חלוקת נתונים לא מאוזנת (stragglers), Dataflow מפצל באופן דינמי את העבודה הנותרת שלא עובדה מה-worker האיטי ומקצה אותה מחדש ל-worker שלא מבצע עבודה. מידע נוסף זמין במאמר בנושא איזון דינמי של עומסי עבודה.
- התאמה דינמית של מספר השרשורים: בצינורות (pipelines) של עיבוד באצווה שמשתמשים ב-Portable Runner, מספר השרשורים המקבילים לעיבוד בכל עובד מותאם אוטומטית על סמך ניצול המעבד והזיכרון. מידע נוסף זמין במאמר בנושא שינוי דינמי של מספר השרשורים.
- שינוי גודל אוטומטי אנכי: ב-Dataflow Prime, Dataflow משנה באופן דינמי את גודל הזיכרון של העובדים ואת משאבי המחשוב כדי למנוע שגיאות של חוסר זיכרון ולשפר את ניצול המשאבים. מידע נוסף מופיע במאמר בנושא שינוי גודל אוטומטי אנכי.
גורמים שמגבילים את ההרצה המקבילית
יכול להיות שצינור עיבוד נתונים לא ישיג את המקביליות הצפויה בגלל המאפיינים הבאים של הנתונים או העיצוב של תרשים צינור עיבוד הנתונים:
מקורות קלט שלא ניתן לפצל
אם אי אפשר לפצל מקור קלט לטווחים עצמאיים, מערכת Dataflow נאלצת לקרוא את המקור באופן רציף עם שרשור עובד יחיד:
- דחיסת קבצים שלא ניתן לפצל: אי אפשר לקרוא במקביל פורמטים כמו
.gz(gzip) או.bzip2(ללא יצירת אינדקס) מקיזוזים שרירותיים של בייטים. קריאה של קובץ דחוס גדול יחיד מגבילה את שלב ההטמעה לשרשור יחיד עד שהנתונים נפרסים ומחולקים מחדש. - פתרון: מאחסנים את הנתונים בפורמטים של קבצים שניתנים לפיצול (כמו Parquet, Avro או פורמטים דחוסים של Snappy) או מפצלים את נתוני הקלט לכמה קבצים קטנים יותר ב-Cloud Storage.
מיזוג של שלבים ו-fan-out רחב
מיזוג שלבים משפר את הביצועים על ידי הפחתת התקורה של הסריאליזציה, אבל הוא עלול להגביל בטעות את המקביליות ולהגביר את העומס על הזיכרון כששלב עם מקביליות נמוכה מייצר מספר גדול של רכיבי פלט (פעולה של 'fan-out' גבוה):
- דוגמה: מקור קורא חמישה קבצים ומתמזג עם טרנספורמציה
FlatMapשיוצרת 1,000,000 רכיבי פלט. אם הטרנספורמציהFlatMapמשולבת עם טרנספורמציות בהמשך הצינור, כל 1,000,000 הרכיבים ממשיכים לפעול על חמישה שרשורים של עובדים לכל היותר, מה שמגביל מאוד את התפוקה בהמשך הצינור. בנוסף, אם טרנספורמציות ביניים מתרחבות באופן משמעותי בזיכרון לפני ביצוע, חבילות גדולות עלולות לגרום לניצול מלא של הזיכרון הזמין של העובד. - רזולוציה: מוסיפים טרנספורמציה
Redistribute(או טרנספורמציה קלאסיתReshuffle) בין השלב של fan-out רחב היקף לבין טרנספורמציות בהמשך כדי לבטל את המיזוג ולחלק מחדש את העבודה בין העובדים. למידע על ניפוי באגים שקשורים לבעיות בזיכרון, אפשר לעיין במאמר פתרון שגיאות אין זיכרון פנוי (OOM).
הטיה של מקשים ומקשי קיצור
פעולות צבירה (GroupByKey, CoGroupByKey, Combine.PerKey) מקבצות רכיבים לפי המפתח המשויך שלהם.
- צוואר בקבוק של מפתח חם: Dataflow מעביר את כל הרכיבים עם אותו מפתח ל-Thread עובד יחיד לצורך צבירה. אם מפתח יחיד מכיל אחוז גבוה מכלל מערך הנתונים, ה-worker הזה הופך ל-straggler, ועובדים במעלה הזרם עלולים לחוות לחץ חוזר. לדוגמה, מפתח ברירת מחדל
nullאו מפתח קטגוריה פופולרי במיוחד. - רזולוציה:
- במקום
GroupByKey, כדאי להשתמש בפונקציות שילוב (CombineFnאוCombine.PerKey) כשזה אפשרי. כך Dataflow יכולה לבצע שילובים מקומיים חלקיים לפני ערבוב הנתונים. - מוסיפים קידומת או סיומת אקראית של מספר שלם למקשי קיצור (key salting) כדי לחלק את מרחב המקשים בין העובדים, ואז מבצעים צבירה בשלב השני כדי למזג את התוצאות המומלחות.
- במקום
ויסות נתונים (throttle) של יעד במורד הזרם
כשכותבים פלט של צינור עיבוד לשירותים חיצוניים כמו מסדי נתונים או ממשקי API של צד שלישי, רמת מקביליות גבוהה עלולה להביא לרוויה של מערכת היעד:
- הגבלת קצב העברת נתונים: מאות שרשורים של עובדים שמבצעים קריאות כתיבה בו-זמנית עלולים לגרום לשגיאות של הגבלת קצב, לפסק זמן בחיבור או לירידה בביצועים של מסד הנתונים.
- רזולוציה:
- כדי להגביל את המקביליות של פעולות הכתיבה, אפשר לקבץ רכיבים עם תפקיד
GroupByKeyאו להשתמש ב-batching sinks עם מקביליות מבוקרת. - צריך להטמיע השהיה מעריכית לפני ניסיון חוזר (exponential backoff) ולוגיקה של ניסיון חוזר בצד הלקוח בהטמעות של יעד
DoFn.
- כדי להגביל את המקביליות של פעולות הכתיבה, אפשר לקבץ רכיבים עם תפקיד
אסטרטגיות אופטימיזציה
כדי לבצע אופטימיזציה של מקביליות במשימות Dataflow, אפשר להיעזר בשיטות הבאות:
- כדי למנוע מיזוג לא רצוי עם
Redistribute:-
Redistribute.arbitrarily(): מפסיק את המיזוג של השלבים ומפיץ מחדש את הרכיבים באופן שווה בין כל העובדים הזמינים. -
Redistribute.byKey(): מאזן מחדש את זוגות המפתח/ערך בשרשורי העובדים תוך שמירה על המיקום של המפתחות. - דוגמאות להטמעה מופיעות במאמר בנושא מניעת מיזוג.
-
- מעקב אחרי חריגים וצווארי בקבוק: אפשר להשתמש בפרטי הביצוע של המסוף Google Cloud כדי לזהות שלבים עם מספר גבוה של חריגים או התקדמות תקועה:
המאמרים הבאים
- מידע על מחזור החיים של צינורות
- התאמה אופקית אוטומטית לעומס
- מידע על איזון דינמי של עומסי עבודה
- כדאי לעיין בשיטות המומלצות לשימוש בפייפליין ב-Dataflow.
- איך לפתור שגיאות שקשורות לזיכרון