Change Data Feed Lakebase

Note

La fonctionnalité de flux de données de modification de Lakebase est disponible en préversion publique.

Qu’est-ce que le flux de modifications de données de Lakebase ?

Lakebase introduit un flux natif de suivi des modifications des données (CDF), rendant vos données opérationnelles exploitables par les pipelines, les modèles et les applications en aval. Chaque insertion, mise à jour et suppression dans une table Postgres de Lakebase est capturée à partir du journal de transactions (write-ahead log) et enregistrée sous la forme d’une nouvelle ligne dans une table Delta gérée par Unity Catalog, regroupée par lots et écrite environ toutes les 15 secondes. L’historique des modifications est stocké dans un format ouvert que n’importe quel moteur de calcul peut lire.

Les tables de destination suivent la même forme que le flux de données de modification delta : chaque ligne porte un _pg_change_typeLSN, un ID de transaction et un horodatage. Les modifications opérationnelles deviennent une source à part entière pour l’ETL, l’audit et les consommateurs en aval, sans avoir à déployer une stack CDC externe.

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

Cas d’utilisation

Lakebase CDF apporte des données opérationnelles dans le lakehouse afin que les pipelines et les applications en aval puissent réagir aux modifications à mesure qu’elles se produisent.

Cas d’utilisation Description
Pipelines ETL Utilisez Lakebase comme source bronze pour les pipelines en médaillon. Créez des pipelines Lakeflow incrémentiels ou des travaux Spark Structured Streaming à partir du flux des modifications, puis mettez à jour les tables Silver et Gold en aval.
Journaux d’audit Conservez un historique complet et interrogeable de chaque insertion, mise à jour et suppression sur une table Lakebase pour la conformité et les analyses. L’historique Delta est immuable.
Systèmes externes Store Lakebase modifie les données dans un format ouvert que n’importe quel moteur peut consommer. Étant donné que la destination est une table Delta dans le catalogue Unity, les systèmes externes et les lecteurs non Databricks peuvent accéder directement au flux.

Activer cette préversion

Un administrateur d’espace de travail doit activer l’aperçu du flux de données modifiées Lakebase à partir de la page Aperçus de l’espace de travail.

Requirements

  • Mise à l’échelle automatique Lakebase : Projet de mise à l’échelle automatique Lakebase exécutant Postgres 17.
  • Base de données source : Les tables doivent résider dans la databricks_postgres base de données dans Lakebase. Chaque projet est créé avec cette base de données par défaut. Il s’agit d'une limitation connue.
  • Unity Catalog: L’identité qui configure CDF a besoin de USE CATALOG, USE SCHEMA et CREATE TABLE sur le catalogue et le schéma de destination. Voir Accorder des autorisations sur un objet.
  • Stockage par défaut : Les catalogues de destination configurés avec le stockage par défaut ne sont pas pris en charge.
  • Projet Lakebase : Votre rôle Postgres nécessite des autorisations CAN MANAGE sur le projet Lakebase. Les propriétaires du projet ont CAN MANAGE par défaut. Consultez Gérer les autorisations de projet.
  • Types de données : Consultez le mappage de type de données. Les types sans équivalent Delta direct sont stockés en tant que CHAÎNE.

Configurer Lakebase CDF

Pour commencer, définissez REPLICA IDENTITY sur FULL sur les tables que vous souhaitez inclure dans le flux (étape 1), puis démarrez le CDF dans l’application Lakebase (étape 2). Vos données apparaissent sous forme lb_<table_name>_history de tables Delta dans le catalogue Unity et le schéma que vous choisissez.

Étape 1 : Définir replica identity full

Pour qu’une table Lakebase participe à la CDF, elle doit avoir REPLICA IDENTITY FULL défini. Par défaut, Postgres journalise 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 à la fois l’état de ligne avant et après dans le journal des transactions, dont le CDF a besoin pour générer un historique complet des modifications.

Vous pouvez exécuter ces commandes dans l’Éditeur SQL Lakebase ou n’importe quel client Postgres.

Table unique

ALTER TABLE <table_name> REPLICA IDENTITY FULL;

Toutes les tables existantes dans un schéma

Pour définir l’identité du réplica sur chaque table existante d’un schéma (public dans cet exemple), exécutez :

DO $$
DECLARE r record;
BEGIN
  FOR r IN
    SELECT table_schema, table_name
    FROM information_schema.tables
    WHERE table_schema = 'public'
      AND table_type = 'BASE TABLE'
  LOOP
    EXECUTE format(
      'ALTER TABLE %I.%I REPLICA IDENTITY FULL;',
      r.table_schema, r.table_name
    );
  END LOOP;
END $$;

Appliquer automatiquement aux tables futures

Pour que chaque table nouvellement créée reçoive REPLICA IDENTITY FULLautomatiquement, installez un déclencheur d’événement Postgres. Il s’exécute après chaque CREATE TABLE et définit l’identité sur la nouvelle table :

CREATE OR REPLACE FUNCTION public.set_full_replica_identity()
RETURNS event_trigger
LANGUAGE plpgsql
AS $$
DECLARE
  obj record;
BEGIN
  FOR obj IN
    SELECT * FROM pg_event_trigger_ddl_commands()
    WHERE command_tag = 'CREATE TABLE'
  LOOP
    EXECUTE format(
      'ALTER TABLE %s REPLICA IDENTITY FULL;',
      obj.object_identity
    );
  END LOOP;
END $$;

CREATE EVENT TRIGGER set_full_replica_identity_on_create
ON ddl_command_end
WHEN TAG IN ('CREATE TABLE')
EXECUTE FUNCTION public.set_full_replica_identity();

Combinez le déclencheur d’événement avec la boucle de l’onglet précédent pour couvrir les tables existantes et futures d’une configuration.

Vérifier quelles tables ont une identité de réplica définie

Pour voir quelles tables d’un schéma ont une identité de réplica configurée, exécutez :

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 le CDF.

Étape 2 : Démarrer le flux de données modifiées

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

  1. Dans votre espace de travail Azure Databricks, ouvrez Lakebase Postgres à partir du sélecteur d’application (en haut à droite).
  2. Sélectionnez votre projet Lakebase et la branche que vous souhaitez utiliser (par exemple, production ou main).
  3. Ouvrez la vue d’ensemble de la branche en cliquant sur le nom de la branche dans la barre de navigation supérieure, puis cliquez sur l’onglet CDF Lakebase .
  4. Cliquez sur Démarrer.
  5. Dans la boîte de dialogue de configuration :
    • Base de données : Valeur par défaut : databricks_postgres.
    • Schéma: Sélectionnez le schéma Postgres source.
    • Vers le catalogue : Sélectionnez le catalogue de destination dans Unity Catalog.
    • Schéma: Sélectionnez le schéma du catalogue Unity de destination.
  6. Cliquez sur Démarrer pour commencer le flux.

Vue d’ensemble de la branche avec l’onglet Lakebase CDF montrant Start et la configuration du 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 de destination et au schéma, puis ouvrez l’onglet Tables .

L’onglet CDF Lakebase comporte deux sous-onglets :

Les sous-onglets affichent le mappage et la progression par table.

  • Schémas: Répertorie chaque schéma source, son catalogue de destination et son schéma dans le catalogue Unity et un état.
  • Tables: Répertorie chaque table source, la table de destination correspondante lb_<table_name>_history, son statut (Streaming ou Snapshotting), le LSN validé (jusqu’où le flux a écrit dans Delta, affiché comme - tant que l’instantané initial est en cours) et la Dernière mise à jour (dernière fois que la table a reçu des modifications).

Vous pouvez également inspecter l’état du flux à partir de Postgres en l’exécutant dans l’éditeur SQL Lakebase :

SELECT * FROM wal2delta.tables;

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

Important

Qu’est-ce que wal2delta ? Lakebase CDF est alimenté par l’extension wal2delta Postgres, qui s’exécute à l’intérieur du calcul Lakebase. Il utilise le décodage logique pour capturer les modifications du journal de transactions en écriture anticipée (WAL) et les consigner dans des 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 de destination et votre schéma. En plus de vos colonnes source, chaque ligne comporte ces colonnes système :

Colonne Catégorie Description
_pg_change_type TEXTE Type d’opération : insert, , deleteupdate_preimageou update_postimage.
_pg_lsn BIGINT Numéro de séquence du journal Postgres.
_pg_xid INTEGER ID de transaction Postgres.
_timestamp TIMESTAMP Horodatage lorsque la modification a été traitée (sans fuseau horaire).
_sort_by BIGINT Clé de tri monotonique utilisée pour classer toutes les modifications.

Modèles de modification courants

  • Capture instantanée initiale : La première fois que CDF s’exécute sur une table Lakebase existante, chaque ligne existante est enregistrée avec _pg_change_type = 'insert'.
  • Updates: 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).
  • Supprime: Une suppression produit une ligne avec _pg_change_type = 'delete'.

Il s’agit des mêmes événements de modification que le flux de données de modification delta, de sorte que les mêmes modèles en aval s’appliquent.

Comportement opérationnel

  • Conflits de noms : Si deux tables sources correspondent au même nom de destination (par exemple, sales.users et marketing.users correspondent toutes deux à lb_users_history), la CDF écrit la première dans lb_users_history et ajoute automatiquement un suffixe à la seconde pour obtenir lb_users_history_1. Vous pouvez renommer l’une ou l’autre table de destination dans le catalogue Unity et le flux continue de fonctionner.
  • Étendue au niveau du schéma : Lorsque vous démarrez 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 comporter au moins une ligne pour apparaître dans la destination.
  • Tables sources supprimées : Si vous supprimez une table dans Lakebase, la table Delta de destination dans le catalogue Unity est conservée.

Générer 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, par ordre de simplicité, du plus simple au plus flexible.

Exemple de scénario. Une application de commerce électronique enregistre les commandes dans une table Postgres orders , chaque ligne portant un item_id et quantity. L’équipe logistique a besoin des niveaux de stock en temps réel. Avec CDF, chaque modification apportée à orders est stockée dans la table Delta lb_orders_history de Unity Catalog. Les pipelines en aval lisent ce flux de modification et mettent à jour une inventory_levels table 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 s’actualise de façon incrémentielle à mesure que de nouveaux événements de modification arrivent, et les consommateurs en aval l’interrogent comme n’importe quelle autre table.

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, à l’exception de la modification nette, de sorte que la somme en cours d’exécution reste correcte lorsque les commandes sont modifiées.

Diffuser les modifications en continu avec Spark Declarative Pipelines

Pour une architecture en médaillon structurée, utilisez pipelines Lakeflow pour déclarer des tables bronze, argent et or. Les pipelines Lakeflow les exécutent sous la forme d’un pipeline connecté, avec les points de contrôle et la gestion des dépendances pris en charge pour vous.

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 génère un delta pour chaque événement. inventory_levels agrège par item_id pour calculer le stock actuel. L’attente supprime les lignes qui feraient passer le stock en valeur négative, signalant un bogue en amont.

Pour obtenir une procédure pas à pas complète de bout en bout, consultez Tutoriel : 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, pour des fusions personnalisées, des effets secondaires ou plusieurs récepteurs — lisez directement la table d’historique avec Spark Structured Streaming et utilisez foreachBatch pour écrire dans votre destination.

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={"on_hand": F.expr("t.on_hand + s.delta")})
        .whenNotMatchedInsert(values={"item_id": "s.item_id", "on_hand": "s.delta"})
        .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 fusionne les deltas nets en inventory_levels.

Conçue pour évoluer par étapes. Chaque lb_<table_name>_history table est une table Delta d’ajout uniquement. Chaque modification source est enregistrée sous la forme d’une nouvelle ligne avec _pg_change_type le marquage de l’opération. Les vues matérialisées de Databricks SQL, les flux des pipelines Lakeflow et les travaux Spark Structured Streaming traitent tous les nouvelles lignes de manière incrémentielle depuis le journal des transactions Delta, de sorte que les pipelines en aval n’exécutent qu’une charge de travail proportionnelle aux changements survenus. Vous n’avez pas besoin d’activer le flux de données modifiées Delta sur la table d’historique, car la sémantique des modifications est déjà encodée dans les données de ligne.

Mappage de 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 CHAÎNE.

Type PostgreSQL type delta Azure Databricks Remarques
BOOLEAN BOOLEAN
INT, SMALLINT, BIGINT INT, SMALLINT, BIGINT
TEXT, VARCHAR, CHAR STRING
JSONB STRING Stocké sous forme de chaîne JSON.
ENUM STRING Stocké sous forme de libellé d’énumération.
NUMÉRIQUE / DÉCIMAL DÉCIMAL OU CHAÎNE Utilise la précision/l’échelle de la source lorsque cela est possible. Effectue une mise à l’échelle sans perte pour des valeurs de précision/d’échelle incompatibles. Utilise STRING par défaut lorsque la précision dépasse 38 ou lorsque la précision ou l’échelle ne sont pas définies (NUMERIC non borné). Toutes les colonnes NUMERIC/DECIMAL sont nullables, car les valeurs NaN sont mappées à NULL. Consultez les types numériques PostgreSQL.
DATE DATE
TIMESTAMP TIMESTAMP_NTZ
TIMESTAMPTZ TIMESTAMP
float, double FLOAT, DOUBLE

Types stockés en tant que CHAÎNE :

  • Geography/Geometry (PostGIS) : Types de l’extension PostGIS (par exemple, geometry, geography).
  • Vector (pgvector) : Type vector de l’extension pgvector.
  • Types composites/structs : Types personnalisés définis avec CREATE TYPE ... AS (field_name type, ...). Il s’agit de types de type ligne avec des champs nommés.
  • Map : Types clé-valeur de type map tels que hstore (à partir de l’extension hstore ). Postgres n’a pas de type de carte intégré. hstore est le moyen courant de stocker des paires clé-valeur dans une colonne.

Gestion des modifications de schéma

  • Le changement de nom d’une table dans Postgres (par exemple) ALTER TABLE users RENAME TO customerspermet 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) déclenchent une nouvelle capture instantanée de la table affectée. CDF lit l’intégralité de la table à partir de Postgres et la réécrit dans la table Delta de destination.

Désactiver lakebase CDF

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

  1. Dans votre espace de travail Azure Databricks, ouvrez Lakebase Postgres à partir du sélecteur d’application (en haut à droite).
  2. Sélectionnez votre projet Lakebase et la branche où vous avez configuré CDF.
  3. Ouvrez la vue d’ensemble de la branche en cliquant sur le nom de la branche dans la barre de navigation supérieure, puis cliquez sur l’onglet CDF Lakebase .
  4. Cliquez sur Désactiver. Dans la boîte de dialogue de confirmation, passez en revue l’avertissement indiquant que les modifications cesseront de circuler vers les tables Delta, puis cliquez à nouveau sur Désactiver pour confirmer.

La désactivation de CDF ne redémarre pas votre calcul.

Limitations et résolution des problèmes

Vous pouvez voir l’état de chaque table (création d’instantané, ignorée ou diffusion en continu) dans l’onglet Lakebase CDF, ou en exécutant la commande suivante dans Lakebase :

SELECT * FROM wal2delta.tables;

Les raisons courantes pour lesquelles une table n’apparaît pas dans le flux :

  • REPLICA IDENTITY FULL n’est pas défini : Exécutez ALTER TABLE <table_name> REPLICA IDENTITY FULL; pour la table. Consultez Étape 1 : Définir replica identity full.
  • Tables partitionnées : Les tables partitionnées Lakebase ne sont pas prises en charge. Un schéma qui contient des tables partitionnée provoque l’échec de ces tables.
  • Tables vides : Une table avec zéro ligne est ignorée jusqu’à ce qu’au moins une ligne existe.
  • Point de terminaison privé sur le stockage de destination : Lakebase CDF n’est pas pris en charge lorsque le stockage géré pour votre catalogue de catalogue Unity de destination est accessible uniquement via un point de terminaison privé (par exemple, lorsque l’accès au réseau public au compte de stockage est désactivé). Pour contourner ce problème, configurez un catalogue dont le stockage managé est accessible publiquement et utilisez ce catalogue comme destination CDF.

Étapes suivantes