קלט/פלט מנוהל של Dataflow ל-Apache Iceberg

Managed I/O תומך ביכולות הבאות של Apache Iceberg:

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

כדי להשתמש בטבלאות BigQuery ל-Apache Iceberg, צריך להשתמש במחבר BigQueryIO עם BigQuery Storage API. הטבלה צריכה כבר להיות קיימת. יצירת טבלה דינמית לא אפשרית.

דרישות

ערכות ה-SDK הבאות תומכות ב-I/O מנוהל ל-Apache Iceberg:

  • ‫Apache Beam SDK for Java בגרסה 2.58.0 ואילך
  • ‫Apache Beam SDK ל-Python בגרסה 2.61.0 ואילך

הגדרות אישיות

‫Managed I/O for Apache Iceberg תומך בפרמטרים הבאים של ההגדרה:

ICEBERG קריאה

הגדרות אישיות סוג תיאור
table str מזהה של טבלת Iceberg.
catalog_name str השם של הקטלוג שמכיל את הטבלה.
catalog_properties map[str, str] מאפיינים שמשמשים להגדרת קטלוג Iceberg.
config_properties map[str, str] מאפיינים שמועברים להגדרות של Hadoop.
ירידה list[str] קבוצת משנה של שמות עמודות שרוצים להחריג מקריאה. אם הערך הוא null או ריק, כל העמודות ייקראו.
מסנ str פרדיקט דמוי SQL לסינון נתונים בזמן הסריקה. דוגמה: "id > 5 AND status = 'ACTIVE'". משתמש בתחביר של Apache Calcite: https://calcite.apache.org/docs/reference.html
keep list[str] קבוצת משנה של שמות עמודות לקריאה בלבד. אם הערך הוא null או ריק, כל העמודות ייקראו.

ICEBERG כתיבה

').

הגדרות אישיות סוג תיאור
table str מזהה טבלה מוגדר במלואו. אפשר גם לספק תבנית לכתיבה למספר יעדים דינמיים, לדוגמה: ‎`dataset.my_{col1}_{col2.nested}_table`‎.
חלוקה אוטומטית של נתונים boolean ההגדרה הזו מאפשרת חלוקה דינמית למקטעים, כדי להתאים באופן אוטומטי את מספר הכותבים המקבילים על סמך נפח הנתונים. הוא מטפל בחלוקת נתונים לא מאוזנת (partition skew) על ידי חלוקה נוספת של המחיצות לכמה מחיצות Shard כדי למנוע צווארי בקבוק במהלך כתיבה של נתונים בהעברה גבוהה. האפשרות הזו זמינה רק במצב ההפצה 'hash'.
catalog_name str השם של הקטלוג שמכיל את הטבלה.
catalog_properties map[str, str] מאפיינים שמשמשים להגדרת קטלוג Iceberg.
config_properties map[str, str] מאפיינים שמועברים להגדרות של Hadoop.
direct_write_byte_limit int32 בצינור עיבוד נתונים (pipeline) של סטרימינג, המגבלה הזו קובעת את מספר החבילות שאפשר להעביר לנתיב הכתיבה הישיר.
distribution_mode str הגדרת ההתפלגות של נתוני הכתיבה. הפצות נתמכות: - none: לא מתבצע ערבוב של השורות (ברירת מחדל) - hash: השורות מעורבבות לפי מפתח המחיצה לפני כתיבת הנתונים
ירידה list[str] רשימה של שמות שדות להסרה מרשומת הקלט לפני הכתיבה. המאפיין הזה מוציא מכלל אפשרות את המאפיינים keep ו-only.
keep list[str] רשימה של שמות שדות שרוצים לשמור ברשומת הקלט. כל שאר השדות מושמטים לפני הכתיבה. הערך הזה לא יכול להיות זהה לערכים drop או only.
רק str השם של שדה רשומה יחיד שצריך לכתוב. הוא לא יכול להיות ביחד עם 'שמירה' ו 'הסרה'.
partition_fields list[str] שדות שמשמשים ליצירת מפרט מחיצה שמוחל כשיוצרים טבלאות. עבור השדה foo, הטרנספורמציות הזמינות של המחיצות הן:
  • foo
  • truncate(foo, N)
  • bucket(foo, N)
  • hour(foo)
  • day(foo)
  • month(foo)
  • year(foo)
  • void(foo)

מידע נוסף על טרנספורמציות של מחיצות זמין בכתובת https://iceberg.apache.org/spec/#partition-transforms.

sort_fields list[str] השדות שמשמשים להגדרת סדר המיון של הטבלה, שמוחל כשהטבלה נוצרת. כל רשומה היא מהצורה <term> [asc|desc] [nulls first|nulls last], כאשר <term> הוא שם שדה או אחת מהטרנספורמציות של החלוקה למחיצות (למשל bucket(col, 4),‏ day(ts)). כיוון ברירת המחדל הוא בסדר עולה, וסדר ברירת המחדל של ערכי null הוא nulls-first בסדר עולה ו-nulls-last בסדר יורד. הערה: הפעולה הזו מגדירה את סדר המיון המוצהר של הטבלה כמטא-נתונים, אבל היא לא גורמת ל-Beam למיין פיזית את הרשומות לפני הכתיבה. מידע נוסף על סדרי מיון זמין בכתובת https://iceberg.apache.org/spec/#sort-orders.
table_properties map[str, str] מאפייני טבלה של Iceberg שיוגדרו בטבלה כשהיא תיווצר. מידע נוסף על מאפייני טבלה זמין בכתובת https://iceberg.apache.org/docs/latest/configuration/#table-properties.
triggering_frequency_seconds int32 בצינור עיבוד זרמי נתונים (stream processing), ההגדרה הזו קובעת את התדירות שבה נוצרות תמונות מצב.
write_properties map[str, str] מאפיינים שמוחלים על כלי הכתיבה של הקובץ הבסיסי (למשל, מאפייני כתיבה של Parquet כמו 'write.parquet.bloom-filter-enabled.column.

המאמרים הבאים

מידע נוסף ודוגמאות קוד זמינים בנושאים הבאים: