Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
Utiliser une table de contrôles pour piloter un
Lorsque vous exécutez le même traitement sur de nombreuses entrées, telles que les marchés, les tables sources, les clients ou les partitions de date, coder cette liste en dur dans votre poste signifie modifier le code et la redéployer à chaque changement de liste. À la place, stockez la liste dans une table de contrôle que le travail lit à l’exécution. Pour ajouter ou retirer du travail, vous mettez à jour une ligne dans la table, et la prochaine exécution du travail reprend le changement sans modifications apportées au travail lui-même. C’est un schéma piloté par les métadonnées : ce sont les données, et non le code, qui contrôlent ce que le poste traite.
Ce tutoriel construit un travail qui utilise ce motif sur le jeu de données d’échantillons Wanderbricks préinstallé, afin que vous puissiez l’exécuter de bout en bout sans créer de données sources. Le scénario est une plateforme de location de vacances qui effectue la même analyse de prix pour chaque segment immobilier (comme Ski Resort ou Urban Year-Round). Une table de contrôle liste les segments à analyser, une tâche SQL lit cette table, et une For each tâche exécute l’analyse une fois par segment, en parallèle.
Fonctionnement
Le travail relie trois tâches en séquence :
| Tâche | Type | Qu’est-ce que cela fait ? |
|---|---|---|
read_segments |
SQL | Lit la table de contrôle et capture les lignes comme un tableau JSON |
process_segments |
Pour chaque | Itère sur le tableau de lignes, lançant la tâche imbriquée une fois par ligne |
run_segment_analysis |
Notebook ou SQL (imbriqué à l’intérieur For each) |
S’exécute une fois par ligne, en utilisant les valeurs de cette ligne pour analyser un segment de propriété |
Le débit est read_segments → process_segments → run_segment_analysis (une fois par rangée). La sortie de la tâche SQL, un tableau JSON d’objets lignes, s’écoule dans le For each champ Inputs de la tâche via la référence {{tasks.read_segments.output.rows}}dynamique de la valeur . La For each tâche 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}}.
Prerequisites
- Un espace de travail Azure Databricks avec la permission de créer des emplois et des carnets.
- Permission de créer des tables dans le catalogue Unity, et permission de créer un schéma dans un catalogue (les
USE CATALOGprivilèges etCREATE SCHEMA) pour contenir la table de contrôle. - Un entrepôt SQL pour exécuter les tâches SQL. Si vous n’en avez pas, voir Créer un entrepôt SQL.
- Le
samplescatalogue, disponible dans tous les espaces de travail compatibles avec Unity Catalog. Le tutoriel lit depuissamples.wanderbricks.properties, donc il n’y a pas de données sources à 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 poste traite. Pour changer ce que fait le poste, vous mettez à jour ce tableau, pas le poste.
Exécutez le SQL suivant dans un notebook Azure 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 minimum de listage à inclure dans l’analyse de ce segment :
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-le <catalog-name> par un catalogue dans lequel vous pouvez créer des schémas, comme votre catalogue d’espace de travail. Utilisez le même catalogue partout où le tutoriel fait référence config.property_segments, y compris la requête de recherche à l’étape 3.
Après cette étape, config.property_segments contient trois rangées, une par segment. Chaque ligne porte les deux valeurs que le travail transmet à chaque itération : l’analyse property_type à et le min_price sol sur lequel filtrer.
Étape 2 : Écrire la logique d’analyse
La tâche imbriquée à l’intérieur de la For each tâche s’exécute une fois par ligne de la table de contrôle, recevant les paramètres de property_type cette ligne et min_price comme paramètres. Vous pouvez écrire cette logique comme une tâche notebook ou SQL. Choisissez en fonction de votre logique métier :
- Utilisez une tâche de 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 d’apprentissage automatique).
- Utilisez une tâche SQL lorsque la logique est une requête ou une transformation unique que vous pouvez exprimer de manière déclarative. Une tâche SQL nécessite un entrepôt SQL.
Les deux variantes ci-dessous produisent le même résultat : pour le segment traité, le nombre d’annonces à son prix plancher ou au-dessus et leur prix moyen.
Tâche de notebook
Créez un notebook dans un chemin d’accès tel que /Workspace/Users/<username>/run_segment_analysis. Ce carnet s’exécute une fois par itération de la For each tâche, recevant un segment différent à chaque fois.
Ajoutez le code suivant au notebook :
# 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={"property_type": property_type, "min_price": min_price},
)
display(result)
Note
Appelez dbutils.widgets.text() avant dbutils.widgets.get(). Si vous appelez get en premier, faire tourner le notebook en dehors d’un travail génère une InputWidgetNotDefined erreur.
Tâche SQL
Une tâche SQL exécute une requête enregistrée, donc créez et enregistrez la requête d’analyse dans l’éditeur SQL dès maintenant. Vous l’attachez à la tâche imbriquée lorsque vous configurez la For each tâche à l’étape 4.
Dans votre espace de travail Azure Databricks, cliquez sur
Nouveau>
Requête pour ouvrir l’éditeur SQL.
Entrez la requête suivante. Les tâches SQL font référence aux paramètres avec la
:param_namesyntaxe, donc la requête lit son segment et son plancher de prix à partir des:property_typeparamètres et:min_price: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;Cliquez sur le titre
New Query <date>dans l’onglet de votre fichier SQL, et donnez-lui le nomrun_segment_analysis. Ensuite, cliquez sur Enregistrer pour le déplacer dans un dossier où vous souhaitez le stocker.
La For each tâche transmet les valeurs de chaque itération aux :property_type paramètres et :min_price nommés à l’exécution. Contrairement aux widgets de carnet, les paramètres nommés SQL ne prennent pas en charge les valeurs par défaut : si un paramètre n’est pas passé, la requête échoue avec une erreur de résolution de paramètre.
Étape 3 : Créer la requête de recherche
La tâche de recherche lit la table de contrôle via une requête enregistrée. Comme à l’étape 2, créez et enregistrez la requête dans l’éditeur SQL dès maintenant, puis attachez-la à la tâche de recherche à l’étape 4.
Dans votre espace de travail Azure Databricks, cliquez sur
Nouveau>
Requête pour ouvrir l’éditeur SQL.
Entrez ce qui suit, en utilisant le même catalogue que vous avez choisi à l’étape 1 :
SELECT property_type, min_price FROM <catalog-name>.config.property_segments;Le nom est entièrement précisé car l’entrepôt SQL qui exécute cette requête peut par défaut se tourner vers un catalogue différent de celui dans lequel vous avez créé la table.
Cliquez sur le titre
New Query <date>dans l’onglet de votre fichier SQL, et donnez-lui le nomread_segments. Ensuite, cliquez sur Enregistrer pour le déplacer dans un dossier où vous souhaitez le stocker.
Étape 4 : Créer et configurer le travail
Avec les deux requêtes sauvegardées, créez le travail et ajoutez ses deux tâches : la tâche de recherche SQL qui lit la table de contrôle, et la For each tâche qui exécute l’analyse pour chaque ligne.
Créez le poste
Dans votre espace de travail Azure Databricks, dans la barre latérale, cliquez Nouveau>
Job. Donnez au poste un nom descriptif, comme
Segment Analysis.
Configurez la tâche de recherche SQL
Cette tâche lit la table de contrôle et rend ses lignes accessibles à la For each tâche en exécutant la read_segments requête que vous avez enregistrée à l’étape 3.
- Cliquez sur la tuile de requête SQL pour configurer la première tâche. Si la tuile de requête SQL n’est pas disponible, cliquez sur Ajouter un autre type de tâche et recherchez requête SQL.
- Définissez le nom de la tâche sur
read_segments. - Si nécessaire, sélectionnez requête SQL dans le menu déroulant Type .
- Dans le champ requête SQL , sélectionnez la
read_segmentsrequête que vous avez enregistrée à l’étape 3. - Définissez SQL warehouse comme un entrepôt dans votre espace de travail.
- Cliquez sur Create task.
Lorsque cette tâche s’exécute, Azure Databricks capture le résultat sous forme de tableau JSON dans tasks.read_segments.output.rows. La sortie de tâche SQL est toujours renvoyée sous forme de tableau JSON, donc vous n’avez pas besoin de 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 la tâche que vous avez définie. Le résultat se présente ainsi :
[
{ "property_type": "Urban Year-Round", "min_price": 150 },
{ "property_type": "Summer Getaway", "min_price": 200 },
{ "property_type": "Ski Resort", "min_price": 250 }
]
Configurez la For each tâche
La For each tâche lit la sortie SQL et lance une tâche imbriquée exécutée par ligne.
Ajouter une tâche et sélectionner Pour chacun.
Définissez le nom de la tâche sur
process_segments.Vérifiez que Depends on est défini à
read_segments.Dans le champ Entrées , entrez le tableau de lignes capturé par la tâche SQL :
{{tasks.read_segments.output.rows}}Réglez la concurrence pour
2exécuter deux itérations en parallèle. Augmentez cette valeur lorsque votre tâche imbriquée prend en charge un parallélisme plus élevé.Pour accomplir cette tâche, cliquez sur Ajouter une tâche à boucler et configurez la tâche imbriquée qui s’exécute sur chaque itération.
La For each tâche et sa tâche imbriquée sont créées ensemble comme une seule tâche. Configurez la tâche imbriquée selon le type que vous avez choisi à l’étape 2 :
Tâche de notebook
Définissez le nom de la tâche sur
run_segment_analysis.Définissez Type sur Notebook.
Définissez le chemin du carnet que vous avez créé à l’étape 2.
Cliquez sur Paramètres, puis cliquez sur Ajouter pour ajouter chaque paramètre :
-
Clé :
property_type, Valeur :{{input.property_type}} -
Clé :
min_price, Valeur :{{input.min_price}}
Chaque
{{input.<key>}}référence se résout dans le champ correspondant de la ligne de l’itération en cours.-
Clé :
Cliquez sur Créer une tâche pour créer la
For eachtâche et sa tâche imbriquée ensemble.
Tâche SQL
Cette tâche exécute la run_segment_analysis requête que vous avez enregistrée à l’étape 2.
Définissez le nom de la tâche sur
run_segment_analysis.Définissez Type sur SQL, puis définissez la tâche SQL sur Requête.
Dans le champ de requête SQL , sélectionnez la
run_segment_analysisrequête que vous avez enregistrée à l’étape 2.Définissez SQL warehouse comme un entrepôt dans votre espace de travail.
Cliquez sur Paramètres, puis cliquez sur Ajouter pour ajouter chaque paramètre :
-
Clé :
property_type, Valeur :{{input.property_type}} -
Clé :
min_price, Valeur :{{input.min_price}}
Chaque
{{input.<key>}}référence se résout dans le champ correspondant de la ligne de l’itération en cours.-
Clé :
Cliquez sur Créer une tâche pour créer la
For eachtâche et sa tâche imbriquée ensemble.
Votre graphe acyclique dirigé (DAG) de travail affiche read_segments désormais le flux dans process_segments, avec la tâche imbriquée à l’intérieur du For each nœud.
Étape 5 : Exécutez le travail et vérifiez
- Cliquez sur Exécuter maintenant pour déclencher le travail.
- Sélectionnez l’onglet Runs pour voir la course. La première exécution d’un travail prend quelques minutes pour commencer le calcul ; Une fois terminé, il apparaît dans la liste.
- Cliquez sur le
process_segmentsnœud pour développer laFor eachtâche. - La page de la course montre un tableau d’itérations, une ligne par segment, chacune avec son statut, son heure de début et sa durée.
- 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 ne pouvez relancer que cette itération depuis la page d’exécution du travail sans relancer toute la tâche.
Étendre le modèle
Pour ajouter un segment à l’analyse, insérez une ligne dans la table de contrôle :
INSERT INTO <catalog-name>.config.property_segments VALUES ('Historical Place', 100);
La prochaine exécution du travail inclut le nouveau segment, sans aucun changement de configuration du travail ni modification du carnet.
Ce même schéma fonctionne dans tous les 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 au client ou livre à des destinations spécifiques à chaque client.
- Ingestion de table : une ligne par nom de table source. La tâche imbriquée lit et ingère chaque table.
- Traitement de renvoi : une ligne pour chaque partition de date. La tâche imbriquée retraite les données historiques de cette partition.
- Exécution pilotée par indicateur de fonctionnalité : une ligne par fonctionnalité ou expérience activée. La tâche imbriquée active la logique correspondante.
Pour arrêter de traiter une ligne sans la supprimer, ajoutez votre propre colonne à la table de contrôle (comme un active drapeau) et filtrez dessus dans la tâche de recherche SQL. C’est une colonne ordinaire que vous définissez et peuplez ; La For each tâche n’a pas de concept intégré. Ajoutez d’abord la colonne, puis définissez les lignes existantes à TRUE:
ALTER TABLE <catalog-name>.config.property_segments ADD COLUMN active BOOLEAN;
UPDATE <catalog-name>.config.property_segments SET active = TRUE;
Ensuite, filtrez dessus dans la read_segments requête pour que seules les lignes actives guident l’itération :
SELECT property_type, min_price FROM <catalog-name>.config.property_segments WHERE active = TRUE;
Ressources supplémentaires
-
Utiliser une tâche pour exécuter une
For eachautre tâche dans une boucle : référence complète pour la configurationFor eachdes tâches, y compris les types de paramètres et les options d’accès concurrentiel -
Utiliser une table de recherche pour les tableaux de paramètres volumineux d’une
For eachtâche : comment gérer des tableaux de paramètres volumineux qui dépassent la limite de valeur de tâche de 48 Ko - Accéder aux valeurs des paramètres à partir d’une tâche : toutes les méthodes d’accès aux valeurs de paramètres dans les notebooks, les scripts Python et les tâches SQL
- Jeu de données Wanderbricks : L’échantillon de jeu utilisé dans ce tutoriel