Aller au contenu principal

Utiliser une table de contrôle pour piloter un For each Job

Lorsque vous ingérez à partir de nombreuses sources, le fait de coder en dur la liste dans le Job signifie la mise à jour du code et le redéploiement à chaque modification de la liste. Utilisez les métadonnées pour résoudre ce problème en stockant la liste des sources dans une table lue à l'exécution. Ajoutez une source en tant que nouvelle ligne et la prochaine exécution de Job la prendra en compte sans aucune modification du Job lui-même.

Ce tutoriel vous montre comment créer un Job en utilisant cette approche. Une tâche SQL lit la table de contrôle, et une tâche For each itère sur chaque ligne en parallèle.

Comment cela fonctionne

Le modèle utilise trois types de tâches câblées ensemble en séquence :

Tâche

Type

Ce qu'il fait

read_markets

SQL

query une table de configuration et capture le résultat sous forme de tableau de lignes

process_markets

Pour chaque

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

run_market_analysis_iteration

Notebook ou SQL (imbriqué dans For each)

S'exécute une fois par ligne, en utilisant les valeurs de ligne passées comme parameters pour exécuter votre logique métier.

Tâche

Type

Ce qu'il fait

read_markets

SQL

query une table de configuration et capture le résultat sous forme de tableau de lignes

process_markets

Pour chaque

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

run_market_analysis_iteration

Notebook ou SQL (imbriqué dans For each)

S'exécute une fois par ligne, en utilisant les valeurs de ligne passées comme parameters pour exécuter votre logique métier.

La sortie de la tâche SQL — un tableau JSON d'objets de ligne — s'écoule directement dans le champ Entrées de la tâche For each en utilisant la référence de valeur dynamique {{tasks.read_markets.output.rows}}. La tâche For each transmet ensuite chaque ligne à la tâche imbriquée en tant que paramètres, disponibles en tant que {{input.market}} et {{input.currency}}.

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 configuration (par exemple, config)
  • Un SQL Warehouse pour exécuter les tâches SQL

Étape 1 : Créer la table de configuration

La table de configuration est la source de vérité pour la liste des valeurs que votre Job traite. Lorsque vous devez ajouter ou supprimer du travail, vous mettez à jour cette table, pas le job.

Exécutez le SQL suivant pour créer une table markets dans votre schéma config :

SQL
CREATE OR REPLACE TABLE config.markets AS
SELECT * FROM VALUES
('NL', 'EUR'),
('UK', 'GBP'),
('US', 'USD')
AS t(market, currency);

Utilisez un Notebook Databricks, l'éditeur SQL ou toute tâche SQL pour exécuter cette instruction. Après cette étape, config.markets contient trois lignes, une par marché, chacune avec son code de devise.

Étape 2 : écrire le code de traitement

La tâche imbriquée à l'intérieur de la tâche For each s'exécute une fois par ligne. Choisissez une tâche de Notebook ou une tâche SQL en fonction de votre logique métier.

Créez un nouveau notebook à un chemin tel que /Workspace/Users/<username>/process_market. Ce notebook s'exécute une fois par itération de la tâche For each, recevant une valeur de marché différente à chaque fois.

Ajoutez le code suivant au Notebook :

Python
# Set default values for testing the notebook outside of a job.
# When the notebook runs inside a For each task, the job overrides these defaults.
dbutils.widgets.text("market", "NL", "Market")
dbutils.widgets.text("currency", "EUR", "Currency")

# Read the parameters passed by the For each task
market = dbutils.widgets.get("market")
currency = dbutils.widgets.get("currency")

print(f"Processing market: {market} ({currency})")

# Your business logic goes here. For example:
df = spark.table("sales.transactions").filter(
f"market = '{market}' AND currency_code = '{currency}'"
)
display(df)

Les appels dbutils.widgets.text() définissent les default valeurs par défaut afin que vous puissiez exécuter le notebook directement dans votre Workspace sans le connecter à un Job. Lorsque le Notebook s'exécute en tant que tâche imbriquée dans une tâche For each, le Job remplace les valeurs par default par les valeurs de parameter réelles pour cette itération.

remarque

Appelez dbutils.widgets.text() avant dbutils.widgets.get(). Si get est appelé avant text, le Notebook génère une erreur InputWidgetNotDefined lorsque vous l'exécutez en dehors d'un Job.

L'utilisation des valeurs par défaut vous permet de tester le Notebook en dehors d'un Job, mais soyez conscient : si la tâche For each est mal configurée et ne transmet pas de parameters, le Notebook utilise les valeurs par défaut et réussit en silence au lieu d'échouer, ce qui peut rendre la mauvaise configuration plus difficile à détecter.

Étape 3 : Créez le Job

Dans votre Workspace Databricks, dans la barre latérale, choisissez Icône Plus. Nouveau > Icône Workflows. Job . Donnez au Job un nom descriptif, par exemple Market Analysis.

Étape 4 : Configurez la tâche de recherche SQL

La tâche SQL exécute votre query de configuration et rend sa sortie disponible pour les tâches en aval.

  1. Dans l'éditeur de Job, cliquez sur Icône Plus. Ajouter une tâche .

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

  3. Set Type to SQL .

  4. Dans le champ SQL , saisissez la query suivante :

    SQL
    SELECT market, currency FROM config.markets
  5. Définissez SQL warehouse sur un warehouse dans votre Workspace.

  6. Cliquez sur Enregistrer la tâche .

Lorsque cette tâche s'exécute, Databricks exécute la query et capture le résultat sous forme de tableau JSON dans tasks.read_markets.output.rows. Le résultat de la tâche SQL est toujours renvoyé sous forme de tableau JSON — aucune configuration supplémentaire n'est requise. La forme générique de cette référence est tasks.<task-name>.output.rows, où <task-name> correspond à la clé de tâche que vous avez définie dans l'éditeur de jobs. Le résultat se présente comme suit :

JSON
[
{ "market": "NL", "currency": "EUR" },
{ "market": "UK", "currency": "GBP" },
{ "market": "US", "currency": "USD" }
]

Étape 5 : Configurer la tâche For each

La tâche For each lit la sortie SQL et lance une exécution de tâche imbriquée par ligne.

  1. Cliquez sur Icône Plus. Ajouter une tâche et définissez Dépend de sur read_markets.

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

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

  4. Dans le champ Inputs , saisissez :


    {{tasks.read_markets.output.rows}}

    Ceci fait référence au tableau de lignes capturé par la tâche SQL.

  5. Définissez 2 la **concurrence** à pour permettre à deux itérations de s'exécuter en parallèle. Augmentez cette valeur lorsque votre tâche imbriquée prend en charge un parallélisme plus élevé.

  6. Cliquez sur Ajouter une tâche à parcourir et configurez la tâche imbriquée en fonction du type que vous avez choisi à l'étape 2 :

  1. Définissez le Nom de la tâche sur run_market_analysis_iteration.

  2. Set Type to Notebook .

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

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

    • Clé : market, Valeur : {{input.market}}
    • Clé : currency, Valeur : {{input.currency}}

    Chaque référence {{input.<key>}} est résolue au champ correspondant de l'objet de ligne de l'itération actuelle.

  5. Cliquez sur Enregistrer la tâche .

Votre DAG de Jobs montre maintenant read_markets s'écoulant dans process_markets, avec la tâche imbriquée visible à l'intérieur du nœud For each.

É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 process_markets pour développer la tâche For each.
  3. La page d'exécution du job affiche un tableau d'itérations, une ligne par valeur marchande, chacune indiquant son statut, l'heure de start et la durée.
  4. Cliquez sur n'importe quelle ligne d'itération pour ouvrir le résultat d'exécution de la tâche et confirmer qu'il a reçu la bonne valeur marchande.

Si une itération spécifique échoue, réexécutez uniquement cette itération depuis la page d'exécution du Job sans réexécuter l'intégralité du Job.

Étendre le modèle

Pour ajouter un nouveau marché, insérez une ligne dans la table de configuration :

SQL
INSERT INTO config.markets VALUES ('DE', 'EUR');

La prochaine exécution du Job inclut automatiquement l'Allemagne, sans aucune modification de configuration du Job ni d'édition de notebook requise.

Ce même schéma fonctionne pour tout cas d'usage où vous souhaitez que les données favorisent l'itération :

  • Traitement par client : Une ligne par ID client. Le Notebook applique des transformations spécifiques aux clients ou livre à des destinations spécifiques aux clients.
  • Ingestion de table : Une ligne par nom de table source. Le Notebook lit et ingère chaque table.
  • Traitement de rattrapage : Une ligne par partition de date. Le notebook retraite les données historiques pour cette partition.
  • Exécution basée sur les indicateurs de fonctionnalité : une ligne par fonctionnalité ou expérimentation activée. Le Notebook active la logique correspondante.

Pour retirer un élément du traitement, supprimez sa ligne ou ajoutez une colonne d'indicateurs active et filtrez dans la query SQL :

SQL
SELECT market, currency FROM config.markets WHERE active = TRUE

Ressources supplémentaires