Aller au contenu principal

Stocker les modifications Postgres dans le lakehouse

remarque

La fonctionnalité Change Data Feed de Lakebase est en aperçu public.

Configurez le flux de données de modification (CDF) de Lakebase sur une table Postgres, puis observez les modifications au niveau des lignes apparaître dans la table Delta de destination.

**Étapes :**Activer la capture des changements → ② Start le flux → ③ Suivre une ligne dans le lakehouse → ④ Modifier la ligne, la voir passer

remarque

Ceci est un démarrage rapide. Pour une documentation complète, consultez Flux de données de changement Lakebase.

Avant de commencer

  • Assurez-vous d'avoir terminé la section Obtenir une base de données Postgres. Vous avez besoin d'un projet Lakebase avec la table d'échantillons playing_with_lakebase.
  • Un catalogue et un schéma Unity Catalog où vous disposez de l'autorisation CREATE TABLE.

Étape 1 : Activer la capture de changement

Postgres a besoin de données de ligne complètes dans le log d'écriture anticipée pour que CDF fonctionne. Définir l'identité de la réplique sur « full » indique à Postgres d'enregistrer l'état de l'ancienne et de la nouvelle ligne pour chaque modification.

Dans l'Éditeur SQL Lakebase, exécutez :

SQL
ALTER TABLE playing_with_lakebase REPLICA IDENTITY FULL;

En savoir plus : Définir l'identité du réplica sur toutes les tables d'un schéma et l'appliquer automatiquement aux nouvelles tables

Étape 2 : start le flux

Lakebase CDF est configuré au niveau du schéma. Chaque table actuelle et future dans le schéma source est incluse automatiquement, de sorte que vous ne sélectionnez pas de tables individuelles.

Depuis votre Branch de production, ouvrez Vue d'ensemble de la Branch en cliquant sur le nom de la Branch dans le fil d'Ariane supérieur, puis ouvrez la Lakebase CDF tab et cliquez sur Start . Choisissez public comme schéma source, puis sélectionnez un catalogue et un schéma de destination Unity Catalog. L'instantané initial commence immédiatement, et lb_playing_with_lakebase_history apparaît comme une table Delta dans votre destination.

start la boîte de dialogue avec la sélection de la source et de la destination.

En savoir plus : start the change data feed

Étape 3 : Suivre une ligne dans le lakehouse

Sélectionnez une ligne depuis Lakebase. Jetez un œil à la ligne id=2:

SQL
SELECT * FROM playing_with_lakebase WHERE id = 2;

Maintenant, trouvez la même ligne dans la table d'historique Delta. Basculez vers un Databricks SQL warehouse ou un Notebook et exécutez :

SQL
SELECT * FROM <catalog>.<schema>.lb_playing_with_lakebase_history
WHERE id = 2;

Remplacez <catalog> et <schema> par la destination que vous avez choisie à l'Étape 2. Vous verrez la ligne id=2 avec les mêmes name et value que dans Lakebase, plus des colonnes supplémentaires. L'instantané initial a écrit chaque ligne existante dans Delta en tant qu'événement insert, ce que cette ligne représente.

Ces colonnes supplémentaires décrivent le type d'événement que représente chaque ligne (_pg_change_type), quand il s'est produit (_timestamp) et les informations de tri Postgres (_pg_lsn, _pg_xid).

En savoir plus : Schéma de la table de destination | Mappage des types de données

Étape 4 : modifiez la ligne et observez-la circuler.

De retour dans l'Éditeur SQL Lakebase, mettez à jour la ligne id=2:

SQL
UPDATE playing_with_lakebase SET value = 55.5 WHERE id = 2;

Patientez quelques secondes pour que la modification apparaisse dans le flux, puis interrogez à nouveau la table d'historique :

SQL
SELECT id, value, _pg_change_type, _timestamp
FROM <catalog>.<schema>.lb_playing_with_lakebase_history
WHERE id = 2
ORDER BY _pg_lsn DESC;

Table d&#39;historique Delta montrant trois lignes pour id=2 : update_preimage, update_postimage et insert

La ligne id=2 apparaît maintenant trois fois : l'original insert, une update_preimage avec l'ancienne valeur et une update_postimage avec la nouvelle valeur. Chaque modification de la ligne devient une nouvelle ligne d'historique, vous disposez donc toujours d'une piste d'audit complète. Les suppressions fonctionnent de la même manière, en ajoutant une ligne avec _pg_change_type = 'delete'.

En savoir plus : Modèles de modification courants | Construire des pipelines en aval

Étapes suivantes