כדי לקרוא אירועים של לכידת נתונים משתנים (CDC) מ-Apache Iceberg באמצעות קטלוג REST של Lakehouse, צריך להשתמש במחבר I/O מנוהל של Apache Beam.
שירות Managed I/O תומך ביכולות הבאות של Apache Iceberg:
| קטלוגים |
|
|---|---|
| יכולות קריאה | קריאה של נתונים באצווה |
| יכולות כתיבה |
|
כדי להשתמש בטבלאות BigQuery ל-Apache Iceberg, צריך להשתמש במחבר BigQueryIO עם BigQuery Storage API. הטבלה צריכה כבר להיות קיימת. יצירת טבלה דינמית לא אפשרית.
מגבלות
- יש תמיכה ב-Apache Iceberg CDC רק באמצעות Managed API. התכונות של שירות השינויים המנוהלים עדיין לא מופעלות. צפויים שינויים שישפיעו על תאימות לדורות קודמים
- CDC Managed API קורא רק תמונות מצב שניתן להוסיף להן נתונים. עדיין אין אפשרות להשתמש ב-CDC מלא.
דרישות מוקדמות
- הגדרת Lakehouse מגדירים את הפרויקט Google Cloud עם ההרשאות הנדרשות לפי ההוראות במאמר שימוש בקטלוג של זמן הריצה של Lakehouse עם קטלוג REST של Iceberg. חשוב לוודא שאתם מבינים את המגבלות של Lakehouse Iceberg REST Catalog שמתוארות בדף הזה.
- יוצרים טבלת Iceberg של מקור. בדוגמה שמוצגת כאן מניחים שיש לכם טבלת Apache Iceberg. כדי ליצור צינור כזה, אפשר להשתמש בצינור שמוצג במאמר Streaming Write to Apache Iceberg with Lakehouse REST Catalog.
תלויות
מוסיפים את יחסי התלות הבאים לפרויקט:
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. מידע נוסף זמין במאמר הגדרת אימות לסביבת פיתוח מקומית.