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,UPDATEetDELETE), 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 :
- Activez la propriété
delta.enableChangeDataFeed(true) sur la table Delta Lake d'entrée. - 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. - 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 :
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.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.
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.enableChangeDataFeeddoit ê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.
- La propriété
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_versionouend_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 indiquerstart_versionoustart_timestamp. L'ID ou le code temporel du commit de début doivent être postérieurs à l'activation de la propriétédelta.enableChangeDataFeeddans 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_versionouend_timestamp) qui correspond au type de limite utilisé pour le début (start_versionoustart_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 :
Dans la console Google Cloud , accédez à la page Studio de BigQuery.
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.
Cliquez sur Exécuter et vérifiez les résultats.
Étapes suivantes
- En savoir plus sur l'importation de tables Delta Lake dans Lakehouse
- Découvrez l'ingestion de la capture des données modifiées dans Lakehouse.
- En savoir plus sur les E/S gérées dans Dataflow