Esegui la migrazione dei dati Change Data Capture (CDC) di Apache Iceberg in Lakehouse utilizzando Dataflow

Mantenere sincronizzate le tabelle Apache Iceberg di destinazione di Lakehouse con le tabelle di origine Apache Iceberg aggiornate di frequente può essere difficile. L'esecuzione di ricaricamenti completi delle tabelle per acquisire le modifiche al set di dati aumenta i costi di calcolo e la latenza di elaborazione, mentre la creazione di pipeline personalizzate per elaborare i log Change Data Capture (CDC) introduce un overhead operativo complesso.

Utilizzando le pipeline CDC batch o di streaming di Dataflow, puoi sincronizzare continuamente o periodicamente le tabelle di origine Apache Iceberg con le tabelle Lakehouse Apache Iceberg. Anziché ricaricare tabelle complete, le pipeline CDC estraggono gli aggiornamenti incrementali direttamente dai log delle modifiche di Apache Iceberg dell'origine e li applicano alle tabelle di destinazione.

La pipeline CDC offre le seguenti funzionalità:

  • Sincronizzazione incrementale conveniente: elabora solo i dati modificati (operazioni INSERT, UPDATE e DELETE), riducendo il sovraccarico di calcolo e la latenza dei dati.
  • Modalità di esecuzione flessibili: supporta le modalità batch e streaming per soddisfare i requisiti aziendali.

La pipeline CDC supporta due modalità di esecuzione:

  • Modalità CDC batch: legge periodicamente gli aggiornamenti incrementali dalla tabella Apache Iceberg di origine tra i limiti degli snapshot di inizio e fine specificati.
  • Modalità CDC in streaming: esegue il polling continuo della tabella Apache Iceberg di origine per i nuovi commit di modifica con una frequenza configurabile (per impostazione predefinita 1 minuto, minimo 1 secondo) e li applica alla tabella lakehouse quasi in tempo reale.

Prima di iniziare

Per configurare la migrazione CDC da Apache Iceberg a Lakehouse, assicurati di disporre di quanto segue:

  1. Abilita le API Dataflow, BigQuery e Lakehouse, se non sono già abilitate.

    Ruoli richiesti per abilitare le API

    Per abilitare le API, devi disporre dell'autorizzazione serviceusage.services.enable. Se hai creato il progetto, probabilmente disponi già di questa autorizzazione tramite il ruolo Proprietario (roles/owner). In caso contrario, puoi ottenere questa autorizzazione tramite il ruolo Amministratore utilizzo servizi (roles/serviceusage.serviceUsageAdmin). Scopri come concedere i ruoli.

    Abilita le API

  2. Per ottenere le autorizzazioni necessarie per creare le risorse, chiedi all'amministratore di concederti i ruoli Identity and Access Management (IAM) richiesti per il tuo progetto.

  3. Una tabella Apache Iceberg di origine valida archiviata in Cloud Storage o registrata in un catalogo REST o metastore Iceberg.

  4. Un catalogo Lakehouse Iceberg esistente per ricevere i dati sincronizzati.

    Se lo spazio dei nomi di destinazione non esiste, la pipeline CDC lo crea automaticamente. Se la tabella Lakehouse di destinazione non esiste, la pipeline CDC la crea automaticamente se fornisci un elenco valido di colonne di uguaglianza.

Supporto e limitazioni

La migrazione CDC da Apache Iceberg a Lakehouse presenta le seguenti considerazioni:

  • Requisito SDK:richiede l'SDK Apache Beam versione 2.77.0 o successive.
  • Creazione della tabella di destinazione:è supportata la creazione della tabella Lakehouse di destinazione. Se la tabella di destinazione non esiste, la pipeline CDC la crea automaticamente, a condizione che tu fornisca un elenco valido di colonne di uguaglianza (equality_columns). È richiesto solo un catalogo Lakehouse esistente; la pipeline CDC crea automaticamente lo spazio dei nomi se non esiste.
  • Evoluzione dello schema:le modifiche automatiche dello schema non sono supportate. Se lo schema della tabella Apache Iceberg di origine cambia, devi aggiornare manualmente lo schema della tabella lakehouse di destinazione in modo che corrisponda prima di eseguire il job di sincronizzazione CDC.
  • Modalità di esecuzione:
    • CDC batch:estrae le modifiche del log delle transazioni una volta per esecuzione del job tra i limiti dello snapshot iniziale e finale specificati (ID snapshot o timestamp).
    • CDC in streaming: esegue il polling continuo della tabella Apache Iceberg di origine per i nuovi log di commit. La frequenza di polling è configurabile (il valore predefinito è 1 minuto, il minimo è 1 secondo).

Crea una pipeline CDC utilizzando Java

Per creare ed eseguire in modo programmatico una pipeline di migrazione batch o di streaming CDC Apache Iceberg utilizzando l'SDK Apache Beam Java, utilizza la trasformazione Managed.read con Managed.ICEBERG_CDC.

I seguenti esempi Java mostrano come configurare ed eseguire pipeline CDC Apache Iceberg per carichi di lavoro batch e di streaming.

Pipeline CDC batch

Il seguente esempio Java configura una pipeline batch che legge gli eventi CDC di Apache Iceberg entro i limiti degli snapshot di inizio e fine specificati:

Java

Aggiungi le seguenti dipendenze al file 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>

Lo snippet Java seguente mostra come configurare ed eseguire una pipeline CDC batch da una tabella Iceberg a 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();
  }
}

Sostituisci quanto segue:

  • SOURCE_CATALOG_TYPE: il tipo di catalogo Iceberg di origine (ad esempio hadoop o hive).
  • SOURCE_WAREHOUSE_LOCATION: la posizione del warehouse del catalogo Iceberg di origine (ad esempio s3://source-warehouse o gs://source-warehouse).
  • SOURCE_NAMESPACE: lo spazio dei nomi della tabella Iceberg di origine.
  • SOURCE_TABLE: il nome della tabella Iceberg di origine.
  • START_SNAPSHOT_ID: l'ID dello snapshot iniziale per la lettura CDC batch.
  • END_SNAPSHOT_ID: l'ID dello snapshot finale per la lettura CDC batch.
  • WAREHOUSE_BUCKET: il nome del bucket Cloud Storage utilizzato come warehouse del catalogo Lakehouse.
  • PROJECT_ID: il tuo ID progetto Google Cloud .
  • TARGET_NAMESPACE: lo spazio dei nomi della tabella Lakehouse di destinazione.
  • TARGET_TABLE: il nome della tabella Lakehouse di destinazione.

Pipeline CDC in modalità flusso

L'esempio Java seguente configura una pipeline di streaming che esegue il polling continuo della tabella Apache Iceberg di origine per nuove modifiche CDC:

Java

Aggiungi le seguenti dipendenze al file 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>

Il seguente snippet Java mostra come configurare ed eseguire una pipeline CDC di streaming da una tabella Iceberg a 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();
  }
}

Sostituisci quanto segue:

  • SOURCE_CATALOG_TYPE: il tipo di catalogo Iceberg di origine (ad esempio hadoop o hive).
  • SOURCE_WAREHOUSE_LOCATION: la posizione del warehouse del catalogo Iceberg di origine (ad esempio s3://source-warehouse o gs://source-warehouse).
  • SOURCE_NAMESPACE: lo spazio dei nomi della tabella Iceberg di origine.
  • SOURCE_TABLE: il nome della tabella Iceberg di origine.
  • POLL_INTERVAL_SECONDS: la frequenza di polling continuo in secondi (valore predefinito 60, valore minimo 1).
  • WAREHOUSE_BUCKET: il nome del bucket Cloud Storage utilizzato come warehouse del catalogo Lakehouse.
  • PROJECT_ID: il tuo ID progetto Google Cloud .
  • TARGET_NAMESPACE: lo spazio dei nomi della tabella Lakehouse di destinazione.
  • TARGET_TABLE: il nome della tabella Lakehouse di destinazione.

Esamina l'output del job

Verifica che i dati CDC siano stati uniti correttamente nella tabella Lakehouse:

  1. Nella console Google Cloud , vai alla pagina BigQuery Studio.

    Vai a BigQuery

  2. Nell'editor di query, esegui una query SQL per verificare i dati sincronizzati:

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

    Sostituisci quanto segue:

    • PROJECT_ID: il tuo ID progetto Google Cloud .
    • CATALOG: il nome del catalogo Lakehouse.
    • NAMESPACE: lo spazio dei nomi della tabella lakehouse.
    • TABLE_NAME: il nome della tabella Lakehouse di destinazione.
  3. Fai clic su Esegui e verifica i risultati.

Passaggi successivi