Copiez plusieurs tables de manière incrémentielle avec une tâche For each.
Lorsque vous devez copier des données de nombreuses tables sources vers des tables Unity Catalog selon un calendrier, la copie de toutes les lignes à chaque exécution est lente et coûteuse. Utilisez un filigrane pour suivre la dernière ligne traitée pour chaque table et copier uniquement les nouvelles lignes à chaque exécution.
Ce tutoriel vous montre comment créer un Job piloté par les métadonnées qui :
- Stocke la liste des tables sources et leur état de filigrane dans une table de contrôle Delta.
- Utilise une tâche
For eachpour traiter chaque table en parallèle - Copie uniquement les lignes ajoutées depuis la dernière exécution réussie
- Met à jour le filigrane après chaque copie réussie.
Comment cela fonctionne
Le Job utilise trois types de tâches enchaînés en séquence :
Tâche | Type | Ce qu'il fait |
|---|---|---|
| SQL | Lit la table de contrôle de filigrane et renvoie une ligne par table source |
| Pour chaque | Itère sur |
| Notebook | Lit les lignes ajoutées depuis le dernier filigrane, les écrit dans la table cible et fait avancer le filigrane |
La sortie de la tâche SQL — un tableau JSON d'objets de ligne — s'intègre dans le champ Inputs de la tâche For each à l'aide de {{tasks.read_watermarks.output.rows}}. Le Notebook imbriqué reçoit source_table, target_table, watermark_column et last_watermark pour chaque itération.
Prérequis
- Un Databricks Workspace avec la permission de créer des Job et des Notebook
- Autorisation de créer des tables dans Unity Catalog
- Un schéma Unity Catalog où vous pouvez créer la table de contrôle et les tables cibles (par exemple,
config) - Un SQL Warehouse pour exécuter des tâches SQL
- Tables sources qui contiennent une colonne à croissance monotone telle qu'un Timestamp ou une séquence d'entiers
Étape 1 : Créer la table de contrôle du filigrane
La table de contrôle du filigrane est la source de vérité pour déterminer les tables à traiter et dans quelle mesure chaque table a été copiée. Chaque ligne représente une table source.
Exécutez le SQL suivant pour créer la table de contrôle et enregistrer deux tables sources :
CREATE OR REPLACE TABLE config.watermarks (
source_table STRING NOT NULL,
target_table STRING NOT NULL,
watermark_column STRING NOT NULL,
last_watermark TIMESTAMP NOT NULL
);
INSERT INTO config.watermarks VALUES
('sales.raw_orders', 'sales.orders', 'updated_at', '1970-01-01'),
('sales.raw_customers', 'sales.customers', 'updated_at', '1970-01-01');
La définition de last_watermark à 1970-01-01 lors de la première exécution entraîne la copie par le notebook de toutes les lignes existantes, agissant comme un chargement initial complet. Les exécutions suivantes copient uniquement les lignes ajoutées ou mises à jour après l'exécution précédente.
Étape 2 : Écrire le notebook de copie
Le Notebook s'exécute une fois par itération de table. Il lit le filigrane, filtre la source, écrit vers la cible et avance le filigrane.
Créez un notebook à un chemin tel que /Workspace/Users/<username>/copy_incremental et ajoutez le code suivant :
# Set defaults for running the notebook outside a job
dbutils.widgets.text("source_table", "sales.raw_orders", "Source table")
dbutils.widgets.text("target_table", "sales.orders", "Target table")
dbutils.widgets.text("watermark_column", "updated_at", "Watermark column")
dbutils.widgets.text("last_watermark", "1970-01-01", "Last watermark")
source_table = dbutils.widgets.get("source_table")
target_table = dbutils.widgets.get("target_table")
watermark_column = dbutils.widgets.get("watermark_column")
last_watermark = dbutils.widgets.get("last_watermark")
# Read only new rows from the source table
new_rows = spark.table(source_table).filter(
f"{watermark_column} > '{last_watermark}'"
)
row_count = new_rows.count()
print(f"Copying {row_count} new rows from {source_table}")
if row_count > 0:
# Append new rows to the target table, creating it if it does not exist
new_rows.write.format("delta").mode("append").saveAsTable(target_table)
# Compute the new high-water mark from the rows just written
from pyspark.sql.functions import max as spark_max
new_watermark = new_rows.agg(spark_max(watermark_column)).collect()[0][0]
# Advance the watermark so the next run starts from here
spark.sql(f"""
UPDATE config.watermarks
SET last_watermark = CAST('{new_watermark}' AS TIMESTAMP)
WHERE source_table = '{source_table}'
""")
print(f"Watermark for {source_table} advanced to {new_watermark}")
else:
print(f"No new rows for {source_table}, watermark unchanged")
Les dbutils.widgets.text() par défaut vous permettent d’exécuter et de tester directement le notebook. Lorsque le notebook s'exécute dans la tâche For each, le Job remplace ces valeurs par default par les valeurs réelles pour chaque itération.
Ce Notebook utilise le mode append, ce qui est approprié lorsque la source ne contient que des insertions. Si votre source contient des mises à jour, utilisez une instruction MERGE au lieu de write.mode("append") pour upsert les lignes dans la table cible. Voir Upsert dans une table Delta Lake à l'aide de l'opération Merge pour la syntaxe de Merge.
Étape 3 : Créez le Job
Dans votre workspace Databricks, cliquez sur **Workflows** dans la barre latérale, puis cliquez sur **Créer un job**. Donnez au Job un nom tel que Incremental table copy.
Étape 4 : Configurer la tâche de recherche de filigrane
La tâche SQL lit la table de contrôle et rend le résultat disponible pour la tâche For each.
-
Cliquez sur Ajouter une tâche .
-
Définissez le Nom de la tâche sur
read_watermarks. -
Set Type to SQL .
-
Dans le champ **SQL**, entrez :
SQLSELECT source_table, target_table, watermark_column, last_watermark
FROM config.watermarks -
Définissez SQL warehouse sur un warehouse dans votre Workspace.
-
Cliquez sur **Créer une tâche**.
Lorsque cette tâche s'exécute, Databricks capture le résultat sous forme de tableau JSON dans tasks.read_watermarks.output.rows:
[
{
"source_table": "sales.raw_orders",
"target_table": "sales.orders",
"watermark_column": "updated_at",
"last_watermark": "2024-06-01T12:00:00.000Z"
},
{
"source_table": "sales.raw_customers",
"target_table": "sales.customers",
"watermark_column": "updated_at",
"last_watermark": "2024-06-01T12:00:00.000Z"
}
]
Étape 5 : Configurer la tâche For each
La tâche For each lit le résultat SQL et lance une exécution de tâche imbriquée par table source.
-
Cliquez sur **Ajouter une tâche** et définissez **Dépend de**
read_watermarksà. -
Définissez le Nom de la tâche sur
copy_tables. -
Définissez **Type** sur **Pour chacun**.
-
Dans le champ Inputs , saisissez :
{{tasks.read_watermarks.output.rows}} -
Définissez
2la **simultanéité** sur pour copier deux tables à la fois. Augmentez cette valeur si votre warehouse peut prendre en charge un parallélisme plus élevé. -
Cliquez sur **Ajouter une tâche à itérer** pour configurer la tâche imbriquée.
-
Définissez le Nom de la tâche sur
copy_incremental. -
Set Type to Notebook .
-
Définissez le chemin d’accès au notebook que vous avez créé à l'étape 2.
-
Cliquez sur Paramètres , puis cliquez sur Ajouter pour ajouter chacun des paramètres suivants :
Clé | Valeur |
|---|---|
|
|
|
|
|
|
|
|
Chaque référence {{input.<key>}} est résolue dans le champ correspondant de la ligne de l'itération actuelle.
11. Cliquez sur **Créer une tâche**.
Étape 6 : Exécutez le job et vérifiez.
- Cliquez sur Run now pour Trigger le Job.
- Sur la page d'exécution du job, cliquez sur le nœud
copy_tablespour développer la tâcheFor each. - La page d'exécution affiche un tableau d'itérations – une ligne par table source –, chacune affichant son statut, son heure de start et sa durée.
- Cliquez sur n'importe quelle itération pour afficher la sortie du Notebook et confirmer le nombre de lignes et la mise à jour du filigrane.
Pour confirmer l'avancement du watermark, exécutez la query suivante une fois le Job terminé :
SELECT source_table, last_watermark FROM config.watermarks;
Chaque valeur last_watermark doit maintenant refléter le timestamp de la ligne la plus récemment copiée. Si une valeur est toujours 1970-01-01, la table source ne contient aucune ligne correspondant au filtre, ou la tâche de copie a rencontré une erreur — vérifiez la sortie d'exécution de la tâche pour plus de détails.
Étendre le modèle
Ajouter une nouvelle table source : Insérer une ligne dans la table de contrôle. La prochaine exécution du job le récupère automatiquement, en commençant par un chargement complet depuis 1970-01-01:
INSERT INTO config.watermarks VALUES
('sales.raw_products', 'sales.products', 'updated_at', '1970-01-01');
Mettre une table en pause : Ajoutez une colonne active et filtrez dans la tâche SQL :
ALTER TABLE config.watermarks ADD COLUMN active BOOLEAN DEFAULT TRUE;
-- In the SQL task:
SELECT source_table, target_table, watermark_column, last_watermark
FROM config.watermarks
WHERE active = TRUE
Backfill a table : Reset son filigrane pour copier à nouveau à partir d'un point spécifique :
UPDATE config.watermarks
SET last_watermark = '2024-01-01'
WHERE source_table = 'sales.raw_orders';
Ressources supplémentaires
- Utilisez une tâche
For eachpour exécuter une autre tâche en boucle — Référence complète pour la configuration des tâchesFor each, y compris les options de concurrence - Utiliser une table de contrôle pour piloter un
For eachJob — Piloter unFor eachJob à partir d'une table de configuration en direct - Mettre à jour/insérer (upsert) dans une Delta Lake table using Merge — Utilisez
MERGEpour mettre à jour/insérer des lignes lorsque les données sources incluent des mises à jour - Les AUTO CDC APIs : simplifiez la capture de données modifiées avec les pipelines — Utilisez la capture de données modifiées pour les sources qui suivent les insertions, les mises à jour et les suppressions