כדי לקרוא אירועים של לכידת נתונים משתנים (CDC) מ-Apache Iceberg באמצעות Lakehouse for Apache Iceberg REST Catalog, צריך להשתמש במחבר מנוהל של Apache Beam I/O.
שירות Managed I/O תומך ביכולות הבאות של Apache Iceberg:
| קטלוגים |
|
|---|---|
| יכולות קריאה | קריאה של נתונים באצווה |
| יכולות כתיבה |
|
כדי להשתמש בטבלאות BigQuery ל-Apache Iceberg, צריך להשתמש במחבר BigQueryIO עם BigQuery Storage API. הטבלה צריכה כבר להיות קיימת. אי אפשר ליצור טבלה דינמית.
מגבלות
- יש תמיכה ב-Apache Iceberg CDC רק באמצעות Managed API. התכונות של שירות השינויים המנוהלים עדיין לא הופעלו. צפויים שינויים שישפיעו על תאימות לדורות קודמים
- CDC Managed API קורא רק תמונות מצב של הוספה בלבד. עדיין לא זמין.
דרישות מוקדמות
- הגדרת Lakehouse for Apache Iceberg מגדירים את ההרשאות הנדרשות בפרויקט Google Cloud באמצעות ההוראות במאמר שימוש בקטלוג של זמן הריצה של Lakehouse עם קטלוג REST של Iceberg. חשוב לוודא שאתם מבינים את המגבלות של Lakehouse for Apache Iceberg Iceberg REST Catalog שמתוארות בדף הזה.
- יוצרים טבלת Iceberg של מקור. בדוגמה שמוצגת כאן מניחים שיש לכם טבלת Apache Iceberg. כדי ליצור צינור כזה, אפשר להשתמש בצינור שמוצג במאמר כתיבת נתונים בסטרימינג ל-Apache Iceberg באמצעות קטלוג REST של Lakehouse for Apache Iceberg.
תלויות
מוסיפים את יחסי התלות הבאים לפרויקט:
Java
<dependency>
<groupId>org.apache.beam</groupId>
<artifactId>beam-sdks-java-managed</artifactId>
<version>${beam.version}</version>
</dependency>
<dependency>
<groupId>org.apache.beam</groupId>
<artifactId>beam-sdks-java-io-iceberg</artifactId>
<version>${beam.version}</version>
</dependency>
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-gcp</artifactId>
<version>${iceberg.version}</version>
</dependency>
דוגמה
בדוגמה הבאה מוצג צינור עיבוד נתונים של סטרימינג שקורא אירועי CDC מטבלה של Apache Iceberg, מצבר קליקים של משתמשים וכותב את התוצאות לטבלה אחרת של Apache Iceberg.
Java
כדי לבצע אימות ב-Dataflow, צריך להגדיר את Application Default Credentials. מידע נוסף זמין במאמר הגדרת אימות לסביבת פיתוח מקומית.