מאפייני הביצועים של צינורות עיבוד נתונים מ-Kafka ל-Iceberg

בדף הזה מתוארות מאפייני הביצועים, החל מגרסה 2.75.0 של Apache Beam, של משימות סטרימינג ב-Dataflow שקוראות מ-Apache Kafka וכותבות לטבלאות Apache Iceberg. הוא בודק את ההבדלים בביצועים בין כתיבה ישירה של Apache Iceberg לבין כתיבה שמועברת דרך Managed BigQuery API, ומשווה את התוצאות האלה לנקודות השוואה בסיסיות מצינורות נתונים מ-Kafka ל-BigQuery. האופטימיזציות של Apache Iceberg I/O עדיין מתבצעות, ולכן מדדי הביצועים האלה עשויים להשתנות.

השוואות לנקודות השוואה זמינות בשלושה הגדרות מיפוי חסרות מצב עיקריות (כלומר, הן קוראות מהמקור, ממירות את ההודעה לרשומה וכותבות ל-sink בלי לעקוב אחרי מצב או להחיל לוגיקה עסקית מורכבת; נקראות map_only או mapping בנקודות השוואה):

  1. ‫Kafka ל-BigQuery (map_only) (נתוני בסיס מתוך ביצועים של Kafka ל-BigQuery)
  2. ‫Kafka to Iceberg Direct ‏ (map_only, ‏ autosharding=false)
  3. ‫Kafka to Iceberg using Managed BigQuery API (map_only)

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

מתודולוגיית הבדיקה

ההשוואות בוצעו באמצעות המשאבים הבאים:

  • אשכול של שירות מנוהל ל-Apache Kafka: התנועה נוצרה באמצעות תבנית מחולל הנתונים של Dataflow Streaming.
    • קצב העברת נתונים של הקלט: 1GBps
    • קצב שליחת ההודעות: כ-1,000,000 הודעות בשנייה
    • פורמט ההודעה: טקסט JSON עם סכימה קבועה (כ-1KB לכל הודעה)
    • מחיצות: 1,000 מחיצות של Kafka
  • יעדים של מאגרי מידע:
    • ‫BigQuery: טבלה רגילה (לא מחולקת למחיצות) שנכתבה באמצעות BigQuery Storage Write API.
    • ‫Apache Iceberg: קטלוג שמגובה על ידי Cloud Storage. הטבלה Direct sink מחולקת למחיצות באמצעות bucket(id, 64) (חלוקה ל-64 רסיסים במפתח הראשי) ומשתמשת במצב ההפצה hash.

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

ביצועי הטמעת נתונים: מיפוי עומסי עבודה

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

הגדרת משרה

הגדרה ‫Kafka ל-BigQuery (map_only) ‫Kafka to Iceberg Direct ‏ (autosharding=false) ‫Kafka ל-Iceberg באמצעות Managed BigQuery API
סוג מכונת Worker e2-standard-2 e2-standard-4 e2-standard-4
מעבדי vCPU לכל Worker 2 4 4
RAM per Worker 8 GB ‫16GB ‫16GB
‫Streaming Engine מופעל מופעל מופעל
התאמה אופקית לעומס (auto-scaling) מופעל מופעל מופעל
תדירות ההפעלה ‫5 שניות ‫60 שניות ‫60 שניות

תפוקה ושימוש במשאבים

כתיבה ישירות לקובצי Parquet פיזיים במאגר אובייקטים כרוכה בתקורה גבוהה יותר של קלט/פלט מאשר הטמעת עדכונים בזמן אמת ב-BigQuery. בהשוואה לכתיבה ישירה ב-Iceberg, ניתוב כתיבה דרך Managed BigQuery API משפר את ניצול ה-CPU של העובדים (כ-70% לעומת כ-60%) ומפחית באופן מתון את הצריכה של Streaming Engine (כ-180 SECU/h לעומת כ-200 SECU/h), אם כי דרישות המחשוב הכוללות של העובדים נשארות דומות (כ-440 vCPUs לעומת כ-450 vCPUs).

מדד ‫Kafka ל-BigQuery (map_only) ‫Kafka to Iceberg Direct ‏ (autosharding=false) ‫Kafka ל-Iceberg באמצעות Managed BigQuery API
קצב העברת נתונים ממוצע של קלט לכל Worker ‫~15 MBps ‫‎~9 MBps ‫‎~9 MBps
ממוצע ניצול המעבד ~70% ‎~60% ~70%
הערכה של מספר ליבות ה-vCPU לנתוני קלט של ‎1 GBps ‎~126 מעבדי vCPU ‫~450 מעבדים וירטואליים ‫~440 מעבדים וירטואליים
מספר העובדים המשוער לנתוני קלט של 1GBps כ-63 עובדים כ-110 עובדים כ-110 עובדים
הערכה של עלות השימוש בשירותי הענן לפי שעה למהירות של 1GBps כ-58 יחידות SECU לשעה ‫~200 SECU/h ‫~180 SECU/h

כתיבה של פרופיל זמן ההמתנה

לפעולות כתיבה ישירות ב-Iceberg יש חביון זנב חמור (P99) בגלל מגבלות של ביצוע שינויים במטא-נתונים של אחסון אובייקטים. השימוש ב-Managed BigQuery API מבטל את העליות החדות בחביון הזנב, תוך שמירה על חביון חציוני נמוך.

זמן האחזור של כתיבה מקצה לקצה ‫Kafka ל-BigQuery ‫Kafka to Iceberg Direct ‏ (autosharding=false) ‫Kafka ל-Iceberg באמצעות Managed BigQuery API
P50 (Median) ‫~1,200 אלפיות השנייה ‫~1,000 אלפיות השנייה ‎~1,000 ms
P95 ‫~3,000 אלפיות השנייה ‫~7,400 אלפיות השנייה ‎~1,900 ms
P99 (Tail) ‫~5,400 אלפיות השנייה ‫~14,000 אלפיות השנייה כ-2,700 אלפיות השנייה

שיקולים לגבי חלוקה אוטומטית של נתונים ושיקולים לגבי עיצוב

בקטע הזה נדון בהשלכות של חלוקה אוטומטית של נתונים על גדלי הקבצים ועל זמן האחזור של צינורות הנתונים כשכותבים ל-Apache Iceberg.

למה נבחרה autosharding=false כנקודת ההשוואה

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

כדי לשמור על פריסות יציבות וצפויות של קובצי Parquet (בגודל של ‎~800 KB בממוצע) ולהבטיח בסיס השוואה הוגן ללא ניקוי מוקדם של הנתונים, נבחרה האפשרות autosharding=false להגדרת יעד ישיר.

מה קורה אם משביתים את הפיצול האוטומטי לעומת השארתו מופעלת?

  • עם autosharding=false (Baseline): גודלי הקבצים הראשוניים גדולים יותר (בממוצע 800KB) בהשוואה לחלוקה אוטומטית. הגודל הזה עדיין קטן בהשוואה לגדלים האידיאליים של קובצי Iceberg ‏ (128-512MB), אבל הוא דורש דחיסה משמעותית פחות בהמשך. עם זאת, יש לכך מחיר: חביון גבוה של זנב הכתיבה (P99 מגיע ל-14.0 שניות בערך) בגלל צווארי בקבוק במטא-נתונים של אחסון האובייקטים.
  • אם מופעלת חלוקה אוטומטית של נתונים: מערכת Dataflow משנה את קנה המידה של שרשורי הכתיבה באופן דינמי כדי לספוג שיאים מקומיים של קצב העברת נתונים, וכך מקטינה את חביון הזנב של הכתיבה. עם זאת, הוא פוגע בשכבת האחסון כי הוא יוצר נפח גדול של קובצי Parquet קטנים ומפוצלים (בגודל של ‎~100 KB או פחות). גודלי הקבצים האלה משתנים מאוד, והם נעים באופן שרירותי בין הרצות (בממוצע מ-39KB בערך עד 100KB בערך), ולכן יש צורך בתחזוקה אגרסיבית של דחיסה במורד הזרם.

התאמה של מחיצות והמלצות

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

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

השלכות על קריאה במורד הזרם: גודל הקובץ ודחיסה

למרות שהמדדים בצד הכתיבה מעדיפים את Managed BigQuery API להטמעת Iceberg, היעילות הכוללת של צינור הנתונים תלויה מאוד בביצועים של קריאת הנתונים בהמשך:

  • יצירת קבצים קטנים ב-Managed BigQuery API: ממשק Managed BigQuery API מבצע לעיתים קרובות פעולות של העברת נתונים מהזיכרון לדיסק כדי להבטיח זמן אחזור נמוך של כתיבה. ההתנהגות הזו מובילה לכתיבה של נפח גבוה של קובצי Parquet קטנים בקטלוג היעד של Iceberg.
  • מידע על ההשפעה של זמן האחזור של שאילתות: מנועי שאילתות (לדוגמה, Starburst/Trino, ‏ Apache Spark, ‏ BigQuery, ‏ Dremio) שקוראים טבלאות עם מיליוני קובצי Parquet קטנים גורמים לעומס כבד על ניתוח מטא-נתונים ולעונשים על סריקת מחיצות.
  • דרישות דחיסה: כדי למנוע ירידה בביצועים של קריאה כשמשתמשים בממשק BigQuery API המנוהל (או אם מופעל פיצול אוטומטי של נתונים בכתיבה ישירה), צריך להריץ באופן קבוע משימות תחזוקה של דחיסת Iceberg (לדוגמה, REWRITE DATA FILES). צריך לקחת בחשבון את התקורה של החישובים לצורך דחיסה כשמתכננים את הארכיטקטורה הכוללת.
  • הפצה של קובץ Direct Write‏ (autosharding=false): כתיבות ישירות של Iceberg עם חלוקה קבועה לשברים יוצרות קובצי Parquet גדולים יותר בממוצע (כ-‎800 KB), וכך מתקבל פריסה פחות מקוטעת לגישה מיידית לשאילתות ללא דרישות דחיסה מיידיות (אבל עדיין מתחת לטווח האידיאלי).

צינורות Iceberg ישירים עם שמירת מצב (groupbykey)

כדי להעריך אסטרטגיות ידניות של אצווה, נערך ניסוי של קיבוץ מפתחות עם שמירת מצב (groupbykey) בהשוואה לצינור הבסיסי Kafka to Iceberg Direct (map_only, autosharding=false). שתי התצורות כותבות קובצי Parquet ישירות לאחסון אובייקטים.

השוואה למחיר המקובל בשוק

מדד / תכונה ערך בסיס של כיור ישיר (autosharding=false) Stateful Direct Sink (groupbykey) השפעה על הביצועים
מספר משוער של vCPU ל-1GBps ‎~450 מעבדי vCPU ‫~520 מעבדים וירטואליים ‫~+16% מחשוב נדרש
ממוצע ניצול המעבד ~60% ‎~50% ~ירידה של 17% ביעילות העובדים
עלות משוערת של SECU לשעה ל-1GBps ~200 SECU/h כ-300 SECU/h ‫~+50% עומס של מנוע סטרימינג
גודל קובץ ממוצע ‫~800 KB ‫~100KB יצירת קבוצות קטנות יותר של קבצים
זמן האחזור של P50 ‎~1,000 ms ‫~1,200 אלפיות השנייה ‫~+20% חציון איטי יותר
זמן האחזור P95 ‫~7,400 אלפיות השנייה ‫~5,500 אלפיות השנייה זמן אחזור נמוך יותר בכ-26%
זמן אחזור P99 ‫~14,000 אלפיות השנייה ‫~13,000 אלפיות השנייה שינוי שולי בזמן האחזור של הזנב

ניתוח פשרות

  1. תקורה של Streaming Engine: הוספה של שלב groupbykey עם שמירת מצב מחייבת את Beam לשמור מצב ביניים בין גבולות החלונות. הגדלת קצב הרענון מגדילה את צריכת יחידות המחשוב של Streaming Engine בשיעור של כ-50% (מ-‎~200 SECU/h ל-‎~300 SECU/h).
  2. זמן האחזור של אגירת הנתונים: צבירת מפתחות ידנית מחייבת אגירת נתונים בחלון, מה שמגדיל את זמן האחזור הממוצע של הכתיבה (P50) ל-‎~1,200 ms ואת זמן האחזור של P95 ל-‎~5.5s.

צינורות הפוכים: סטרימינג מ-Iceberg ל-Kafka

כדי להעריך את היכולות של lakehouse דו-כיווני, נערכו גם בדיקות השוואה לנתוני סטרימינג שזורמים בכיוון ההפוך – קריאת סטרימינג של נתונים שנוספו בלבד מטבלת Apache Iceberg ופרסום בחזרה ב-Apache Kafka.

הגדרת משימות ויעילות

בניגוד לצינורות להעברת נתונים שצריכים להתמודד עם כתיבות כבדות של קבצים בחנות אובייקטים או עם צווארי בקבוק של אישור מטא-נתונים, קריאה של שינויים מ-Iceberg והזרמתם מתבצעות ביעילות גבוהה:

מדד ‫Iceberg ל-Kafka (הוספה בלבד, בדיוק פעם אחת)
סוג מכונת Worker e2-standard-4
הערכה של מספר ליבות ה-vCPU לנתוני קלט של 1GBps ‫~30 מעבדים וירטואליים
מספר העובדים המשוער לנתוני קלט של 1GBps כ-7 עובדים
הערכה של עלות השימוש בשירותי הענן לפי שעה למהירות של 1GBps ‫~0.2 יחידות SECU לשעה

מסקנות עיקריות לגבי צינורות הפוכים

  • תקורה נמוכה משמעותית של מחשוב: קריאה של זרמי CDC מ-Iceberg והקרנה שלהם דורשת הרבה פחות משאבי מחשוב (כ-30 vCPU לעומת כ-450 vCPU לכתיבות ישירות), כי היא חוסכת את העבודה המאומצת של חלוקה למחיצות, קידוד ומסירה של נפחים גדולים של קובצי Parquet לאחסון אובייקטים.
  • יעילות משאבים: צריכה או שכפול מבוססי-אירועים במורד הזרם מפורמט lakehouse בחזרה לשכבות סטרימינג הם יעילים מאוד בהשוואה לנתיבי הטמעה נכנסים.

סיכום המלצות אדריכליות

דפוס ארכיטקטוני זמן האחזור של כתיבת נתונים ב-P99 פריסת הקובץ שיקולים לגבי קריאה במורד הזרם
‫Kafka to BigQuery (map_only) ‫~5.4 שניות לא רלוונטי אופטימלי (מנוע אחסון מנוהל של BigQuery)
‫Kafka to Iceberg using Managed BigQuery API ‫~2.7s קבצים קטנים באופן שרירותי נדרשת דחיסה תקופתית לקריאות בנפח גבוה
‫Kafka to Iceberg Direct ‏ (autosharding=false) ‫~14.0 שניות ‫~800 KB טוב (גודלי קבצים ראשוניים גדולים יותר, דרישה נמוכה יותר לדחיסה)
‫Kafka to Iceberg Direct ‏ (groupbykey) ‫~13.0 שניות ‫~100KB בינונית (עלות תקורה גבוהה יותר של מחשוב ומצב)

הערכת עלויות

כדי לאמוד את עלות הבסיס של צינור השוואה משלכם עם חיוב לפי משאבים, אתם יכולים להשתמש ב Google Cloud מחשבון התמחור באופן הבא:

  1. פותחים את מחשבון העלויות.
  2. לוחצים על הוספה לאומדן.
  3. בוחרים באפשרות Dataflow.
  4. בקטע סוג השירות, בוחרים באפשרות Dataflow Classic.
  5. לוחצים על הגדרות מתקדמות כדי לראות את כל האפשרויות.
  6. בוחרים את המיקום שבו העבודה תפעל.
  7. בקטע סוג העבודה, בוחרים באפשרות 'סטרימינג'.
  8. בוחרים באפשרות הפעלת מנוע הסטרימינג.
  9. מזינים את המידע על שעות ההפעלה של המשימה, צמתי העובדים, מכונות העובדים ואחסון Persistent Disk.
  10. מזינים את המספר המשוער של יחידות החישוב של Streaming Engine.

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

לדוגמה, נניח שאתם מריצים צינור עיבוד נתונים באמצעות הארכיטקטורה Kafka to Iceberg Direct (autosharding=false), עם קצב נתוני קלט של 100MBps. על סמך תוצאות ההשוואה לצינור עיבוד נתונים של 1GBps, אפשר לאמוד את דרישות המשאבים באופן הבא:

  • מקדם קנה מידה: (‎100 MBps) / (‎1,024 MBps) = ‎~0.1
  • צמתים משוערים של עובדים: 110 עובדים × 0.1 = ‎~11 עובדים
  • המספר המשוער של יחידות חישוב של Streaming Engine לשעה: ‎200 × 0.1 = ~20 units per hour

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

הרצת צינור לבדיקה

כדי לפרוס משימת סטרימינג של Apache Iceberg באמצעות תבנית Flex של Dataflow, משתמשים בפקודה gcloud dataflow flex-template run.

gcloud dataflow flex-template run JOB_NAME \
  --project=PROJECT_ID \
  --region=REGION \
  --template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_To_Iceberg_Yaml \
  --enable-streaming-engine \
  --parameters ^@^bootstrapServers="KAFKA_BOOTSTRAP_ADDRESS"\
@topic="KAFKA_TOPIC"\
@table="ICEBERG_TABLE_IDENTIFIER"\
@catalogName="CATALOG_NAME"\
@catalogProperties='{"type":"CATALOG_TYPE","warehouse":"gs://BUCKET_NAME/warehouse/"}'\
@triggeringFrequencySeconds=60\
@schema='SCHEMA_DEFINITION'

מחליפים את מה שכתוב בשדות הבאים:

  • ‫JOB_NAME: השם של משימת Dataflow
  • ‫PROJECT_ID: מזהה הפרויקט ב- Google Cloud
  • ‫REGION: האזור Google Cloud שבו מופעלת העבודה (לדוגמה, us-central1)
  • ‫KAFKA_BOOTSTRAP_ADDRESS: כתובת ה-bootstrap של אשכול Apache Kafka
  • ‫KAFKA_TOPIC: השם של נושא Kafka
  • ‫ICEBERG_TABLE_IDENTIFIER: המזהה של טבלת היעד של Iceberg
  • ‫CATALOG_NAME: השם של קטלוג Iceberg
  • ‫CATALOG_TYPE: סוג הקטלוג לשימוש (לדוגמה, hadoop או bigquery)
  • ‫BUCKET_NAME: השם של הקטגוריה ב-Cloud Storage למיקום המחסן
  • ‫SCHEMA_DEFINITION: הגדרת הסכימה של נתוני הנושא ב-Kafka (לדוגמה, {"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]})