Il peut être difficile de synchroniser les tables Apache Iceberg Lakehouse cibles avec les tables sources Apache Iceberg 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 capture des données modifiées (CDC) introduit une surcharge opérationnelle complexe.
En utilisant des pipelines CDC par lot ou par flux Dataflow, vous pouvez synchroniser en continu ou périodiquement vos tables sources Apache Iceberg avec les tables Lakehouse Apache Iceberg. Au lieu de recharger des tables entières, les pipelines CDC extraient les mises à jour incrémentielles directement à partir des journaux des modifications Apache Iceberg de votre source et les appliquent à vos tables cibles.
Le pipeline CDC offre les fonctionnalités suivantes :
- Synchronisation incrémentielle économique : ne traite que les données modifiées (opérations
INSERT,UPDATEetDELETE), ce qui réduit la surcharge de calcul et la latence des données. - Modes d'exécution flexibles : compatibles avec les modes par lot et par flux pour répondre à vos besoins commerciaux.
Le pipeline CDC est compatible avec deux modes d'exécution :
- Mode CDC par lot : lit périodiquement les mises à jour incrémentielles de votre table Apache Iceberg source entre les limites d'instantané de début et de fin spécifiées.
- Mode CDC de streaming : interroge en continu votre table Apache Iceberg source pour détecter les nouveaux commits de modification à une fréquence configurable (par défaut, 1 minute, minimum 1 seconde) et les applique à votre table Lakehouse en temps quasi réel.
Avant de commencer
Pour configurer la migration CDC d'Apache Iceberg vers Lakehouse, 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 Apache Iceberg source valide stockée dans Cloud Storage ou enregistrée dans un catalogue REST ou un metastore Iceberg.
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 d'Apache Iceberg 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, à condition que vous fournissiez 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 Apache Iceberg 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.
- Modes d'exécution :
- CDC par lot : extrait les modifications du journal des transactions une fois par exécution du job, entre les limites d'instantané de début et de fin spécifiées (ID ou code temporel de l'instantané).
- CDC en flux continu : interroge en continu la table Apache Iceberg source pour obtenir de nouveaux journaux de commit. La fréquence d'interrogation est configurable (par défaut, elle est de 1 minute, avec un minimum de 1 seconde).
Créer un pipeline CDC à l'aide de Java
Pour créer et exécuter de manière programmatique un pipeline de migration par lot ou par flux CDC Apache Iceberg à l'aide du SDK Java Apache Beam, utilisez la transformation Managed.read avec Managed.ICEBERG_CDC.
Les exemples Java suivants montrent comment configurer et exécuter des pipelines CDC Apache Iceberg pour les charges de travail par lot et de streaming.
Pipeline CDC par lot
L'exemple Java suivant configure un pipeline par lot qui lit les événements CDC Apache Iceberg dans des limites de début et de fin d'instantané spécifiées :
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 à partir d'une table Iceberg 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.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();
}
}
Remplacez les éléments suivants :
SOURCE_CATALOG_TYPE: type du catalogue Iceberg source (par exemple,hadoopouhive).SOURCE_WAREHOUSE_LOCATION: emplacement de l'entrepôt du catalogue Iceberg source (par exemple,s3://source-warehouseougs://source-warehouse).SOURCE_NAMESPACE: espace de noms de la table Iceberg source.SOURCE_TABLE: nom de la table Iceberg source.START_SNAPSHOT_ID: ID de l'instantané de départ pour la lecture CDC par lot.END_SNAPSHOT_ID: ID de l'instantané de fin pour la lecture CDC par lot.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.
Pipeline CDC de flux de données
L'exemple Java suivant configure un pipeline de streaming qui interroge en continu la table source Apache Iceberg pour détecter les nouvelles modifications CDC :
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 de streaming à partir d'une table Iceberg 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.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();
}
}
Remplacez les éléments suivants :
SOURCE_CATALOG_TYPE: type du catalogue Iceberg source (par exemple,hadoopouhive).SOURCE_WAREHOUSE_LOCATION: emplacement de l'entrepôt du catalogue Iceberg source (par exemple,s3://source-warehouseougs://source-warehouse).SOURCE_NAMESPACE: espace de noms de la table Iceberg source.SOURCE_TABLE: nom de la table Iceberg source.POLL_INTERVAL_SECONDS: fréquence d'interrogation continue en secondes (60par défaut,1minimum).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.
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 Iceberg 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