Aller au contenu principal

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

Lorsque vous exécutez le même traitement sur de nombreuses entrées, telles que des marchés, des tables sources, des clients ou des partitions de date, le codage en dur de cette liste dans votre job signifie que vous devez modifier le code et redéployer à chaque fois que la liste change. Au lieu de cela, stockez la liste dans une table de contrôle que le job lit au moment de l'exécution. Pour ajouter ou supprimer du travail, mettez à jour une ligne dans la table ; la prochaine exécution du job prendra en compte la modification sans que vous ayez à modifier le job lui-même. Il s’agit d’un modèle piloté par les métadonnées : ce sont les données, et non le code, qui contrôlent ce que le job traite.

Ce tutoriel crée un job qui utilise ce modèle sur le dataset d’exemple Wanderbricks préinstallé, afin que vous puissiez l’exécuter de bout en bout sans créer de données source. Le scénario est une plateforme de location de vacances qui exécute la même analyse de prix pour chaque segment de propriété (tel que Ski Resort ou Urban Year-Round). Une table de contrôle répertorie les segments à analyser, une tâche SQL lit cette table et une tâche For each exécute l’analyse une fois par segment, en parallèle.

Comment cela fonctionne

Le Job relie trois tâches ensemble en séquence :

Tâche

Type

Ce qu'il fait

read_segments

SQL

Lit la table de contrôle et capture les lignes sous forme de tableau JSON

process_segments

Pour chaque

Itère sur le tableau de lignes, en lançant la tâche imbriquée une fois par ligne

run_segment_analysis

Notebook ou SQL (imbriqué dans For each)

S'exécute une fois par ligne, en utilisant les valeurs de cette ligne pour analyser un segment de propriété

Tâche

Type

Ce qu'il fait

read_segments

SQL

Lit la table de contrôle et capture les lignes sous forme de tableau JSON

process_segments

Pour chaque

Itère sur le tableau de lignes, en lançant la tâche imbriquée une fois par ligne

run_segment_analysis

Notebook ou SQL (imbriqué dans For each)

S'exécute une fois par ligne, en utilisant les valeurs de cette ligne pour analyser un segment de propriété

Le flux est read_segmentsprocess_segmentsrun_segment_analysis (une fois par ligne). La sortie de la tâche SQL, un tableau JSON d'objets de ligne, circule dans le champ Entrées de la tâche For each via la référence de valeur dynamique {{tasks.read_segments.output.rows}}. La tâche For each transmet ensuite les champs de chaque ligne à la tâche imbriquée sous forme de paramètres, disponibles en tant que {{input.property_type}} et {{input.min_price}}.

Prérequis

  • Un workspace Databricks avec l'autorisation de créer des jobs et des notebooks.
  • Autorisation de créer des tables dans Unity Catalog, et autorisation de créer un schéma dans un catalogue (les privilèges USE CATALOG et CREATE SCHEMA) pour contenir la table de contrôle.
  • Un SQL Warehouse pour exécuter les tâches SQL. Si vous n'en avez pas, consultez Créer un SQL Warehouse.
  • Le catalogue samples, qui est disponible dans chaque Workspace activé pour Unity Catalog. Le tutoriel effectue une lecture à partir de samples.wanderbricks.properties, il n'y a donc aucune donnée source à configurer.

Étape 1 : créer la table de contrôle

La table de contrôle est la source de vérité pour la liste des segments que votre job traite. Pour modifier ce que fait le job, vous devez mettre à jour cette table, et non le job.

Exécutez le SQL suivant dans un Notebook Databricks ou dans l’éditeur SQL. La première instruction crée un schéma pour contenir la table de contrôle, et la seconde crée la table avec une ligne par segment de propriété et le prix de liste minimum à inclure dans l'analyse de ce segment :

SQL
USE CATALOG <catalog-name>;

CREATE SCHEMA IF NOT EXISTS config;

CREATE OR REPLACE TABLE config.property_segments AS
SELECT * FROM VALUES
('Urban Year-Round', 150),
('Summer Getaway', 200),
('Ski Resort', 250)
AS t(property_type, min_price);

Remplacez <catalog-name> par un catalogue dans lequel vous pouvez créer des schémas, comme votre catalogue de workspace. Utilisez le même catalogue partout où le tutoriel fait référence à config.property_segments, y compris la query de recherche à l'étape 3.

Après cette étape, config.property_segments contient trois lignes, une par segment. Chaque ligne transporte les deux valeurs que le Job transmet à chaque itération : le property_type à analyser et le seuil min_price sur lequel filtrer.

Étape 2 : Rédiger la logique d’analyse

La tâche imbriquée à l'intérieur de la tâche For each s'exécute une fois par ligne de la table de contrôle, en recevant les property_type et min_price de cette ligne en tant que paramètres. Vous pouvez écrire cette logique sous forme de tâche de notebook ou de tâche SQL. Faites votre choix en fonction de votre logique métier :

  • Utilisez une tâche Notebook lorsque la logique par itération nécessite du code procédural, plusieurs langages ou des bibliothèques (par exemple, une étape de Data Science ou de machine learning).
  • Utilisez une tâche SQL lorsque la logique est une query unique ou une Transformation que vous pouvez exprimer de manière déclarative. Une tâche SQL nécessite un SQL Warehouse.

Les deux variantes ci-dessous produisent le même résultat : pour le segment en cours de traitement, le nombre d’annonces au prix plancher ou au-dessus et leur prix moyen.

Créez un nouveau notebook à un chemin tel que /Workspace/Users/<username>/run_segment_analysis. Ce notebook s’exécute une fois par itération de la tâche For each, en recevant un segment différent à chaque fois.

Ajoutez le code suivant au Notebook :

Python
# Set default values so you can run the notebook on its own while developing.
# When the notebook runs inside a For each task, the job overrides these defaults.
dbutils.widgets.text("property_type", "Ski Resort", "Property type")
dbutils.widgets.text("min_price", "250", "Minimum price")

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

result = spark.sql(
"""
SELECT :property_type AS property_type,
COUNT(*) AS property_count,
ROUND(AVG(base_price), 2) AS avg_price
FROM samples.wanderbricks.properties
WHERE property_type = :property_type
AND base_price >= :min_price
""",
args={&quot;property_type&quot;: property_type, &quot;min_price&quot;: min_price},
)
display(result)
remarque

Appelez dbutils.widgets.text() avant dbutils.widgets.get(). Si vous appelez get en premier, l'exécution du notebook en dehors d'un job génère une erreur InputWidgetNotDefined.

Étape 3 : Créer la query de recherche

La tâche de recherche lit la table de contrôle via une query enregistrée. Comme à l'étape 2, créez et enregistrez la query dans l'éditeur SQL maintenant, puis attachez-la à la tâche de recherche à l'étape 4.

  1. Dans votre Workspace Databricks, cliquez sur Icône Plus. Nouveau > Icône de query. query pour ouvrir l’éditeur SQL.

  2. Saisissez les informations suivantes, en utilisant le même catalogue que celui choisi à l’étape 1 :

    SQL
    SELECT property_type, min_price FROM <catalog-name>.config.property_segments;

    Le nom est entièrement qualifié car le SQL Warehouse qui exécute cette query peut utiliser default un catalogue différent de celui dans lequel vous avez créé la table.

  3. Cliquez sur le titre New Query <date> dans l’en-tête de tab de votre fichier SQL, et donnez-lui le nom read_segments. Cliquez ensuite sur Enregistrer pour le déplacer vers un dossier où vous souhaitez le stocker.

Étape 4 : créer et configurer le job

Une fois les deux queries enregistrées, créez le job et ajoutez ses deux tâches : la tâche de recherche SQL qui lit la table de contrôle, et la tâche For each qui exécute l'analyse pour chaque ligne.

Créer le Job

Dans votre workspace Databricks, dans la barre latérale, cliquez sur Icône Plus. New > Icône Workflows. Job . Donnez au job un nom descriptif, tel que Segment Analysis.

Configurez la tâche de recherche SQL

Cette tâche lit la table de contrôle et rend ses lignes disponibles pour la tâche For each en exécutant la query read_segments que vous avez enregistrée à l’étape 3.

  1. Cliquez sur la vignette query SQL pour configurer la première tâche. Si la vignette SQL query n’est pas disponible, cliquez sur Ajouter un autre type de tâche et recherchez SQL query .
  2. Définissez le Nom de la tâche sur read_segments.
  3. Si nécessaire, sélectionnez SQL query dans le menu déroulant Type .
  4. Dans le champ SQL query , sélectionnez la query read_segments que vous avez enregistrée à l'étape 3.
  5. Définissez SQL warehouse sur un warehouse dans votre Workspace.
  6. 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_segments.output.rows. La sortie de la tâche SQL est toujours renvoyée sous forme de tableau JSON ; vous n'avez donc besoin d'aucune configuration supplémentaire. La forme générale de la référence est tasks.<task-name>.output.rows, où <task-name> correspond au nom de tâche que vous avez défini. La sortie se présente comme suit :

JSON
[
{ "property_type": "Urban Year-Round", "min_price": 150 },
{ "property_type": "Summer Getaway", "min_price": 200 },
{ "property_type": "Ski Resort", "min_price": 250 }
]

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 sélectionnez Pour chaque .

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

  3. Vérifiez que Depends on est défini sur read_segments.

  4. Dans le champ Inputs , saisissez le tableau de lignes capturé par la tâche SQL :


    {{tasks.read_segments.output.rows}}
  5. Définissez Concurrency sur 2 pour exécuter deux itérations en parallèle. Augmentez cette valeur lorsque votre tâche imbriquée prend en charge un parallélisme plus élevé.

  6. Pour terminer cette tâche, cliquez sur Ajouter une tâche à parcourir et configurez la tâche imbriquée qui s’exécute à chaque itération.

La tâche For each et sa tâche imbriquée sont créées ensemble en tant que tâche unique. 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_segment_analysis.

  2. Set Type to Notebook .

  3. Définissez Path sur le notebook que vous avez créé à l’étape 2.

  4. Cliquez sur Parameters , puis sur Add pour ajouter chaque paramètre :

    • Clé : property_type, Valeur : {{input.property_type}}
    • Clé : min_price, Valeur : {{input.min_price}}

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

  5. Cliquez sur Create task pour créer la tâche For each et sa tâche imbriquée ensemble.

Votre graphe orienté acyclique (DAG) de job montre désormais read_segments s’écoulant vers process_segments, avec la tâche imbriquée à l’intérieur du nœud For each.

Étape 5 : Exécutez le job et vérifiez

  1. Cliquez sur Run now pour Trigger le Job.
  2. Sélectionnez le Runs tab pour voir l'exécution. La première exécution d'un Job prend quelques minutes pour start le compute ; une fois terminée, elle apparaît dans la liste.
  3. Cliquez sur le nœud process_segments pour développer la tâche For each.
  4. La page d’exécution affiche un tableau des itérations, une ligne par segment, chacune avec son statut, son heure de start et sa durée.
  5. Cliquez sur n’importe quelle ligne d’itération pour ouvrir sa sortie et confirmer qu’elle a analysé le segment attendu.

Vous pouvez voir les résultats de chaque itération indépendamment. Si une itération spécifique échoue, vous pouvez réexécuter uniquement cette itération depuis la page d'exécution du job sans avoir à réexécuter l'intégralité du job.

Étendre le modèle

Pour ajouter un segment à l'analyse, insérez une ligne dans la table de contrôle :

SQL
INSERT INTO <catalog-name>.config.property_segments VALUES ('Historical Place', 100);

La prochaine exécution du job inclut le nouveau segment, sans modification de la configuration du job ni modification du notebook.

Ce même schéma fonctionne pour tout cas où vous souhaitez que les données pilotent l’itération :

  • Traitement par client : une ligne par ID client. La tâche imbriquée applique des Transformations spécifiques aux clients ou effectue des livraisons vers des destinations spécifiques aux clients.
  • Ingestion de table : une ligne par nom de table source. La tâche imbriquée lit et ingère chaque table.
  • Traitement de reconstitution : une ligne par partition de date. La tâche imbriquée retraite les données historiques pour cette partition.
  • Exécution pilotée par indicateur de fonctionnalité : une ligne par fonctionnalité ou experimentation activée. La tâche imbriquée active la logique correspondante.

Pour arrêter le traitement d'une ligne sans la supprimer, ajoutez votre propre colonne à la table de contrôle (telle qu'un indicateur active) et filtrez dessus dans la tâche de recherche SQL. Il s'agit d'une colonne ordinaire que vous définissez et renseignez ; la tâche For each n'a aucun concept intégré de celle-ci. Ajoutez d'abord la colonne, puis définissez les lignes existantes sur TRUE:

SQL
ALTER TABLE <catalog-name>.config.property_segments ADD COLUMN active BOOLEAN;
UPDATE <catalog-name>.config.property_segments SET active = TRUE;

Filtrez ensuite dessus dans la query read_segments afin que seules les lignes actives pilotent l'itération :

SQL
SELECT property_type, min_price FROM <catalog-name>.config.property_segments WHERE active = TRUE;

Ressources supplémentaires