Apache Iceberg-CDC-Daten (Change Data Capture) mit Dataflow in Lakehouse migrieren

Es kann schwierig sein, Ziel-Lakehouse-Apache Iceberg-Tabellen mit häufig aktualisierten Apache Iceberg-Quelltabellen zu synchronisieren. Wenn Sie vollständige Tabellen neu laden, um Dataset-Änderungen zu erfassen, steigen die Rechenkosten und die Verarbeitungslatenz. Wenn Sie benutzerdefinierte Pipelines erstellen, um CDC-Logs (Change Data Capture) zu verarbeiten, entsteht ein komplexer Betriebsaufwand.

Mit Dataflow-Batch- oder Streaming-CDC-Pipelines können Sie Ihre Apache Iceberg-Quelltabelle kontinuierlich oder regelmäßig mit Lakehouse Apache Iceberg-Tabellen synchronisieren. Anstatt vollständige Tabellen neu zu laden, werden bei CDC-Pipelines inkrementelle Updates direkt aus den Apache Iceberg-Änderungsprotokollen Ihrer Quelle extrahiert und auf Ihre Zieltabelle angewendet.

Die CDC-Pipeline bietet die folgenden Funktionen:

  • Kostengünstige inkrementelle Synchronisierung: Es werden nur geänderte Daten verarbeitet (INSERT-, UPDATE- und DELETE-Vorgänge). Dadurch werden der Rechenaufwand und die Datenlatenz reduziert.
  • Flexible Ausführungsmodi: Unterstützt sowohl Batch- als auch Streamingmodi, um Ihren geschäftlichen Anforderungen gerecht zu werden.

Die CDC-Pipeline unterstützt zwei Ausführungsmodi:

  • Batch-CDC-Modus:Inkrementelle Updates werden in regelmäßigen Abständen aus der Apache Iceberg-Quelltabelle zwischen den angegebenen Start- und End-Snapshot-Grenzen gelesen.
  • Streaming-CDC-Modus:Die Apache Iceberg-Quelltabelle wird in einem konfigurierbaren Intervall (standardmäßig 1 Minute, mindestens 1 Sekunde) kontinuierlich nach neuen Änderungs-Commits durchsucht und diese werden nahezu in Echtzeit auf Ihre Lakehouse-Tabelle angewendet.

Hinweis

Für die Einrichtung der CDC-Migration von Apache Iceberg zu Lakehouse benötigen Sie Folgendes:

  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 Apache Iceberg-Quelltabelle, die in Cloud Storage gespeichert oder in einem Iceberg-REST-Katalog oder Metastore registriert ist.

  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

Bei der CDC-Migration von Apache Iceberg zu Lakehouse sind folgende Aspekte zu beachten:

  • 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 Apache Iceberg-Quelltabelle ändert, müssen Sie das Schema der Lakehouse-Zieltabelle manuell aktualisieren, damit es übereinstimmt, bevor Sie den CDC-Synchronisierungsjob ausführen.
  • Ausführungsmodi:
    • Batch-CDC:Änderungen im Transaktionslog werden einmal pro Jobausführung zwischen den angegebenen Snapshot-Grenzen (Snapshot-ID oder Zeitstempel) abgerufen.
    • Streaming-CDC:Die Quell-Apache Iceberg-Tabelle wird kontinuierlich nach neuen Commit-Logs abgefragt. Die Abfragehäufigkeit ist konfigurierbar (Standardeinstellung: 1 Minute, Minimum: 1 Sekunde).

CDC-Pipeline mit Java erstellen

Wenn Sie eine Apache Iceberg-CDC-Batch- oder Streaming-Migrationspipeline mit dem Apache Beam Java SDK programmieren und ausführen möchten, verwenden Sie die Managed.read-Transformation mit Managed.ICEBERG_CDC.

Die folgenden Java-Beispiele zeigen, wie Sie Apache Iceberg-CDC-Pipelines für Batch- und Streaming-Arbeitslasten konfigurieren und ausführen.

Batch-CDC-Pipeline

Im folgenden Java-Beispiel wird eine Batchpipeline konfiguriert, die Apache Iceberg-CDC-Ereignisse innerhalb der angegebenen Snapshot-Grenzwerte für Start und Ende liest:

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 Sie eine Batch-CDC-Pipeline von einer Iceberg-Tabelle zu Lakehouse konfigurieren und ausführen:

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

Ersetzen Sie Folgendes:

  • SOURCE_CATALOG_TYPE: Der Typ des Iceberg-Quellkatalogs (z. B. hadoop oder hive).
  • SOURCE_WAREHOUSE_LOCATION: Der Warehouse-Standort des Iceberg-Quellkatalogs (z. B. s3://source-warehouse oder gs://source-warehouse).
  • SOURCE_NAMESPACE: der Namespace der Iceberg-Quelltabelle.
  • SOURCE_TABLE: Der Name der Iceberg-Quelltabelle.
  • START_SNAPSHOT_ID: Die ID des Start-Snapshots für das Batch-CDC-Lesen.
  • END_SNAPSHOT_ID: die ID des End-Snapshots für den Batch-CDC-Lesevorgang.
  • 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.

Streaming-CDC-Pipeline

Im folgenden Java-Beispiel wird eine Streamingpipeline konfiguriert, die die Apache Iceberg-Quelltabelle kontinuierlich nach neuen CDC-Änderungen abfragt:

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 Sie eine Streaming-CDC-Pipeline von einer Iceberg-Tabelle zu Lakehouse konfigurieren und ausführen:

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

Ersetzen Sie Folgendes:

  • SOURCE_CATALOG_TYPE: Der Typ des Iceberg-Quellkatalogs (z. B. hadoop oder hive).
  • SOURCE_WAREHOUSE_LOCATION: Der Warehouse-Standort des Iceberg-Quellkatalogs (z. B. s3://source-warehouse oder gs://source-warehouse).
  • SOURCE_NAMESPACE: der Namespace der Iceberg-Quelltabelle.
  • SOURCE_TABLE: Der Name der Iceberg-Quelltabelle.
  • POLL_INTERVAL_SECONDS: die Häufigkeit des kontinuierlichen Polling in Sekunden (Standardwert: 60, Mindestwert: 1).
  • 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.

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