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

Mantenere sincronizzate le tabelle Apache Iceberg di Lakehouse di destinazione con le tabelle di origine Delta Lake aggiornate di frequente può essere difficile. L'esecuzione di ricariche complete 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 delle transazioni introduce un overhead operativo complesso.

Utilizzando le pipeline di acquisizione dei dati delle modifiche (CDC) batch pianificate di Dataflow, puoi sincronizzare continuamente le tabelle di origine Delta Lake con le tabelle Lakehouse Apache Iceberg. Invece di ricaricare le tabelle complete, i job batch pianificati estraggono gli aggiornamenti incrementali direttamente dai log delle transazioni Delta Lake e li applicano alle tabelle di destinazione in base a una pianificazione configurabile.

La pipeline CDC batch pianificata offre le seguenti funzionalità:

  • Sincronizzazione incrementale conveniente: elabora periodicamente solo i dati modificati (operazioni INSERT, UPDATE e DELETE), riducendo il sovraccarico di calcolo e la latenza dei dati.
  • Mappatura automatica delle modifiche: analizza i log del feed di dati delle modifiche di Delta Lake e mappa le operazioni di modifica alla tabella Iceberg di Lakehouse di destinazione.
  • Pianificazione flessibile dei job: esegue i job Dataflow in modalità batch in base a una pianificazione (ad esempio, intervalli di 15 minuti) adatta alle tue esigenze aziendali.

Workflow di migrazione continua

Per configurare la migrazione continua end-to-end da Delta Lake a Lakehouse, completa la seguente sequenza di operazioni:

  1. Attiva la proprietà delta.enableChangeDataFeed (true) nella tabella Delta Lake di input.
  2. Esegui una migrazione iniziale della tabella a un ID commit o un timestamp dopo aver attivato la proprietà delta.enableChangeDataFeed. Per istruzioni, vedi Importare tabelle Delta Lake in Lakehouse utilizzando Dataflow.
  3. Esegui job batch CDC pianificati per aggiornamenti successivi dei dati con una frequenza (intervallo di ID commit o timestamp) scelta in base ai requisiti del tuo workload.

Prima di iniziare

Per configurare la migrazione CDC batch pianificata, 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 Delta Lake valida archiviata in un bucket Cloud Storage. La directory della tabella deve contenere i file di dati e la directory dei log delle transazioni _delta_log/. Inoltre, la tabella deve soddisfare i seguenti requisiti:

    • La proprietà delta.enableChangeDataFeed deve essere abilitata (true) nella tabella Delta Lake.
    • L'ID commit iniziale o il timestamp configurato per la pipeline CDC deve essere successivo all'attivazione della proprietà delta.enableChangeDataFeed.
  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 batch pianificata da Delta Lake 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 se fornisci 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 Delta Lake 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 batch:CDC di Delta Lake funziona in modalità batch. La pipeline estrae i log delle transazioni una volta per esecuzione del job tra i limiti di commit iniziale e finale specificati (ID commit o timestamp). Se non viene specificato un commit finale (end_version o end_timestamp), la pipeline legge fino alla versione del commit più recente.
  • Archiviazione di origine supportata:i dati di origine devono essere una tabella Delta Lake valida archiviata in Cloud Storage (gs://). Amazon S3 non è supportato per le origini delle tabelle Delta Lake.
  • Tabelle basate sul catalogo:le tabelle Delta Lake basate sul catalogo (ad esempio Unity Catalog) non sono supportate.
  • Versioni di Delta Lake non supportate: la versione 1.2.1 o precedente di Delta Lake non è supportata perché Change Data Feed (CDF) non è supportato per queste versioni.
  • Feed di dati sulle modifiche automatiche:il feed di dati sulle modifiche automatiche non è supportato perché non scrive file di log delle modifiche CDF.

Crea una pipeline CDC in modo programmatico

Per creare ed eseguire in modo programmatico una pipeline di migrazione batch CDC Delta Lake, utilizza la trasformazione Managed.read nella pipeline Apache Beam.

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 batch da Delta Lake 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.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();
  }
}

Sostituisci quanto segue:

  • BUCKET_NAME: il nome del bucket Cloud Storage che contiene la tabella Delta Lake.
  • TABLE_NAME: il nome della directory della tabella Delta Lake di origine.
  • START_COMMIT_ID: numero di versione del commit iniziale per la lettura degli eventi CDC di Delta Lake. È necessario fornire start_version o start_timestamp. L'ID commit iniziale o il timestamp deve essere successivo all'attivazione della proprietà delta.enableChangeDataFeed nella tabella di origine.
  • END_COMMIT_ID: numero di versione del commit finale per la lettura degli eventi CDC Delta Lake. Come best practice, specifica un limite finale (end_version o end_timestamp) che corrisponda al tipo di limite utilizzato per l'inizio (start_version o start_timestamp). Se un limite finale viene omesso, la pipeline legge fino alla versione del commit più recente.
  • 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.
  • PRIMARY_KEY_COLUMN: il nome della colonna della chiave primaria utilizzata per identificare le righe per gli aggiornamenti e le eliminazioni CDC.

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