Exactly-once ב-Dataflow

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

סקירה כללית

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

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

  • הרשומה מעובדת ולא אובדת.
  • כל תוצאות העיבוד שנשארות בצינור העיבוד משתקפות לכל היותר פעם אחת.

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

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

נתונים מאוחרים

עיבוד בדיוק פעם אחת מבטיח את הדיוק של צינור עיבוד הנתונים: אם צינור עיבוד הנתונים מעבד רשומה, Dataflow מוודא שהרשומה משתקפת בפלט ושלא נוצר עותק כפול של הרשומה.

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

באופן כללי, אין דרך להבטיח שלמות בצינור להעברת נתונים בזמן אמת, כי באופן תיאורטי רשומות יכולות להגיע באיחור. במקרה הקיצוני, תצטרכו לחכות לנצח כדי לקבל תוצאה. מבחינה מעשית, Apache Beam מאפשרת להגדיר את הסף להשמטת נתונים מאוחרים ואת המועד להפקת תוצאות מצטברות. מידע נוסף זמין במאמר Watermarks and late data (סימני מים ונתונים מאוחרים) במסמכי התיעוד של Apache Beam.

תופעות לוואי

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

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

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

רישום ביומן

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

סטרימינג של כל נתון בדיוק פעם אחת

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

ארגון נתונים בזמן אמת של Dataflow

משימות Dataflow להזרמת נתונים פועלות במקביל בהרבה עובדים שונים, על ידי הקצאת טווחי עבודה לכל עובד. למרות שההקצאות עשויות להשתנות לאורך זמן בתגובה לכשלים של עובדים, שינוי גודל אוטומטי או אירועים אחרים, אחרי כל GroupByKey טרנספורמציה, כל הרשומות עם אותו מפתח מעובדות באותו עובד. הטרנספורמציה GroupByKey משמשת לעיתים קרובות בטרנספורמציות מורכבות, כמו Count או FileIO. כדי להבטיח שרשומות של מפתח נתון יגיעו לאותו עובד, עובדי Dataflow מבצעים ערבוב של נתונים ביניהם באמצעות קריאות לפרוצדורות מרוחקות (RPC).

כדי לוודא שהרשומות לא יאבדו במהלך הערבוב, מערכת Dataflow משתמשת בגיבוי upstream. בגיבוי במעלה הזרם, תהליך העבודה ששולח את הרשומות מנסה שוב את קריאות ה-RPC עד שהוא מקבל אישור חיובי שהרשומה התקבלה. תופעות הלוואי של עיבוד הרשומה נשמרות לאחסון מתמיד בהמשך. אם העובד ששולח את הרשומות לא זמין, Dataflow ממשיך לנסות שוב את קריאות ה-RPC, וכך מוודא שכל רשומה מועברת לפחות פעם אחת.

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

איך מבטיחים זמן אחזור קצר

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

פילטרים של פריחה

מסנני Bloom הם מבני נתונים קומפקטיים שמאפשרים לבצע במהירות בדיקות של חברות בקבוצה. ב-Dataflow, כל עובד שומר מסנן בלום של כל מזהה שהוא רואה. כשמגיע מזהה רשומה חדש, העובד מחפש את המזהה במסנן. אם המסנן מחזיר false, הרשומה הזו לא כפולה, וה-worker לא מחפש את המזהה באחסון יציב.

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

איסוף אשפה

כדי להימנע ממצב שבו האחסון מתמלא במזהי רשומות, Dataflow משתמשת באיסוף כדי להסיר רשומות ישנות. ‫Dataflow משתמש בחותמת הזמן של המערכת כדי לחשב את סימן המים של איסוף האשפה.

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

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

מקורות לא דטרמיניסטיים

‫Dataflow משתמש ב-Apache Beam SDK כדי לקרוא נתונים לתוך צינורות עיבוד נתונים. אם העיבוד נכשל, יכול להיות ש-Dataflow ינסה שוב לקרוא ממקור. במצב כזה, Dataflow צריך לוודא שכל רשומה ייחודית שנוצרת על ידי מקור מתועדת בדיוק פעם אחת. במקורות דטרמיניסטיים, כמו Pub/Sub Lite או Kafka, הרשומות נקראות על סמך היסטוריית ההזזה שנרשמה, ולכן אין צורך לבצע את השלב הזה.

מכיוון ש-Dataflow לא יכול להקצות מזהי רשומות באופן אוטומטי, מקורות לא דטרמיניסטיים צריכים לציין ל-Dataflow מהם מזהי הרשומות כדי למנוע כפילויות. אם המקור מספק מזהים ייחודיים לכל רשומה, המחבר משתמש בערבוב בצינור כדי להסיר כפילויות. רשומות עם אותו מזהה מסוננות. דוגמה לאופן שבו Dataflow מטמיע עיבוד בדיוק פעם אחת כשמשתמשים ב-Pub/Sub כמקור מופיעה בקטע עיבוד בדיוק פעם אחת בדף בנושא סטרימינג עם Pub/Sub.

כשמבצעים DoFns מותאמים אישית כחלק מצינור הנתונים, Dataflow לא מבטיח שהקוד הזה יופעל רק פעם אחת לכל רשומה. כדי להבטיח עיבוד של לפחות פעם אחת במקרה של כשלים ב-worker, יכול להיות ש-Dataflow יעביר רשומה מסוימת דרך טרנספורמציה כמה פעמים, או שהוא יפעיל את אותה רשומה בו-זמנית בכמה workers. אם כוללים בקוד של צינור עיבוד הנתונים פעולות כמו יצירת קשר עם שירות חיצוני, יכול להיות שהפעולות יופעלו יותר מפעם אחת עבור רשומה מסוימת.

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

שליחת פלט בדיוק פעם אחת

ה-SDK של Apache Beam כולל יעדי נתונים מובנים שנועדו להבטיח שלא ייווצרו כפילויות. כשהדבר אפשרי, מומלץ להשתמש באחד ממאגרי הנתונים המובנים האלה.

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

לדוגמה, בצבירה בחלון, יכול להיות שקבוצת הרשומות בחלון לא תהיה דטרמיניסטית. באופן ספציפי, יכול להיות שהחלון ינסה להפעיל את הרכיבים e0, e1, e2. יכול להיות שה-worker יקרוס לפני שהוא יבצע את העיבוד של החלון, אבל לא לפני שהאלמנטים האלה יישלחו כתוצאת לוואי. כשמפעילים מחדש את ה-worker, החלון מופעל שוב, ופריט מאוחר e3 מגיע. מכיוון שהרכיב הזה מגיע לפני שהחלון מתבצע, הוא לא נספר כנתונים שהגיעו באיחור, ולכן הפונקציה DoFn נקראת שוב עם הרכיבים e0,‏ e1,‏ e2 ו-e3. לאחר מכן, הרכיבים האלה נשלחים לפעולת תופעת הלוואי. אידמפוטנטיות לא עוזרת בתרחיש הזה, כי בכל פעם נשלחות קבוצות שונות של רשומות לוגיות.

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

c.apply(Window.<..>into(FixedWindows.of(Duration.standardMinutes(1))))
 .apply(GroupByKey.<..>.create())
 .apply(new PrepareOutputData())
 .apply(Reshuffle.<..>of())
 .apply(WriteToSideEffect());

כדי לוודא שהרצת Dataflow תדע שהאלמנטים צריכים להיות יציבים לפני ההפעלה של DoFn, מוסיפים את ההערה RequiresStableInput ל-DoFn.

מידע נוסף