העברת נתונים של סימון נתונים שהשתנו (CDC) מ-Apache Iceberg אל Lakehouse באמצעות Dataflow

יכול להיות שיהיה קשה לשמור על סנכרון בין טבלאות היעד של Lakehouse Apache Iceberg לבין טבלאות המקור של Apache Iceberg שמתעדכנות לעיתים קרובות. ביצוע טעינות מחדש של טבלאות מלאות כדי לתעד שינויים במערך הנתונים מגדיל את עלויות החישוב ואת זמן האחזור של העיבוד, ובניית צינורות מותאמים אישית לעיבוד יומני סימון נתונים שהשתנו (CDC) יוצרת תקורה תפעולית מורכבת.

באמצעות צינורות עיבוד נתונים של Dataflow מסוג CDC (לעיבוד נתונים באצווה או בסטרימינג), אתם יכולים לסנכרן באופן רציף או תקופתי את טבלאות המקור של Apache Iceberg עם טבלאות Lakehouse Apache Iceberg. במקום לטעון מחדש טבלאות מלאות, צינורות CDC מחלצים עדכונים מצטברים ישירות מיומני השינויים של Apache Iceberg במקור ומחילים אותם על טבלאות היעד.

צינור ה-CDC מספק את היכולות הבאות:

  • סנכרון מצטבר חסכוני: מעבד רק נתונים שהשתנו (פעולות INSERT, UPDATE ו-DELETE), וכך מצמצם את התקורה של המחשוב ואת זמן האחזור של הנתונים.
  • מצבי הפעלה גמישים: תמיכה במצב אצווה ובמצב סטרימינג כדי להתאים לדרישות העסקיות שלכם.

צינור ה-CDC תומך בשני מצבי ביצוע:

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

לפני שמתחילים

כדי להגדיר העברת CDC מ-Apache Iceberg ל-Lakehouse, צריך לוודא שיש לכם:

  1. מפעילים את ממשקי ה-API של Dataflow,‏ BigQuery ו-Lakehouse, אם הם עדיין לא מופעלים.

    תפקידים שנדרשים להפעלת ממשקי API

    כדי להפעיל ממשקי API, צריך את ההרשאה serviceusage.services.enable. אם יצרתם את הפרויקט, סביר להניח שכבר יש לכם את ההרשאה הזו דרך התפקיד 'בעלים' (roles/owner). אחרת, תוכלו לקבל את ההרשאה הזו דרך התפקיד 'אדמין בממשק Service Usage' (roles/serviceusage.serviceUsageAdmin). איך מקצים תפקידים

    הפעלת ממשקי ה-API

  2. כדי לקבל את ההרשאות שדרושות ליצירת המשאבים, צריך לבקש מהאדמין להקצות לכם את התפקידים הנדרשים בניהול הזהויות והרשאות הגישה (IAM) בפרויקט.

  3. טבלה תקינה של Apache Iceberg שמאוחסנת ב-Cloud Storage או רשומה בקטלוג REST או במאגר מטא-נתונים של Iceberg.

  4. קטלוג קיים של Iceberg Lakehouse שאליו יסונכרנו הנתונים.

    אם מרחב השמות של היעד לא קיים, צינור ה-CDC יוצר אותו באופן אוטומטי. אם טבלת Lakehouse של היעד לא קיימת, צינור ה-CDC יוצר אותה באופן אוטומטי אם מספקים רשימה תקינה של עמודות שוויון.

תמיכה ומגבלות

כשמבצעים העברה של CDC מ-Apache Iceberg ל-Lakehouse, צריך לקחת בחשבון את הנקודות הבאות:

  • דרישות SDK: נדרשת גרסת Apache Beam SDK 2.77.0 ואילך.
  • יצירת טבלת יעד: יש תמיכה ביצירת טבלת יעד ב-Lakehouse. אם טבלת היעד לא קיימת, צינור ה-CDC יוצר אותה באופן אוטומטי, בתנאי שסיפקתם רשימה תקינה של עמודות שוויון (equality_columns). נדרש רק קטלוג Lakehouse קיים; צינור ה-CDC יוצר את מרחב השמות באופן אוטומטי אם הוא לא קיים.
  • התפתחות סכימה: לא ניתן לבצע שינויים אוטומטיים בסכימה. אם הסכימה של טבלת המקור ב-Apache Iceberg משתנה, צריך לעדכן באופן ידני את הסכימה של טבלת היעד ב-Lakehouse כך שתהיה זהה לסכימה של טבלת המקור לפני שמריצים את משימת הסנכרון של CDC.
  • מצבי הפעלה:
    • CDC של קבוצות: שולף שינויים ביומן העסקאות פעם אחת בכל הפעלה של משימה בין גבולות תמונת מצב של התחלה וסיום שצוינו (מזהה תמונת מצב או חותמת זמן).
    • הזרמת CDC: סקר רציף של טבלת המקור Apache Iceberg כדי למצוא יומני ביצוע חדשים. אפשר להגדיר את תדירות הבדיקה (ברירת המחדל היא דקה אחת, המינימום הוא שנייה אחת).

יצירת צינור עיבוד נתונים של CDC באמצעות Java

כדי ליצור ולהפעיל באופן פרוגרמטי צינור העברה של נתונים בסטרימינג או ברצף (batch) של Apache Iceberg CDC באמצעות Apache Beam Java SDK, משתמשים בטרנספורמציה Managed.read עם Managed.ICEBERG_CDC.

בדוגמאות הבאות ב-Java מוסבר איך להגדיר ולהפעיל צינורות עיבוד נתונים של Apache Iceberg CDC לעומסי עבודה של אצווה וסטרימינג.

פייפליין של CDC באצווה

בדוגמה הבאה של Java מוגדר צינור להעברת נתונים באצווה שקורא אירועי CDC של Apache Iceberg בתוך גבולות ספציפיים של תמונת מצב של התחלה וסיום:

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 באצווה מטבלת Iceberg אל Lakehouse:

import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.iceberg.cdc.IcebergCdcMetadataColumns;
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 IcebergBatchCdcToLakehouse {
  public static void main(String[] args) {
    Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());

    // Configure source catalog properties
    Map<String, String> sourceCatalogProps = new HashMap<>();
    sourceCatalogProps.put("type", "SOURCE_CATALOG_TYPE");
    sourceCatalogProps.put("warehouse", "SOURCE_WAREHOUSE_LOCATION");
    // Add additional properties if required by your catalog type (for example, uri)

    // Configure batch Iceberg CDC reader with snapshot bounds
    Map<String, Object> readConfig = new HashMap<>();
    readConfig.put("table", "SOURCE_NAMESPACE.SOURCE_TABLE");
    readConfig.put("catalog_properties", sourceCatalogProps);
    readConfig.put("start_snapshot_id", START_SNAPSHOT_ID);
    readConfig.put("end_snapshot_id", END_SNAPSHOT_ID);
    readConfig.put(
        "include_metadata_columns",
        Arrays.asList(
            IcebergCdcMetadataColumns.CHANGE_TYPE,
            IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER));

    // Read batch CDC events from source Iceberg table
    PCollection<Row> cdcRows =
        pipeline.apply(
            "ReadFromIcebergCDC",
            Managed.read(Managed.ICEBERG_CDC).withConfig(readConfig)).getSinglePCollection();

    // Configure destination Lakehouse REST Catalog properties for Lakehouse Iceberg
    Map<String, String> destinationCatalogProps = new HashMap<>();
    destinationCatalogProps.put("type", "rest");
    destinationCatalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
    destinationCatalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
    destinationCatalogProps.put("header.x-goog-user-project", "PROJECT_ID");
    destinationCatalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
    destinationCatalogProps.put("rest.auth.type", "google");

    // Configure Lakehouse Iceberg writer
    Map<String, Object> writeConfig = new HashMap<>();
    writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
    writeConfig.put("catalog_properties", destinationCatalogProps);
    writeConfig.put("mode", "merge-on-read");
    writeConfig.put("change_type_column", IcebergCdcMetadataColumns.CHANGE_TYPE);

    // Write CDC rows to target Lakehouse Iceberg table
    cdcRows.apply(
        "WriteToLakehouse",
        Managed.write(Managed.ICEBERG).withConfig(writeConfig));

    pipeline.run();
  }
}

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

  • ‫SOURCE_CATALOG_TYPE: סוג קטלוג Iceberg של המקור (לדוגמה, hadoop או hive).
  • ‫SOURCE_WAREHOUSE_LOCATION: מיקום מחסן הנתונים של קטלוג המקור של Iceberg (לדוגמה, s3://source-warehouse או gs://source-warehouse).
  • ‫SOURCE_NAMESPACE: מרחב השמות של טבלת Iceberg של המקור.
  • ‫SOURCE_TABLE: השם של טבלת Iceberg של המקור.
  • ‫START_SNAPSHOT_ID: מזהה קובץ ה-snapshot ההתחלתי לקריאת ה-CDC של הקבוצה.
  • ‫END_SNAPSHOT_ID: מזהה קובץ ה-snapshot הסופי לקריאת ה-CDC של האצווה.
  • ‫WAREHOUSE_BUCKET: השם של קטגוריית Cloud Storage שמשמשת כמחסן של קטלוג Lakehouse.
  • ‫PROJECT_ID: מזהה הפרויקט ב- Google Cloud .
  • ‫TARGET_NAMESPACE: מרחב השמות של טבלת Lakehouse של היעד.
  • ‫TARGET_TABLE: השם של טבלת היעד ב-Lakehouse.

פייפליין CDC בסטרימינג

בדוגמה הבאה של Java מוגדר פייפליין זורם שבודק באופן רציף את טבלת המקור Apache Iceberg כדי לזהות שינויים חדשים ב-CDC:

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 בסטרימינג מטבלת Iceberg אל Lakehouse:

import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.iceberg.cdc.IcebergCdcMetadataColumns;
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 IcebergStreamingCdcToLakehouse {
  public static void main(String[] args) {
    Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());

    // Configure source catalog properties
    Map<String, String> sourceCatalogProps = new HashMap<>();
    sourceCatalogProps.put("type", "SOURCE_CATALOG_TYPE");
    sourceCatalogProps.put("warehouse", "SOURCE_WAREHOUSE_LOCATION");
    // Add additional properties if required by your catalog type (for example, uri)

    // Configure streaming Iceberg CDC reader with continuous polling
    Map<String, Object> readConfig = new HashMap<>();
    readConfig.put("table", "SOURCE_NAMESPACE.SOURCE_TABLE");
    readConfig.put("catalog_properties", sourceCatalogProps);
    readConfig.put("streaming", true);
    readConfig.put("poll_interval_seconds", POLL_INTERVAL_SECONDS);
    readConfig.put(
        "include_metadata_columns",
        Arrays.asList(
            IcebergCdcMetadataColumns.CHANGE_TYPE,
            IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER));

    // Read streaming CDC events from source Iceberg table
    PCollection<Row> cdcRows =
        pipeline.apply(
            "ReadFromIcebergCDC",
            Managed.read(Managed.ICEBERG_CDC).withConfig(readConfig)).getSinglePCollection();

    // Configure destination BigLake REST Catalog properties for Lakehouse Iceberg
    Map<String, String> destinationCatalogProps = new HashMap<>();
    destinationCatalogProps.put("type", "rest");
    destinationCatalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
    destinationCatalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
    destinationCatalogProps.put("header.x-goog-user-project", "PROJECT_ID");
    destinationCatalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
    destinationCatalogProps.put("rest.auth.type", "google");

    // Configure Lakehouse Iceberg writer
    Map<String, Object> writeConfig = new HashMap<>();
    writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
    writeConfig.put("catalog_properties", destinationCatalogProps);
    writeConfig.put("mode", "merge-on-read");
    writeConfig.put("change_type_column", IcebergCdcMetadataColumns.CHANGE_TYPE);

    // Write CDC rows to target Lakehouse Iceberg table
    cdcRows.apply(
        "WriteToLakehouse",
        Managed.write(Managed.ICEBERG).withConfig(writeConfig));

    pipeline.run();
  }
}

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

  • ‫SOURCE_CATALOG_TYPE: סוג קטלוג Iceberg של המקור (לדוגמה, hadoop או hive).
  • ‫SOURCE_WAREHOUSE_LOCATION: מיקום מחסן הנתונים של קטלוג המקור של Iceberg (לדוגמה, s3://source-warehouse או gs://source-warehouse).
  • ‫SOURCE_NAMESPACE: מרחב השמות של טבלת Iceberg של המקור.
  • ‫SOURCE_TABLE: השם של טבלת Iceberg של המקור.
  • ‫POLL_INTERVAL_SECONDS: תדירות הסקר הרציף בשניות (ברירת מחדל 60, מינימום 1).
  • ‫WAREHOUSE_BUCKET: השם של קטגוריית Cloud Storage שמשמשת כמחסן של קטלוג Lakehouse.
  • ‫PROJECT_ID: מזהה הפרויקט ב- Google Cloud .
  • ‫TARGET_NAMESPACE: מרחב השמות של טבלת Lakehouse של היעד.
  • ‫TARGET_TABLE: השם של טבלת היעד ב-Lakehouse.

בדיקת הפלט של המשרה

מוודאים שנתוני ה-CDC מוזגו בהצלחה לטבלת Lakehouse:

  1. במסוף Google Cloud , עוברים לדף Studio של BigQuery.

    כניסה ל-BigQuery

  2. בעורך השאילתות, מריצים שאילתת SQL כדי לאמת את הנתונים שסונכרנו:

    SELECT * FROM `PROJECT_ID`.`CATALOG`.`NAMESPACE`.`TABLE_NAME` LIMIT 10;
    

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

    • ‫PROJECT_ID: מזהה הפרויקט ב- Google Cloud .
    • ‫CATALOG: השם של קטלוג Lakehouse.
    • ‫NAMESPACE: מרחב השמות של טבלת Lakehouse.
    • ‫TABLE_NAME: השם של טבלת היעד ב-Lakehouse.
  3. לוחצים על הפעלה ומאמתים את התוצאות.

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