Migrer les données de capture des données modifiées (CDC) Delta Lake vers Lakehouse à l'aide de Dataflow

Il peut être difficile de synchroniser les tables Apache Iceberg Lakehouse cibles avec les tables sources Delta Lake fréquemment mises à jour. L'exécution de rechargements complets des tables pour capturer les modifications apportées aux ensembles de données augmente les coûts de calcul et la latence de traitement. La création de pipelines personnalisés pour traiter les journaux de transactions introduit une surcharge opérationnelle complexe.

En utilisant des pipelines de capture des données modifiées (CDC) par lot planifiés Dataflow, vous pouvez synchroniser en continu vos tables sources Delta Lake avec les tables Lakehouse Apache Iceberg. Au lieu de recharger des tables entières, les jobs par lot planifiés extraient les mises à jour incrémentielles directement de vos journaux de transactions Delta Lake et les appliquent à vos tables cibles selon une planification configurable.

Le pipeline CDC par lot planifié offre les fonctionnalités suivantes :

  • Synchronisation incrémentielle économique : traite périodiquement uniquement les données modifiées (opérations INSERT, UPDATE et DELETE), ce qui réduit la surcharge de calcul et la latence des données.
  • Mappage automatique des modifications : analyse les journaux Change Data Feed de Delta Lake et mappe les opérations de modification à la table Iceberg Lakehouse cible.
  • Planification flexible des jobs : exécute des jobs Dataflow en mode batch selon un calendrier (par exemple, à intervalles de 15 minutes) qui correspond à vos besoins commerciaux.

Workflow de migration continue

Pour configurer une migration continue de bout en bout de Delta Lake vers Lakehouse, effectuez les opérations suivantes :

  1. Activez la propriété delta.enableChangeDataFeed (true) sur la table Delta Lake d'entrée.
  2. Effectuez une migration initiale des tables vers un ID de commit ou un code temporel après avoir activé la propriété delta.enableChangeDataFeed. Pour obtenir des instructions, consultez Importer des tables Delta Lake dans Lakehouse à l'aide de Dataflow.
  3. Exécutez des jobs CDC par lot planifiés pour les mises à jour de données ultérieures avec une fréquence (ID de commit ou plage de code temporel) choisie en fonction des exigences de votre charge de travail.

Avant de commencer

Pour configurer la migration CDC par lot planifiée, assurez-vous de disposer des éléments suivants :

  1. Activez les API Dataflow, BigQuery et Lakehouse, si ce n'est pas déjà fait.

    Rôles requis pour activer les API

    Pour activer les API, vous devez disposer de l'autorisation serviceusage.services.enable. Si vous avez créé le projet, vous disposez probablement déjà de cette autorisation grâce au rôle Propriétaire (roles/owner). Sinon, vous pouvez obtenir cette autorisation grâce au rôle Administrateur Service Usage (roles/serviceusage.serviceUsageAdmin). Découvrez comment attribuer des rôles.

    Activer les API

  2. Pour obtenir les autorisations nécessaires pour créer les ressources, demandez à votre administrateur de vous accorder les rôles IAM (Identity and Access Management) requis sur votre projet.

  3. Table Delta Lake valide stockée dans un bucket Cloud Storage. Le répertoire de tables doit contenir vos fichiers de données et le répertoire du journal des transactions _delta_log/. De plus, le tableau doit répondre aux exigences suivantes :

    • La propriété delta.enableChangeDataFeed doit être activée (true) sur la table Delta Lake.
    • L'ID ou le code temporel du commit de début configuré pour le pipeline CDC doit être postérieur à l'activation de la propriété delta.enableChangeDataFeed.
  4. Un catalogue Iceberg Lakehouse existant pour recevoir les données synchronisées. Si l'espace de noms cible n'existe pas, le pipeline CDC le crée automatiquement. Si la table Lakehouse cible n'existe pas, le pipeline CDC la crée automatiquement si vous fournissez une liste valide de colonnes d'égalité.

Compatibilité et limites

La migration CDC par lot planifiée de Delta Lake vers Lakehouse présente les considérations suivantes :

  • Exigences concernant le SDK : nécessite le SDK Apache Beam version 2.77.0 ou ultérieure.
  • Création de la table de destination : la création de la table Lakehouse de destination est acceptée. Si la table cible n'existe pas, le pipeline CDC la crée automatiquement si vous fournissez une liste valide de colonnes d'égalité (equality_columns). Seul un catalogue Lakehouse existant est requis. Le pipeline CDC crée automatiquement l'espace de noms s'il n'existe pas.
  • Évolution du schéma : les modifications automatiques du schéma ne sont pas acceptées. Si le schéma de la table Delta Lake source change, vous devez mettre à jour manuellement le schéma de la table Lakehouse cible pour qu'il corresponde avant d'exécuter le job de synchronisation CDC.
  • Mode d'exécution par lot : la CDC Delta Lake fonctionne en mode par lot. Le pipeline extrait les journaux de transactions une fois par exécution de job entre les limites de validation de début et de fin spécifiées (ID de validation ou code temporel). Si aucune limite de commit de fin (end_version ou end_timestamp) n'est spécifiée, le pipeline lit jusqu'à la dernière version du commit.
  • Stockage source compatible : les données sources doivent être une table Delta Lake valide stockée dans Cloud Storage (gs://). Amazon S3 n'est pas compatible avec les sources de tables Delta Lake.
  • Tables basées sur le catalogue : les tables Delta Lake basées sur le catalogue (par exemple, Unity Catalog) ne sont pas acceptées.
  • Versions de Delta Lake non compatibles : la version 1.2.1 de Delta Lake ou les versions antérieures ne sont pas compatibles, car le flux de données de modification (CDF, Change Data Feed) n'est pas disponible pour ces versions.
  • Flux de données modifiées automatiques : ce flux n'est pas compatible, car il n'écrit pas de fichiers journaux des modifications CDF.

Créer un pipeline CDC de manière programmatique

Pour créer et exécuter de manière programmatique un pipeline de migration par lot CDC Delta Lake, utilisez la transformation Managed.read dans votre pipeline Apache Beam.

Java

Ajoutez les dépendances suivantes à votre fichier 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>

L'extrait de code Java suivant montre comment configurer et exécuter un pipeline CDC par lot de Delta Lake vers 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();
  }
}

Remplacez les éléments suivants :

  • BUCKET_NAME : nom du bucket Cloud Storage contenant la table Delta Lake.
  • TABLE_NAME : nom du répertoire de la table Delta Lake source.
  • START_COMMIT_ID : numéro de version du commit de départ pour la lecture des événements CDC Delta Lake. Vous devez indiquer start_version ou start_timestamp. L'ID ou le code temporel du commit de début doivent être postérieurs à l'activation de la propriété delta.enableChangeDataFeed dans la table source.
  • END_COMMIT_ID : numéro de version de commit de fin pour la lecture des événements CDC Delta Lake. Nous vous recommandons de spécifier une limite de fin (end_version ou end_timestamp) qui correspond au type de limite utilisé pour le début (start_version ou start_timestamp). Si une limite de fin est omise, le pipeline lit jusqu'à la dernière version de commit.
  • WAREHOUSE_BUCKET : nom du bucket Cloud Storage utilisé comme entrepôt de catalogue Lakehouse.
  • PROJECT_ID : ID de votre projet Google Cloud .
  • TARGET_NAMESPACE : espace de noms de la table Lakehouse cible.
  • TARGET_TABLE : nom de la table Lakehouse cible.
  • PRIMARY_KEY_COLUMN : nom de la colonne de clé primaire utilisée pour identifier les lignes pour les mises à jour et les suppressions CDC.

Examiner le résultat du job

Vérifiez que les données CDC ont bien été fusionnées dans votre table Lakehouse :

  1. Dans la console Google Cloud , accédez à la page Studio de BigQuery.

    Accéder à BigQuery

  2. Dans l'éditeur de requête, exécutez une requête SQL pour vérifier les données synchronisées :

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

    Remplacez les éléments suivants :

    • PROJECT_ID : ID de votre projet Google Cloud .
    • CATALOG : nom de votre catalogue Lakehouse.
    • NAMESPACE : espace de noms de votre table Lakehouse.
    • TABLE_NAME : nom de votre table Lakehouse cible.
  3. Cliquez sur Exécuter et vérifiez les résultats.

Étapes suivantes