Aller au contenu principal

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 each pour 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

read_watermarks

SQL

Lit la table de contrôle de filigrane et renvoie une ligne par table source

copy_tables

Pour chaque

Itère sur {{tasks.read_watermarks.output.rows}}, exécutant la tâche imbriquée une fois par table source

copy_incremental (imbriqué)

Notebook

Lit les lignes ajoutées depuis le dernier filigrane, les écrit dans la table cible et fait avancer le filigrane

Tâche

Type

Ce qu'il fait

read_watermarks

SQL

Lit la table de contrôle de filigrane et renvoie une ligne par table source

copy_tables

Pour chaque

Itère sur {{tasks.read_watermarks.output.rows}}, exécutant la tâche imbriquée une fois par table source

copy_incremental (imbriqué)

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 schémas et des tables dans Unity Catalog
  • Un SQL Warehouse pour exécuter des tâches SQL

Ce tutoriel lit à partir du dataset d’exemple samples.wanderbricks et écrit dans un schéma nommé par les variables catalog et schema en haut de chaque bloc de code. Ces variables ont pour valeur par default main.example_output. Pour écrire ailleurs, modifiez les deux valeurs de manière cohérente dans chaque bloc. L’exemple crée le schéma s’il n’existe pas.

É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 code SQL suivant pour créer la table de contrôle et enregistrer deux tables sources. L’exemple enregistre samples.wanderbricks.users et samples.wanderbricks.properties, en copiant chacun dans une table cible de votre propre schéma. Les variables catalog et schema définissent à la fois l'emplacement de la table de contrôle et les noms de la table cible :

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

CREATE SCHEMA IF NOT EXISTS IDENTIFIER(catalog || '.' || schema);

CREATE OR REPLACE TABLE IDENTIFIER(catalog || '.' || schema || '.watermarks') (
source_table STRING NOT NULL,
target_table STRING NOT NULL,
watermark_column STRING NOT NULL,
last_watermark TIMESTAMP NOT NULL
);

INSERT INTO IDENTIFIER(catalog || '.' || schema || '.watermarks') VALUES
('samples.wanderbricks.users', catalog || '.' || schema || '.users', 'created_at', '1970-01-01'),
('samples.wanderbricks.properties', catalog || '.' || schema || '.properties', 'created_at', '1970-01-01');

Les deux tables sources utilisent created_at comme colonne de filigrane. Il s’agit d’un Timestamp d’insertion qui ne fait qu’augmenter à mesure que de nouvelles lignes arrivent. La définition de last_watermark sur 1970-01-01 lors de la première exécution entraîne la copie de toutes les lignes existantes par le notebook. Cela agit comme un chargement complet initial. Les exécutions ultérieures ne copient que les lignes ajoutées après l’exécution précédente.

remarque

La recréation de la table de contrôle Reset chaque last_watermark sur 1970-01-01, mais ne vide pas les tables cibles. Si vous réexécutez cette étape puis relancez le job, le Notebook traite les deux sources comme n'ayant jamais été copiées et ajoute chaque ligne historique une seconde fois. Pour start proprement, supprimez les tables cibles lorsque vous recréez la table de contrôle.

É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. Les paramètres par default du widget vous permettent d'exécuter et de tester le Notebook directement. Lorsqu'il s'exécute à l'intérieur de la tâche For each, le job les remplace par les valeurs de chaque itération, y compris catalog et schema qui localisent la table de contrôle.

Le code lit uniquement les lignes ajoutées depuis le dernier filigrane et les ajoute à la table cible, en la créant si elle n’existe pas. Il calcule ensuite la marque de niveau supérieur (high water mark) des lignes qu'il vient d'écrire et fait avancer la table de contrôle afin que la prochaine exécution start à partir de là :

Python
from pyspark.sql.functions import max as spark_max

# Widget defaults let you run the notebook directly; the For each task overrides them per iteration
dbutils.widgets.text("catalog", "main", "Catalog")
dbutils.widgets.text("schema", "example_output", "Schema")
dbutils.widgets.text("source_table", "samples.wanderbricks.users", "Source table")
dbutils.widgets.text("target_table", "main.example_output.users", "Target table")
dbutils.widgets.text("watermark_column", "created_at", "Watermark column")
dbutils.widgets.text("last_watermark", "1970-01-01", "Last watermark")

catalog = dbutils.widgets.get("catalog")
schema = dbutils.widgets.get("schema")
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 rows newer than the last watermark. A strict > can skip rows that share the
# stored high-water timestamp; for insert-only sources with distinct timestamps this is safe.
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 the new rows, creating the target table on the first run
new_rows.write.format("delta").mode("append").saveAsTable(target_table)

# Compute the high-water mark from the rows just written
new_watermark = new_rows.agg(spark_max(watermark_column)).collect()[0][0]

# Advance the control table so the next run starts from here
spark.sql(f"""
UPDATE {catalog}.{schema}.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")
remarque

Ce Notebook utilise le mode append, qui est approprié lorsque la source ne contient que des insertions, comme c'est le cas pour samples.wanderbricks.users et samples.wanderbricks.properties. Si votre source contient des mises à jour, utilisez un filigrane sur le Timestamp de mise à jour et utilisez une instruction MERGE au lieu de write.mode("append") pour effectuer des upserts de lignes dans la table cible. Consultez Upsert into a Delta Lake table using 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. Comme la query déclare les variables catalog et schema qui localisent la table de contrôle, elle doit être exécutée en tant que fichier SQL à instructions multiples plutôt que dans le champ SQL en ligne de la tâche.

  1. Créez un fichier SQL dans votre Workspace, tel que /Workspace/Users/<username>/read_watermarks.sql, avec le contenu suivant. Définissez catalog et schema sur les mêmes valeurs que celles utilisées à l'étape 1 :

    SQL
    DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
    DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

    SELECT source_table, target_table, watermark_column, last_watermark
    FROM IDENTIFIER(catalog || '.' || schema || '.watermarks');
  2. Dans le job, cliquez sur Ajouter une tâche .

  3. Définissez le Nom de la tâche sur read_watermarks.

  4. Définissez Type sur SQL , puis définissez Tâche SQL sur Fichier .

  5. Définissez Chemin sur le fichier SQL que vous avez créé.

  6. Définissez SQL warehouse sur un warehouse dans votre Workspace.

  7. 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. Après un chargement complet initial, chaque last_watermark reflète la ligne la plus récente copiée à partir de cette source :

JSON
[
{
"source_table": "samples.wanderbricks.users",
"target_table": "main.example_output.users",
"watermark_column": "created_at",
"last_watermark": "2025-07-30T23:05:18.000Z"
},
{
"source_table": "samples.wanderbricks.properties",
"target_table": "main.example_output.properties",
"watermark_column": "created_at",
"last_watermark": "2025-07-30T00: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.

  1. Cliquez sur **Ajouter une tâche** et définissez **Dépend de** read_watermarksà.

  2. Définissez le Nom de la tâche sur copy_tables.

  3. Définissez **Type** sur **Pour chacun**.

  4. Dans le champ Inputs , saisissez :


    {{tasks.read_watermarks.output.rows}}
  5. Définissez 2 la **simultanéité** sur pour copier deux tables à la fois. Augmentez cette valeur si votre warehouse peut prendre en charge un parallélisme plus élevé.

  6. Cliquez sur **Ajouter une tâche à itérer** pour configurer la tâche imbriquée.

  7. Définissez le Nom de la tâche sur copy_incremental.

  8. Set Type to Notebook .

  9. Définissez le chemin d’accès au notebook que vous avez créé à l'étape 2.

  10. Cliquez sur Paramètres , puis cliquez sur Ajouter pour ajouter chacun des paramètres suivants :

Clé

Valeur

catalog

main

schema

example_output

source_table

{{input.source_table}}

target_table

{{input.target_table}}

watermark_column

{{input.watermark_column}}

last_watermark

{{input.last_watermark}}

Clé

Valeur

catalog

main

schema

example_output

source_table

{{input.source_table}}

target_table

{{input.target_table}}

watermark_column

{{input.watermark_column}}

last_watermark

{{input.last_watermark}}

Définissez catalog et schema sur les mêmes valeurs que celles utilisées à l’étape 1 afin que le notebook fasse avancer la table de contrôle lue par la tâche SQL. Chaque référence {{input.<key>}} correspond au 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.

  1. Cliquez sur Run now pour Trigger le Job.
  2. Sur la page d'exécution du job, cliquez sur le nœud copy_tables pour développer la tâche For each.
  3. 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.
  4. 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 que le filigrane a progressé, exécutez la query suivante une fois le Job terminé. Définissez catalog et schema sur les mêmes valeurs que celles utilisées lors des étapes précédentes :

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

SELECT source_table, last_watermark
FROM IDENTIFIER(catalog || '.' || schema || '.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

Chaque extrait déclare les mêmes variables catalog et schema utilisées dans les étapes précédentes. Définissez-les sur les valeurs qui localisent votre table de contrôle.

Ajouter une nouvelle table source : insérez 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. Cet extrait suppose que vous avez déjà ajouté la colonne active à partir de Suspendre une table ci-dessous, il définit donc active sur true. Les deux extensions peuvent être appliquées dans n’importe quel ordre ; si vous n’avez pas encore ajouté la colonne, supprimez la dernière valeur et sa colonne :

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

INSERT INTO IDENTIFIER(catalog || '.' || schema || '.watermarks')
(source_table, target_table, watermark_column, last_watermark, active)
VALUES
('samples.wanderbricks.hosts', catalog || '.' || schema || '.hosts', 'joined_at', '1970-01-01', TRUE);

Pause a table : ajoutez une colonne active, remplissez-la avec true pour les lignes existantes, puis filtrez-la dans la tâche de fichier SQL. Delta nécessite l’ajout de la colonne et la définition de sa valeur dans des instructions distinctes :

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

ALTER TABLE IDENTIFIER(catalog || '.' || schema || '.watermarks') ADD COLUMN active BOOLEAN;

UPDATE IDENTIFIER(catalog || '.' || schema || '.watermarks') SET active = TRUE;

Ajoutez ensuite WHERE active = TRUE au SELECT dans votre fichier read_watermarks.sql afin que le job ignore les tables en pause.

Backfill a table : Reset son filigrane pour copier à nouveau à partir d'un point spécifique :

SQL
DECLARE OR REPLACE VARIABLE catalog STRING DEFAULT 'main';
DECLARE OR REPLACE VARIABLE schema STRING DEFAULT 'example_output';

UPDATE IDENTIFIER(catalog || '.' || schema || '.watermarks')
SET last_watermark = '2025-01-01'
WHERE source_table = 'samples.wanderbricks.users';

Ressources supplémentaires