Aller au contenu principal

Utiliser une table de recherche pour les grands tableaux de paramètres dans une tâche For each

For each Les tâches transmettent des tableaux de paramètres aux tâches imbriquées de manière itérative, afin de donner à chaque tâche imbriquée les informations nécessaires à son exécution. Les tableaux de paramètres sont limités à 5 000 caractères, ou 48 Ko si vous utilisez des références de valeurs de tâche pour les transmettre. Lorsque vous avez une plus grande quantité de données à transmettre aux tâches imbriquées, vous ne pouvez pas utiliser directement l'entrée, les valeurs de tâche ou les paramètres de job pour transmettre ces données.

Une alternative à la transmission des données complètes consiste à stocker les données de tâche sous forme de fichier JSON, et à transmettre une clé de recherche (dans les données JSON) via l’entrée de tâche au lieu des données complètes. Les tâches imbriquées peuvent utiliser la clé pour récupérer les données spécifiques nécessaires à chaque itération.

L'exemple suivant montre un exemple de fichier de configuration JSON, et comment transmettre des parameter à une tâche imbriquée qui recherche les valeurs dans la configuration JSON.

Exemple de configuration JSON

Cette configuration d’exemple est une liste d’étapes, avec des paramètres (args) pour chaque itération (seules trois étapes sont présentées pour cet exemple). Supposons que ce fichier JSON est enregistré sous /Workspace/Users/<user>/copy-filtered-table-config.json. Nous faisons référence à cela dans la tâche imbriquée.

JSON
{
"steps": [
{
"key": "table_1",
"args": {
"catalog": "my-catalog",
"schema": "my-schema",
"source_table": "raw_data_table_1",
"destination_table": "filtered_table_1",
"filter_column": "col_a",
"filter_value": "value_1"
}
},
{
"key": "table_2",
"args": {
"catalog": "my-catalog",
"schema": "my-schema",
"source_table": "raw_data_table_2",
"destination_table": "filtered_table_2",
"filter_column": "col_b",
"filter_value": "value_2"
}
},
{
"key": "table_3",
"args": {
"catalog": "my-catalog",
"schema": "my-schema",
"source_table": "raw_data_table_3",
"destination_table": "filtered_table_3",
"filter_column": "col_c",
"filter_value": "value_3"
}
}
]
}

Exemple de tâche For each

La tâche For each de votre Job comprend une entrée avec les clés pour chaque itération. Cet exemple montre une tâche appelée copy-filtered-tables avec les entrées définies sur ["table_1","table_2","table_3"]. Cette liste est limitée à 5 000 caractères, mais comme vous ne transmettez que des clés, elle est bien plus petite que les données complètes.

Dans cet exemple, les étapes ne dépendent pas d'autres étapes ou tâches, nous pouvons donc définir une concurrence supérieure à 1 pour accélérer l'exécution de la tâche.

Pour chaque tâche, indiquant les entrées et la concurrence

Exemple de tâche imbriquée

La tâche imbriquée reçoit l'entrée de la tâche parente For each. Dans ce cas, nous configurons l'entrée à utiliser comme Key pour le fichier de configuration. L'image suivante montre la tâche imbriquée, y compris la configuration d'un Parameter appelé key avec la valeur {{input}}.

Pour chaque tâche imbriquée, montrant comment utiliser les entrées

Cette tâche est un Notebook qui contient du code. Dans votre Notebook, vous pouvez utiliser le code Python suivant pour lire l'entrée et l'utiliser comme clé dans le fichier JSON de configuration. Les données du fichier JSON sont utilisées pour lire, filtrer et écrire des données d'une table.

remarque

Pour voir un exemple de pilotage d'une tâche For each à partir d'une table Delta en direct au lieu d'un fichier JSON statique, consultez Utiliser une table de contrôle pour piloter un job For each.

Python
# copy-filtered-table (iteratable task code to read a table, filter by a value, and write as a new table)

from pyspark.sql.functions import expr
from types import SimpleNamespace

import json

# If the notebook is run outside of a job with a key parameter, this provides
# a default. This allows testing outside of a For each task
dbutils.widgets.text("key", "table_1", "key")

# load configuration (note that the path must be set to valid configuration file)
config_path = "/Workspace/Users/<user>/copy-filtered-table-config.json"
with open(config_path, "r") as file:
config = json.loads(file.read())

# look up step and arguments
key = dbutils.widgets.get("key")
current_step = next((step for step in config['steps'] if step['key'] == key), None)
if current_step is None:
raise ValueError(f"Could not find step '{key}' in the configuration")

args = SimpleNamespace(**current_step["args"])

# read the source table defined for the step, and filter it
df = spark.read.table(f"{args.catalog}.{args.schema}.{args.source_table}") \
.filter(expr(f"{args.filter_column} like '%{args.filter_value}%'"))

# write the filtered table to the destination
df.write.mode("overwrite").saveAsTable(f"{args.catalog}.{args.schema}.{args.destination_table}")