Delta Lake-CDC-Daten (Change Data Capture) mit Dataflow in Lakehouse migrieren

Es kann schwierig sein, Apache Iceberg-Zieltabelle im Lakehouse mit häufig aktualisierten Delta Lake-Quelltabelle zu synchronisieren. Das vollständige Neuladen von Tabellen, um Dataset-Änderungen zu erfassen, erhöht die Rechenkosten und die Verarbeitungslatenz. Das Erstellen benutzerdefinierter Pipelines zur Verarbeitung von Transaktionslogs führt zu einem komplexen Betriebsaufwand.

Mit geplanten Batch-CDC-Pipelines (Change Data Capture) von Dataflow können Sie Ihre Delta Lake-Quelltabelle kontinuierlich mit Lakehouse Apache Iceberg-Tabellen synchronisieren. Anstatt vollständige Tabellen neu zu laden, werden bei geplanten Batchjobs inkrementelle Aktualisierungen direkt aus Ihren Delta Lake-Transaktionsprotokollen extrahiert und nach einem konfigurierbaren Zeitplan auf Ihre Zieltabelle angewendet.

Die geplante Batch-CDC-Pipeline bietet die folgenden Funktionen:

  • Kostengünstige inkrementelle Synchronisierung: Es werden regelmäßig nur geänderte Daten verarbeitet (INSERT-, UPDATE- und DELETE-Vorgänge), wodurch der Rechenaufwand und die Datenlatenz reduziert werden.
  • Automatische Änderungszuordnung: Analysiert die Protokolle des Change Data Feed von Delta Lake und ordnet Änderungsvorgänge der Iceberg-Zieltabelle des Lakehouse zu.
  • Flexible Jobplanung: Batch-Dataflow-Jobs werden nach einem Zeitplan ausgeführt (z. B. alle 15 Minuten), der Ihren geschäftlichen Anforderungen entspricht.

Workflow für die kontinuierliche Migration

So richten Sie die kontinuierliche End-to-End-Migration von Delta Lake zu Lakehouse ein:

  1. Aktivieren Sie die Eigenschaft delta.enableChangeDataFeed (true) für die Delta Lake-Eingabetabelle.
  2. Führen Sie nach dem Aktivieren des Attributs delta.enableChangeDataFeed eine erste Tabellenmigration zu einer Commit-ID oder einem Zeitstempel durch. Eine Anleitung finden Sie unter Delta Lake-Tabellen mit Dataflow in Lakehouse importieren.
  3. Führen Sie geplante Batch-CDC-Jobs für nachfolgende Datenaktualisierungen mit einer Häufigkeit (Commit-ID oder Zeitstempelbereich) aus, die für Ihre Arbeitslastanforderungen ausgewählt wurde.

Hinweis

So richten Sie die geplante Batch-Migration für CDC ein:

  1. Aktivieren Sie die Dataflow-, BigQuery- und Lakehouse-APIs, falls sie noch nicht aktiviert sind.

    Rollen, die zum Aktivieren von APIs erforderlich sind

    Zum Aktivieren von APIs benötigen Sie die Berechtigung serviceusage.services.enable. Wenn Sie das Projekt erstellt haben, haben Sie diese Berechtigung wahrscheinlich bereits über die Rolle „Inhaber“ (roles/owner). Andernfalls können Sie diese Berechtigung über die Rolle „Service Usage-Administrator“ (roles/serviceusage.serviceUsageAdmin) erhalten. Informationen zum Zuweisen von Rollen

    APIs aktivieren

  2. Bitten Sie Ihren Administrator, Ihnen die erforderlichen IAM-Rollen (Identity and Access Management) für Ihr Projekt zuzuweisen, um die Berechtigungen zu erhalten, die Sie zum Erstellen der Ressourcen benötigen.

  3. Eine gültige Delta Lake-Tabelle, die in einem Cloud Storage-Bucket gespeichert ist. Das Tabellenverzeichnis muss Ihre Datendateien und das _delta_log/-Transaktionslogverzeichnis enthalten. Außerdem muss die Tabelle die folgenden Anforderungen erfüllen:

    • Das Attribut delta.enableChangeDataFeed muss für die Delta Lake-Tabelle aktiviert sein (true).
    • Die für die CDC-Pipeline konfigurierte Start-Commit-ID oder der Startzeitstempel muss nach der Aktivierung der Property delta.enableChangeDataFeed liegen.
  4. Ein vorhandener Lakehouse Iceberg-Katalog, in dem die synchronisierten Daten gespeichert werden. Wenn der Ziel-Namespace nicht vorhanden ist, wird er von der CDC-Pipeline automatisch erstellt. Wenn die Ziel-Lakehouse-Tabelle nicht vorhanden ist, wird sie von der CDC-Pipeline automatisch erstellt, sofern Sie eine gültige Liste von Gleichheitsspalten angeben.

Unterstützung und Einschränkungen

Für die geplante Batch-CDC-Migration von Delta Lake zu Lakehouse gelten die folgenden Überlegungen:

  • SDK-Anforderung:Erfordert das Apache Beam SDK in der Version 2.77.0 oder höher.
  • Erstellung der Zieltabelle:Das Erstellen der Ziel-Lakehouse-Tabelle wird unterstützt. Wenn die Zieltabelle nicht vorhanden ist, wird sie von der CDC-Pipeline automatisch erstellt, sofern Sie eine gültige Liste von Gleichheitsspalten (equality_columns) angeben. Es ist nur ein vorhandener Lakehouse-Katalog erforderlich. Der Namespace wird von der CDC-Pipeline automatisch erstellt, falls er nicht vorhanden ist.
  • Schemaentwicklung:Automatische Schemaänderungen werden nicht unterstützt. Wenn sich das Schema der Delta Lake-Quelltabelle ändert, müssen Sie das Schema der Lakehouse-Zieltabelle manuell anpassen, bevor Sie den CDC-Synchronisierungsjob ausführen.
  • Batchausführungsmodus:Delta Lake CDC wird im Batchmodus ausgeführt. Die Pipeline ruft Transaktionslogs einmal pro Jobausführung zwischen den angegebenen Start- und End-Commit-Grenzen (Commit-ID oder Zeitstempel) ab. Wenn keine End-Commit-Grenze (end_version oder end_timestamp) angegeben ist, werden in der Pipeline Daten bis zur neuesten Commit-Version gelesen.
  • Unterstützter Quellspeicher:Die Quelldaten müssen eine gültige Delta Lake-Tabelle sein, die in Cloud Storage (gs://) gespeichert ist. Amazon S3 wird für Delta Lake-Tabellenquellen nicht unterstützt.
  • Katalogbasierte Tabellen:Katalogbasierte Delta Lake-Tabellen (z. B. Unity Catalog) werden nicht unterstützt.
  • Nicht unterstützte Delta Lake-Versionen:Delta Lake-Version 1.2.1 oder früher wird nicht unterstützt, da CDF (Change Data Feed) für diese Versionen nicht unterstützt wird.
  • Automatischer Änderungsdatenfeed:Der automatische Änderungsdatenfeed wird nicht unterstützt, da keine CDF-Änderungsprotokolldateien geschrieben werden.

CDC-Pipeline programmatisch erstellen

Wenn Sie eine Batchmigrationspipeline für Delta Lake CDC programmatisch erstellen und ausführen möchten, verwenden Sie die Transformation Managed.read in Ihrer Apache Beam-Pipeline.

Java

Fügen Sie der Datei pom.xml die folgenden Abhängigkeiten hinzu:

<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>

Das folgende Java-Snippet zeigt, wie eine Batch-CDC-Pipeline von Delta Lake zu Lakehouse konfiguriert und ausgeführt wird:

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();
  }
}

Ersetzen Sie Folgendes:

  • BUCKET_NAME: Der Name des Cloud Storage-Buckets, der die Delta Lake-Tabelle enthält.
  • TABLE_NAME: Name des Quellverzeichnisses der Delta Lake-Tabelle.
  • START_COMMIT_ID: Startversionsnummer des Commits zum Lesen von Delta Lake-CDC-Ereignissen. Es muss entweder start_version oder start_timestamp angegeben werden. Die Start-Commit-ID oder der Start-Zeitstempel muss nach der Aktivierung der delta.enableChangeDataFeed-Property in der Quelltabelle liegen.
  • END_COMMIT_ID: Die Commit-Versionsnummer, mit der das Lesen von Delta Lake-CDC-Ereignissen beendet wird. Als Best Practice sollten Sie eine Endgrenze (end_version oder end_timestamp) angeben, die dem für den Start verwendeten Grenztyp (start_version oder start_timestamp) entspricht. Wenn eine Endgrenze weggelassen wird, liest die Pipeline bis zur neuesten Commit-Version.
  • WAREHOUSE_BUCKET: Der Name des Cloud Storage-Bucket, der als Lakehouse-Katalog-Data Warehouse verwendet wird.
  • PROJECT_ID: Projekt-ID in Google Cloud .
  • TARGET_NAMESPACE: der Namespace der Ziel-Lakehouse-Tabelle.
  • TARGET_TABLE: der Name der Lakehouse-Zieltabelle.
  • PRIMARY_KEY_COLUMN: Der Name der Primärschlüsselspalte, mit der Zeilen für CDC-Aktualisierungen und ‑Löschungen identifiziert werden.

Jobausgabe prüfen

Prüfen Sie, ob die CDC-Daten erfolgreich in Ihre Lakehouse-Tabelle zusammengeführt wurden:

  1. Rufen Sie in der Google Cloud Console die Seite BigQuery Studio auf.

    BigQuery aufrufen

  2. Führen Sie im Abfrageeditor eine SQL-Abfrage aus, um die synchronisierten Daten zu prüfen:

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

    Ersetzen Sie Folgendes:

    • PROJECT_ID: Projekt-ID in Google Cloud .
    • CATALOG: der Name Ihres Lakehouse-Katalogs.
    • NAMESPACE: Der Namespace Ihrer Lakehouse-Tabelle.
    • TABLE_NAME: der Name Ihrer Ziellakehouse-Tabelle.
  3. Klicken Sie auf Ausführen und prüfen Sie die Ergebnisse.

Nächste Schritte