Aller au contenu principal

Lakebase Change Data Feed

remarque

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

Qu'est-ce que le flux de données de modification Lakebase ?

Lakebase introduit un flux de données de modification (CDF) natif, déverrouillant vos données opérationnelles pour les pipelines, les modèles et les applications en aval. Chaque insertion, mise à jour et suppression sur une table Postgres Lakebase est capturée à partir du log de pré-écriture et stockée sous forme de nouvelle ligne dans une table Delta gérée par Unity Catalog, traitée par lots et vidée toutes les ~ 15 secondes. L'historique des modifications est stocké dans un format ouvert que n'importe quel moteur de compute peut lire.

Les tables de destination suivent la même forme que le Delta Change Data Feed: chaque ligne contient un _pg_change_type, un LSN, un ID de transaction et un Timestamp. Les changements opérationnels deviennent une source de premier ordre pour l'ETL, l'audit et les consommateurs en aval — sans la mise en place d'une pile CDC externe.

Flux de données Lakebase CDF de Postgres via wal2delta vers les tables Delta dans Unity Catalog.

Cas d'utilisation

Lakebase CDF apporte les données opérationnelles au lakehouse afin que les pipelines et applications en aval puissent réagir aux changements au fur et à mesure qu'ils se produisent.

Cas d'usage

Description

Pipelines ETL

Utilisez Lakebase comme source bronze pour les pipelines en médaillon. Créez des pipelines Lakeflow incrémentiels ou des Jobs Spark Structured Streaming sur le flux de changements et mettez à jour les tables silver et Gold en aval.

Journaux d'audit

Maintenez un historique complet et interrogeable de chaque insertion, mise à jour et suppression sur une table Lakebase pour la conformité et l'analyse forensique. L'historique est Delta immuable.

Systèmes externes

Stockez les données de modification Lakebase dans un format ouvert que tout moteur peut consommer. Puisque la destination est une table Delta dans Unity Catalog, les systèmes externes et les lecteurs non Databricks peuvent accéder directement au flux.

Cas d'usage

Description

Pipelines ETL

Utilisez Lakebase comme source bronze pour les pipelines en médaillon. Créez des pipelines Lakeflow incrémentiels ou des Jobs Spark Structured Streaming sur le flux de changements et mettez à jour les tables silver et Gold en aval.

Journaux d'audit

Maintenez un historique complet et interrogeable de chaque insertion, mise à jour et suppression sur une table Lakebase pour la conformité et l'analyse forensique. L'historique est Delta immuable.

Systèmes externes

Stockez les données de modification Lakebase dans un format ouvert que tout moteur peut consommer. Puisque la destination est une table Delta dans Unity Catalog, les systèmes externes et les lecteurs non Databricks peuvent accéder directement au flux.

Activez cet aperçu

Un administrateur de workspace doit activer l'aperçu Lakebase Change Data Feed à partir de la page d'aperçus du workspace.

Exigences

  • Lakebase Autoscaling : un projet Lakebase Autoscaling exécutant Postgres 17.
  • Base de données source : Les tables doivent résider dans la base de données databricks_postgres dans Lakebase. Chaque projet est créé avec cette base de données default. Ceci est une limitation connue.
  • Unity Catalog : l'identité configurant le CDF a besoin de USE CATALOG , USE SCHEMA et CREATE TABLE sur le catalogue et le schéma de destination. Consultez Accorder des autorisations sur un objet.
  • Stockage default : les catalogues de destination configurés avec le stockage default ne sont pas pris en charge.
  • Projet Lakebase : Votre rôle Postgres nécessite les autorisations CAN MANAGE sur le projet Lakebase. Les propriétaires de projet disposent de l'autorisation CAN MANAGE by default. Consultez Gérer les autorisations du projet.
  • Types de données : voir le mappage des types de données. Les types sans équivalent Delta direct sont stockés en tant que STRING.

Configurer Lakebase CDF

Pour commencer, définissez l'identité complète de la réplique sur les tables que vous souhaitez dans le flux (Étape 1), puis start CDF dans l'application Lakebase (Étape 2). Vos données apparaissent sous forme de lb_<table_name>_history tables Delta dans le catalogue et le schéma Unity Catalog que vous choisissez.

Étape 1 : Définir l’identité complète de la réplique

Pour qu'une table Lakebase participe à CDF, elle doit avoir REPLICA IDENTITY FULL défini. Par default, Postgres logs uniquement la clé primaire lorsqu’une ligne est mise à jour ou supprimée. La définition de l'identité complète indique à Postgres d'enregistrer l'état des lignes avant et après dans le journal de préécriture (write-ahead log), ce dont CDF a besoin pour construire un historique complet des modifications.

Vous pouvez exécuter ces commandes dans l'éditeur SQL Lakebase ou tout client Postgres.

SQL
ALTER TABLE <table_name> REPLICA IDENTITY FULL;

Vérifiez quelles tables ont une identité de réplique définie.

Pour savoir quelles tables d'un schéma ont une identité de réplique configurée, exécutez :

SQL
SELECT n.nspname AS table_schema,
c.relname AS table_name,
CASE c.relreplident
WHEN 'd' THEN 'default'
WHEN 'n' THEN 'nothing'
WHEN 'f' THEN 'full'
WHEN 'i' THEN 'index'
END AS replica_identity
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relkind = 'r'
AND n.nspname = 'public'
ORDER BY n.nspname, c.relname;

Seules les lignes avec replica_identity = 'full' sont prêtes pour CDF.

Étape 2 : start the change data feed

Lakebase CDF est configuré au niveau du schéma. Une fois start, chaque table actuelle et future dans le schéma source est incluse dans le flux.

  1. Dans votre Workspace Databricks, ouvrez Lakebase Postgres depuis le sélecteur d'applications (en haut à droite).

  2. Sélectionnez votre projet Lakebase et la Branch que vous souhaitez utiliser (par exemple, **production** ou **main**).

  3. Ouvrez l' Aperçu de la Branch en cliquant sur le nom de la Branch dans le fil d'Ariane supérieur, puis cliquez sur l'onglet Lakebase CDF .

  4. Click start .

  5. Dans la boîte de dialogue de configuration :

    • Base de données : La valeur par default est databricks_postgres.
    • Schéma : Sélectionnez le schéma Postgres source.
    • Vers le catalogue : Sélectionnez le catalogue Unity Catalog de destination.
    • Schéma : sélectionnez le schéma Unity Catalog de destination.
  6. Cliquez sur Start pour démarrer le flux.

Aperçu de Branch avec l&#39;onglet Lakebase CDF affichant Start et configuration de schéma.

Les tables apparaissent dans la destination en tant que lb_<table_name>_history. Pour les trouver, ouvrez Catalogue dans la barre latérale, accédez au catalogue et au schéma de destination, et ouvrez l'onglet tab .

Le tab **Lakebase CDF** contient deux sous-tabs :

Les sous-onglets affichent la correspondance et la progression par table.

  • Schémas : Répertorie chaque schéma source, son catalogue et son schéma de destination dans Unity Catalog, ainsi qu'un statut.
  • Tables : Liste chaque table source, sa table de destination lb_<table_name>_history, son état (Streaming ou Snapshotting), le LSN validé (la distance parcourue par le flux dans Delta, affichée comme - pendant l'instantané initial) et la dernière mise à jour (la dernière fois que la table a reçu des modifications).

Vous pouvez également inspecter l'état du flux depuis Postgres en exécutant ceci dans l'éditeur SQL Lakebase:

SQL
SELECT * FROM wal2delta.tables;

Le résultat inclut table_oid, status (STREAMING ou SNAPSHOTTING), committed_lsn et last_write_time par table.

info

Qu'est-ce que wal2delta ? Lakebase CDF est alimenté par l'extension Postgres **wal2delta**, qui s'exécute à l'intérieur du compute Lakebase. Il utilise le décodage logique pour capturer les modifications du journal de transactions (WAL) et les écrit dans les tables Delta d'Unity Catalog.

Schéma de table de destination

CDF écrit une table Delta par table source, nommée lb_<table_name>_history dans votre catalogue et schéma de destination. En plus de vos colonnes source, chaque ligne contient ces colonnes système :

Colonne

Type

Description

_pg_change_type

Texte

Type d’opération : insert, delete, update_preimage ou update_postimage.

_pg_lsn

BIGINT

Numéro de séquence des Logs Postgres.

_pg_xid

INTEGER

ID de transaction Postgres.

_timestamp

Horodatage

Timestamp auquel la modification a été traitée (sans fuseau horaire).

_sort_by

BIGINT

Clé de tri monotone utilisée pour ordonner toutes les modifications.

Colonne

Type

Description

_pg_change_type

Texte

Type d’opération : insert, delete, update_preimage ou update_postimage.

_pg_lsn

BIGINT

Numéro de séquence des Logs Postgres.

_pg_xid

INTEGER

ID de transaction Postgres.

_timestamp

Horodatage

Timestamp auquel la modification a été traitée (sans fuseau horaire).

_sort_by

BIGINT

Clé de tri monotone utilisée pour ordonner toutes les modifications.

Modèles de changement courants

  • Instantanné initial : La première fois que CDF s’exécute sur une table Lakebase existante, chaque ligne existante est écrite avec _pg_change_type = 'insert'.
  • Mises à jour : Une mise à jour produit deux lignes : une avec _pg_change_type = 'update_preimage' (ancienne ligne) et une avec _pg_change_type = 'update_postimage' (nouvelle ligne).
  • Suppressions : Une suppression produit une ligne avec _pg_change_type = 'delete'.

Ce sont les mêmes événements de changement que le Delta Change Data Feed, ainsi les mêmes modèles en aval s'appliquent.

Comportement opérationnel

  • Conflits de noms : si deux tables sources devaient correspondre au même nom de destination (par exemple, sales.users et marketing.users correspondant toutes deux à lb_users_history), CDF écrit la première dans lb_users_history et ajoute automatiquement un suffixe à la seconde dans lb_users_history_1. Vous pouvez renommer n'importe quelle table de destination dans Unity Catalog et le flux continue de fonctionner.
  • Portée au niveau du schéma : Lorsque vous start CDF sur un schéma Lakebase, chaque table actuelle et future de ce schéma est incluse. Les tables vides sont ignorées — une table doit avoir au moins une ligne pour apparaître dans la destination.
  • Tables sources abandonnées : Si vous supprimez une table dans Lakebase, la table Delta de destination dans Unity Catalog est préservée.

Créer des pipelines en aval

Lakebase CDF est conçu pour les pipelines en aval qui réagissent aux changements opérationnels. Les modèles ci-dessous montrent trois façons de consommer le flux, classées de la plus simple à la plus flexible.

**Scénario d'exemple.** Une application d'e-commerce enregistre les commandes dans une table Postgres orders, chaque ligne comportant un item_id et un quantity. L'équipe logistique a besoin des niveaux d'inventaire en temps réel. Avec CDF, chaque modification apportée à orders est stockée dans la table Delta lb_orders_history dans Unity Catalog. Les pipelines en aval lisent ce flux de modifications et mettent à jour une table inventory_levels chaque fois qu'une commande est passée, modifiée ou annulée.

Calculer l'inventaire actuel avec une vue matérialisée

Le modèle le plus simple est une vue matérialisée SQL sur la table d'historique. La MV se refresh de manière incrémentielle à mesure que de nouveaux événements de changement arrivent, et les consommateurs en aval l'interrogent comme n'importe quelle autre table.

SQL
CREATE MATERIALIZED VIEW inventory_levels AS
SELECT
item_id,
SUM(
CASE
-- New orders (and the "new half" of updates) decrement inventory
WHEN _pg_change_type IN ('insert', 'update_postimage') THEN -quantity
-- Cancellations (and the "old half" of updates) restore inventory
WHEN _pg_change_type IN ('delete', 'update_preimage') THEN quantity
ELSE 0
END
) AS current_inventory,
MAX(_timestamp) AS last_transaction_ts,
MAX(_pg_lsn) AS last_lsn
FROM lb_orders_history
GROUP BY item_id;

Les deux lignes produites pour chaque mise à jour s'annulent mutuellement, à l'exception de la variation nette, de sorte que la somme courante reste correcte lorsque les commandes sont modifiées.

Stream changes with Spark Declarative Pipelines

Pour une architecture en médaillon structurée, utilisez les LakeFlow Pipelines pour déclarer les tables bronze, silver et Gold. Les LakeFlow pipelines les exécutent comme un pipeline connecté avec des points de contrôle et une gestion des dépendances pris en charge pour vous.

Python
import dlt
from pyspark.sql import functions as F

@dlt.table
def inventory_adjustments():
return (
spark.readStream.table("<catalog>.<schema>.lb_orders_history")
.withColumn(
"delta",
F.when(F.col("_pg_change_type").isin("insert", "update_postimage"), -F.col("quantity"))
.when(F.col("_pg_change_type").isin("delete", "update_preimage"), F.col("quantity"))
.otherwise(0),
)
.select("item_id", "delta", "_timestamp")
)

@dlt.expect_or_drop("non_negative_stock", "on_hand >= 0")
@dlt.table
def inventory_levels():
return (
spark.read.table("LIVE.inventory_adjustments")
.groupBy("item_id")
.agg(F.sum("delta").alias("on_hand"))
)

inventory_adjustments lit lb_orders_history de manière incrémentielle avec readStream et produit un delta par événement. inventory_levels agrège par item_id pour compute le stock actuel. L'attente supprime les lignes qui feraient passer le stock en négatif, signalant un bug en amont.

Pour une présentation complète de bout en bout, consultez Didacticiel : Créer un pipeline ETL à l'aide de la capture de données modifiées.

Traitement personnalisé avec Spark Structured Streaming

Lorsque vous avez besoin d'un contrôle total — par exemple, des Merge personnalisés, des effets secondaires ou plusieurs récepteurs —, lisez directement la table d'historique avec Spark Structured Streaming et utilisez foreachBatch pour écrire vers votre destination.

Python
from pyspark.sql import functions as F
from delta.tables import DeltaTable

def update_inventory(batch_df, batch_id):
deltas = (
batch_df
.withColumn(
"delta",
F.when(F.col("_pg_change_type").isin("insert", "update_postimage"), -F.col("quantity"))
.when(F.col("_pg_change_type").isin("delete", "update_preimage"), F.col("quantity"))
.otherwise(0),
)
.groupBy("item_id")
.agg(F.sum("delta").alias("delta"))
)

target = DeltaTable.forName(spark, "<catalog>.<schema>.inventory_levels")
(target.alias("t")
.merge(deltas.alias("s"), "t.item_id = s.item_id")
.whenMatchedUpdate(set={&quot;on_hand&quot;: F.expr(&quot;t.on_hand + s.delta&quot;)})
.whenNotMatchedInsert(values={&quot;item_id&quot;: &quot;s.item_id&quot;, &quot;on_hand&quot;: &quot;s.delta&quot;})
.execute())

(spark.readStream.table("<catalog>.<schema>.lb_orders_history")
.writeStream
.foreachBatch(update_inventory)
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/checkpoints/inventory_levels")
.start())

Chaque microbatch agrège les événements de modification par item_id et Merge les deltas nets dans inventory_levels.

Incrémentiel par conception. Chaque table lb_<table_name>_history est une table Delta en ajout seul. Chaque modification de la source est enregistrée comme une nouvelle ligne, avec _pg_change_type marquant l'opération. Les vues matérialisées Databricks SQL, les flux LakeFlow Pipelines et les jobs Spark Structured Streaming traitent tous les nouvelles lignes de manière incrémentielle à partir du log des transactions Delta, de sorte que les pipelines en aval ne traitent que ce qui a changé. Il n’est pas nécessaire d’activer le Flux de données de modification Delta sur la table d’historique, car la sémantique de modification est déjà encodée dans les données de ligne.

Mappage des types de données

CDF prend en charge la plupart des types primitifs PostgreSQL standard. Les types sans équivalent Delta direct sont stockés en tant que STRING.

Type PostgreSQL

Type de Databricks Delta

Notes

Booléen

Booléen

INT, SMALLINT, BIGINT

INT, SMALLINT, BIGINT

TEXT, VARCHAR, CHAR

CHAÎNE

JSONB

CHAÎNE

Stocké sous forme de chaîne JSON.

Énumération

CHAÎNE

Stocké en tant que libellé d'énumération.

NUMÉRIQUE / DÉCIMAL

DECIMAL ou STRING

Utilise la précision/l'échelle source si possible. Effectue une mise à l'échelle sans perte pour les valeurs de précision/échelle incompatibles. Revient au type STRING lorsque la précision dépasse 38 ou lorsque la précision/l'échelle ne sont pas définies (NUMERIC illimité). Toutes les colonnes NUMERIC/DECIMAL sont nullables, car les valeurs NaN sont mappées à NULL. Consultez les types numériques PostgreSQL.

Date

Date

Horodatage

TIMESTAMP_NTZ

TIMESTAMPTZ

Horodatage

Présentation libre, DOUBLE

Présentation libre, DOUBLE

Type PostgreSQL

Type de Databricks Delta

Notes

Booléen

Booléen

INT, SMALLINT, BIGINT

INT, SMALLINT, BIGINT

TEXT, VARCHAR, CHAR

CHAÎNE

JSONB

CHAÎNE

Stocké sous forme de chaîne JSON.

Énumération

CHAÎNE

Stocké en tant que libellé d'énumération.

NUMÉRIQUE / DÉCIMAL

DECIMAL ou STRING

Utilise la précision/l'échelle source si possible. Effectue une mise à l'échelle sans perte pour les valeurs de précision/échelle incompatibles. Revient au type STRING lorsque la précision dépasse 38 ou lorsque la précision/l'échelle ne sont pas définies (NUMERIC illimité). Toutes les colonnes NUMERIC/DECIMAL sont nullables, car les valeurs NaN sont mappées à NULL. Consultez les types numériques PostgreSQL.

Date

Date

Horodatage

TIMESTAMP_NTZ

TIMESTAMPTZ

Horodatage

Présentation libre, DOUBLE

Présentation libre, DOUBLE

Types stockés sous forme de STRING :

  • Géographie/Géométrie (PostGIS) : types de l’extension PostGIS (par exemple, geometry, geography).
  • Vecteur (pgvector) : le type vector de l'extension pgvector.
  • Types composites/structurés : Types personnalisés définis avec CREATE TYPE ... AS (field_name type, ...). Ce sont des types de lignes avec des champs nommés.
  • Map : types clé-valeur de type Map tels que hstore (issu de l'extension hstore). Postgres ne possède pas de type Map intégré. hstore est la méthode courante pour stocker des paires clé-valeur dans une colonne.

Gestion des changements de schéma

  • Le renommage d'une table dans Postgres (par exemple, ALTER TABLE users RENAME TO customers) permet au flux de continuer. Le nom de la table Delta de destination ne change pas — il reste lb_users_history.
  • Les modifications de schéma (ajout d'une colonne, suppression d'une colonne ou modification du type de données d'une colonne) trigger une nouvelle capture instantanée de la table concernée. CDF relit la table entière de Postgres et la réécrit dans la table Delta de destination.

Désactiver Lakebase CDF

La désactivation du CDF arrête le flux pour tous les schémas Lakebase du projet.

  1. Dans votre Workspace Databricks, ouvrez Lakebase Postgres depuis le sélecteur d'applications (en haut à droite).
  2. Sélectionnez votre projet Lakebase et la branch où vous avez configuré CDF.
  3. Ouvrez l' Aperçu de la Branch en cliquant sur le nom de la Branch dans le fil d'Ariane supérieur, puis cliquez sur l'onglet Lakebase CDF .
  4. Cliquez sur Désactiver . Dans la boîte de dialogue de confirmation, examinez l'avertissement selon lequel les modifications cesseront d'être transmises aux tables Delta, puis cliquez de nouveau sur **Désactiver** pour confirmer.

La désactivation du CDF ne redémarre pas votre compute.

Limitations et dépannage

Vous pouvez voir le statut par table (instantané, ignoré ou streaming) dans l’onglet **Lakebase CDF**, ou en exécutant ceci dans Lakebase :

SQL
SELECT * FROM wal2delta.tables;

Raisons courantes pour lesquelles une table n'apparaît pas dans le flux :

  • REPLICA IDENTITY FULL non défini : Exécutez ALTER TABLE <table_name> REPLICA IDENTITY FULL; pour la table. Voir Étape 1 : Définir l'identité complète du réplica.
  • Tables partitionnées : Les tables partitionnées Lakebase ne sont pas prises en charge. Un schéma qui contient des tables partitionnées entraîne l'échec de ces tables.
  • **Tables vides :** Une table avec zéro ligne est ignorée jusqu’à ce qu’au moins une ligne existe.

Étapes suivantes