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- undDELETE-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:
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 RollenBitten 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.
Eine gültige Apache Iceberg-Quelltabelle, die in Cloud Storage gespeichert oder in einem Iceberg-REST-Katalog oder Metastore registriert ist.
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.hadoopoderhive).SOURCE_WAREHOUSE_LOCATION: Der Warehouse-Standort des Iceberg-Quellkatalogs (z. B.s3://source-warehouseodergs://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.hadoopoderhive).SOURCE_WAREHOUSE_LOCATION: Der Warehouse-Standort des Iceberg-Quellkatalogs (z. B.s3://source-warehouseodergs://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:
Rufen Sie in der Google Cloud Console die Seite BigQuery Studio auf.
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.
Klicken Sie auf Ausführen und prüfen Sie die Ergebnisse.
Nächste Schritte
- Weitere Informationen zum Importieren von Iceberg-Tabellen in Lakehouse
- Change Data Capture-Aufnahme in Lakehouse
- Weitere Informationen zu verwalteten E/A in Dataflow