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- undDELETE-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:
- Aktivieren Sie die Eigenschaft
delta.enableChangeDataFeed(true) für die Delta Lake-Eingabetabelle. - Führen Sie nach dem Aktivieren des Attributs
delta.enableChangeDataFeedeine erste Tabellenmigration zu einer Commit-ID oder einem Zeitstempel durch. Eine Anleitung finden Sie unter Delta Lake-Tabellen mit Dataflow in Lakehouse importieren. - 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:
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 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.enableChangeDataFeedmuss 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.enableChangeDataFeedliegen.
- Das Attribut
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_versionoderend_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 entwederstart_versionoderstart_timestampangegeben werden. Die Start-Commit-ID oder der Start-Zeitstempel muss nach der Aktivierung derdelta.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_versionoderend_timestamp) angeben, die dem für den Start verwendeten Grenztyp (start_versionoderstart_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:
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 Delta Lake-Tabellen in Lakehouse
- Change Data Capture-Aufnahme in Lakehouse
- Weitere Informationen zu verwalteten E/A in Dataflow