יכול להיות שיהיה קשה לשמור על סנכרון בין טבלאות היעד של Lakehouse Apache Iceberg לבין טבלאות המקור של Delta Lake שמתעדכנות לעיתים קרובות. ביצוע טעינות מחדש של טבלאות מלאות כדי לתעד שינויים במערך הנתונים מגדיל את עלויות המחשוב ואת זמן האחזור של העיבוד, ויצירת צינורות נתונים בהתאמה אישית לעיבוד יומני עסקאות יוצרת תקורה תפעולית מורכבת.
באמצעות צינורות עיבוד נתונים של Dataflow ללכידת נתונים (CDC) של שינויים במנות שנקבעו מראש, אתם יכולים לסנכרן באופן רציף את טבלאות המקור של Delta Lake עם טבלאות Lakehouse Apache Iceberg. במקום לטעון מחדש טבלאות מלאות, משימות אצווה מתוזמנות מחלצות עדכונים מצטברים ישירות מיומני הטרנזקציות של Delta Lake ומחילות אותם על טבלאות היעד בלוח זמנים שניתן להגדרה.
צינור ה-CDC של אצווה מתוזמנת מספק את היכולות הבאות:
- סנכרון מצטבר חסכוני: המערכת מעבדת מעת לעת רק את הנתונים שהשתנו (פעולות
INSERT,UPDATEו-DELETE), וכך מצמצמת את התקורה של המחשוב ואת זמן האחזור של הנתונים. - מיפוי שינויים אוטומטי: ניתוח של יומני Change Data Feed של Delta Lake ומיפוי של פעולות שינוי לטבלת Iceberg של Lakehouse היעד.
- תזמון גמיש של משימות: הפעלה של משימות Dataflow במצב אצווה לפי לוח זמנים (למשל, במרווחי זמן של 15 דקות) שמתאים לצרכים העסקיים שלכם.
תהליך עבודה של העברה רציפה
כדי להגדיר העברה רציפה מקצה לקצה מ-Delta Lake ל-Lakehouse, צריך לבצע את רצף הפעולות הבא:
- מפעילים את המאפיין
delta.enableChangeDataFeed(true) בטבלת הקלט של Delta Lake. - אחרי שמפעילים את המאפיין
delta.enableChangeDataFeed, מבצעים העברה ראשונית של הטבלה למזהה commit או לחותמת זמן. הוראות מפורטות זמינות במאמר ייבוא טבלאות Delta Lake ל-Lakehouse באמצעות Dataflow. - מריצים משימות מתוזמנות של CDC בקבוצות כדי לעדכן את הנתונים בהמשך בתדירות (מזהה commit או טווח חותמות זמן) שנבחרה בהתאם לדרישות עומס העבודה.
לפני שמתחילים
כדי להגדיר העברה מתוזמנת של נתוני CDC במקבצים, צריך לוודא שיש לכם:
מפעילים את ממשקי ה-API של Dataflow, BigQuery ו-Lakehouse, אם הם עדיין לא מופעלים.
תפקידים שנדרשים להפעלת ממשקי API
כדי להפעיל ממשקי API, צריך את ההרשאה
serviceusage.services.enable. אם יצרתם את הפרויקט, סביר להניח שכבר יש לכם את ההרשאה הזו דרך התפקיד 'בעלים' (roles/owner). אחרת, תוכלו לקבל את ההרשאה הזו דרך התפקיד 'אדמין בממשק Service Usage' (roles/serviceusage.serviceUsageAdmin). איך מקצים תפקידיםכדי לקבל את ההרשאות שדרושות ליצירת המשאבים, צריך לבקש מהאדמין להקצות לכם את התפקידים הנדרשים בניהול הזהויות והרשאות הגישה (IAM) בפרויקט.
טבלת Delta Lake תקינה שמאוחסנת בקטגוריה של Cloud Storage. ספריית הטבלה צריכה להכיל את קובצי הנתונים ואת ספריית יומן העסקאות
_delta_log/. בנוסף, הטבלה צריכה לעמוד בדרישות הבאות:- צריך להפעיל את המאפיין
delta.enableChangeDataFeed(true) בטבלת Delta Lake. - מזהה השמירה או חותמת הזמן של ההתחלה שהוגדרו לצינור ה-CDC צריכים להיות אחרי שהנכס
delta.enableChangeDataFeedהופעל.
- צריך להפעיל את המאפיין
קטלוג קיים של Iceberg Lakehouse שאליו יסונכרנו הנתונים. אם מרחב השמות של היעד לא קיים, צינור ה-CDC יוצר אותו באופן אוטומטי. אם טבלת Lakehouse של היעד לא קיימת, צינור ה-CDC יוצר אותה באופן אוטומטי אם מספקים רשימה תקינה של עמודות שוויון.
תמיכה ומגבלות
השיקולים הבאים רלוונטיים להעברת CDC באצווה מתוזמנת מ-Delta Lake אל Lakehouse:
- דרישות SDK: נדרשת גרסת Apache Beam SDK 2.77.0 ואילך.
- יצירת טבלת יעד: יש תמיכה ביצירת טבלת יעד ב-Lakehouse. אם טבלת היעד לא קיימת, צינור ה-CDC יוצר אותה באופן אוטומטי אם מספקים רשימה תקינה של עמודות שוויון (
equality_columns). נדרש רק קטלוג Lakehouse קיים. צינור ה-CDC יוצר את מרחב השמות באופן אוטומטי אם הוא לא קיים. - התפתחות סכימה: לא ניתן לבצע שינויים אוטומטיים בסכימה. אם הסכימה של טבלת המקור Delta Lake משתנה, צריך לעדכן באופן ידני את סכימת טבלת היעד Lakehouse כך שתתאים לה לפני שמריצים את משימת הסנכרון של CDC.
- מצב הפעלה של אצווה: תהליך ה-CDC של Delta Lake פועל במצב אצווה.
הפייפליין שולף יומני טרנזקציות פעם אחת בכל הפעלה של משימה בין גבולות התחלה וסיום מוגדרים (מזהה קומיט או חותמת זמן). אם לא מציינים מחויבות לסיום (
end_versionאוend_timestamp), צינור הנתונים קורא עד לגרסה האחרונה של המחויבות. - אחסון מקור נתמך: נתוני המקור צריכים להיות טבלה תקינה של Delta Lake שמאוחסנת ב-Cloud Storage (
gs://). לא ניתן להשתמש ב-Amazon S3 כמקור לטבלאות Delta Lake. - טבלאות שמבוססות על קטלוג: אין תמיכה בטבלאות Delta Lake שמבוססות על קטלוג (לדוגמה, Unity Catalog).
- גרסאות Delta Lake שלא נתמכות: גרסה 1.2.1 או גרסאות קודמות של Delta Lake לא נתמכות כי לא קיימת תמיכה בפיד נתוני שינויים (CDF) בגרסאות האלה.
- פיד אוטומטי של נתוני שינויים: אין תמיכה בפיד אוטומטי של נתוני שינויים כי הוא לא כותב קובצי יומן שינויים של CDF.
יצירת צינור עיבוד נתונים של CDC באופן פרוגרמטי
כדי ליצור ולהפעיל באופן פרוגרמטיבי צינור עיבוד נתונים להעברת נתונים ברצף (batch) של Delta Lake CDC, משתמשים בטרנספורמציה Managed.read בצינור עיבוד הנתונים של Apache Beam.
Java
מוסיפים את יחסי התלות הבאים לקובץ pom.xml:
<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>
בקטע הקוד הבא ב-Java אפשר לראות איך מגדירים ומריצים צינור עיבוד נתונים של CDC באצווה מ-Delta Lake אל Lakehouse:
import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.managed.Managed;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.Row;
public class DeltaLakeCdcToLakehouse {
public static void main(String[] args) {
Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());
// 1. Configure Delta Lake CDC reader
Map<String, Object> readConfig = new HashMap<>();
readConfig.put("table", "gs://BUCKET_NAME/TABLE_NAME");
// Starting commit bound (either start_version or start_timestamp must be provided):
readConfig.put("start_version", START_COMMIT_ID);
// Optional ending commit bound (end_version or end_timestamp, matching the start bound type):
// readConfig.put("end_version", END_COMMIT_ID);
readConfig.put("include_metadata_columns", Arrays.asList("_commit_version", "_change_type"));
// Read CDC events from Delta Lake
PCollection<Row> cdcRows =
pipeline.apply(
"ReadFromDeltaCDC",
Managed.read(Managed.DELTA_LAKE_CDC).withConfig(readConfig)).getSinglePCollection();
// 2. Configure BigLake REST Catalog properties for Lakehouse Iceberg
Map<String, String> catalogProps = new HashMap<>();
catalogProps.put("type", "rest");
catalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
catalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
catalogProps.put("header.x-goog-user-project", "PROJECT_ID");
catalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
catalogProps.put("rest.auth.type", "org.apache.iceberg.gcp.auth.GoogleAuthManager");
// 3. Map Delta Lake change types to Iceberg CDC change types
Map<String, String> changeTypeMap = new HashMap<>();
changeTypeMap.put("insert", "INSERT");
changeTypeMap.put("delete", "DELETE");
changeTypeMap.put("update_preimage", "UPDATE_BEFORE");
changeTypeMap.put("update_postimage", "UPDATE_AFTER");
// 4. Configure Lakehouse Iceberg CDC writer
Map<String, Object> writeConfig = new HashMap<>();
writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
writeConfig.put("catalog_name", "lakehouse");
writeConfig.put("catalog_properties", catalogProps);
writeConfig.put("mode", "merge-on-read");
writeConfig.put("sequence_number_column", "_commit_version");
writeConfig.put("equality_columns", Arrays.asList("PRIMARY_KEY_COLUMN"));
// change_type_column and change_type_map are only required for Portable Runner
writeConfig.put("change_type_column", "_change_type");
writeConfig.put("change_type_map", changeTypeMap);
// Write CDC rows to target Lakehouse Iceberg table
cdcRows.apply(
"WriteToLakehouse",
Managed.write(Managed.ICEBERG).withConfig(writeConfig));
pipeline.run();
}
}
מחליפים את מה שכתוב בשדות הבאים:
-
BUCKET_NAME: השם של קטגוריית Cloud Storage שמכילה את טבלת Delta Lake. -
TABLE_NAME: השם של ספריית הטבלה של Delta Lake במקור. -
START_COMMIT_ID: מספר גרסת הקומיט הראשונית לקריאת אירועי CDC של Delta Lake. חובה לציין אתstart_versionאו אתstart_timestamp. מזהה השמירה או חותמת הזמן של ההתחלה צריכים להיות אחרי שהמאפייןdelta.enableChangeDataFeedהופעל בטבלת המקור. -
END_COMMIT_ID: מספר הגרסה של הקומיט האחרון לקריאת אירועי CDC של Delta Lake. מומלץ לציין גבול סיום (end_versionאוend_timestamp) שתואם לסוג הגבול שמשמש להתחלה (start_versionאוstart_timestamp). אם לא מציינים גבול סיום, צינור עיבוד הנתונים קורא עד לגרסת הקומיט האחרונה. -
WAREHOUSE_BUCKET: השם של קטגוריית Cloud Storage שמשמשת כמחסן של קטלוג Lakehouse. -
PROJECT_ID: מזהה הפרויקט ב- Google Cloud . -
TARGET_NAMESPACE: מרחב השמות של טבלת Lakehouse של היעד. -
TARGET_TABLE: השם של טבלת היעד ב-Lakehouse. -
PRIMARY_KEY_COLUMN: שם העמודה של המפתח הראשי שמשמש לזיהוי שורות לעדכונים ולמחיקות של CDC.
בדיקת הפלט של המשרה
מוודאים שנתוני ה-CDC מוזגו בהצלחה לטבלת Lakehouse:
במסוף Google Cloud , עוברים לדף Studio של BigQuery.
בעורך השאילתות, מריצים שאילתת SQL כדי לאמת את הנתונים שסונכרנו:
SELECT * FROM `PROJECT_ID`.`CATALOG`.`NAMESPACE`.`TABLE_NAME` LIMIT 10;מחליפים את מה שכתוב בשדות הבאים:
-
PROJECT_ID: מזהה הפרויקט ב- Google Cloud . -
CATALOG: השם של קטלוג Lakehouse. -
NAMESPACE: מרחב השמות של טבלת Lakehouse. -
TABLE_NAME: השם של טבלת היעד ב-Lakehouse.
-
לוחצים על הפעלה ומאמתים את התוצאות.
המאמרים הבאים
- מידע נוסף על ייבוא טבלאות Delta Lake ל-Lakehouse
- מידע נוסף על הטמעה של נתונים לשינוי ב-Lakehouse
- מידע נוסף על קלט/פלט מנוהל ב-Dataflow