Orchestrez les notebooks Databricks et modularisez le code
Vous pouvez orchestrer les Notebooks Databricks et modulariser le code à l'aide de Lakeflow Jobs, dbutils.notebook.run(), des fichiers de l'Workspace et %run. Choisissez une méthode selon votre besoin en matière de planification, de transmission de paramètre et de contrôle de version.
Méthodes d’orchestration et de modularisation du code
Le tableau suivant compare les méthodes disponibles pour orchestrer les Notebook et modulariser le code dans les Notebook.
Méthode | Cas d'usage | Notes |
|---|---|---|
Orchestration de Notebook (recommandée) | Méthode recommandée pour l’orchestration de Notebook. Prend en charge les workflows complexes avec des dépendances de tâches, la planification et des triggers. Offre une approche robuste et évolutive pour les charges de travail de production, mais nécessite une installation et une configuration. | |
Orchestration de Notebook | Utilisez Démarre un nouveau Job éphémère pour chaque appel, ce qui peut augmenter les frais généraux et ne dispose pas de fonctionnalités de planification avancées. | |
Modularisation du code (recommandé) | Méthode recommandée pour modulariser le code. Modularisez le code en fichiers de code réutilisables stockés dans le Workspace. Prend en charge le contrôle de version avec les dépôts et l’intégration aux IDEs pour un meilleur debugging et des tests unitaires. Nécessite une configuration supplémentaire pour gérer les chemins de fichiers et les dépendances. | |
Modularisation du code | Utilisez Importez des fonctions ou des variables d'autres Notebooks en les exécutant en ligne. Utile pour le prototypage, mais peut conduire à un code étroitement couplé, plus difficile à maintenir. Ne prend pas en charge le passage de paramètres ou le contrôle de version. |
%run vs. dbutils.notebook.run()
La commande %run vous permet d'inclure un autre Notebook dans un Notebook. Vous pouvez utiliser %run pour modulariser votre code en plaçant les fonctions de support dans un Notebook distinct. Vous pouvez également l'utiliser pour concaténer des Notebooks qui implémentent les étapes d'une analyse. Lorsque vous utilisez %run, le Notebook appelé est immédiatement exécuté et les fonctions et variables qui y sont définies deviennent disponibles dans le Notebook appelant.
L'API dbutils.notebook complète %run car elle vous permet de passer des paramètres à un notebook et d'en renvoyer des valeurs. Cela vous permet de créer des workflows et des pipelines complexes avec des dépendances. Par exemple, vous pouvez obtenir une liste de fichiers dans un répertoire et transmettre les noms à un autre Notebook, ce qui est impossible avec %run. Vous pouvez également créer des workflows conditionnels (if-then-else) basés sur les valeurs de retour.
Contrairement à %run, la méthode dbutils.notebook.run() start un nouveau Job pour exécuter le Notebook.
Comme toutes les APIs dbutils, ces méthodes sont disponibles uniquement en Python et Scala. Cependant, vous pouvez utiliser dbutils.notebook.run() pour appeler un notebook R.
Utiliser %run pour importer un Notebook.
Dans cet exemple, le premier Notebook définit une fonction, reverse, qui est disponible dans le deuxième Notebook après que vous ayez utilisé le magic %run pour exécuter shared-code-notebook.


Comme les deux Notebooks se trouvent dans le même répertoire du Workspace, utilisez le préfixe ./ dans ./shared-code-notebook pour indiquer que le chemin doit être résolu par rapport au Notebook en cours d'exécution. Vous pouvez organiser les Notebooks en répertoires, tels que %run ./dir/notebook, ou utiliser un chemin absolu comme %run /Users/username@organization.com/directory/notebook.
%rundoit se trouver dans une cellule *à elle seule*, car il exécute l'intégralité du Notebook en ligne.- Vous ne pouvez *pas* utiliser
%runpour exécuter un fichier Python etimportles entités définies dans ce fichier dans un Notebook. Pour importer depuis un fichier Python, consultez Modulariser votre code à l'aide de fichiers. Ou, conditionnez le fichier en une bibliothèque Python, créez une bibliothèque Databricks à partir de cette bibliothèque Python, et installez la bibliothèque dans le cluster que vous utilisez pour exécuter votre Notebook. - Lorsque vous utilisez
%runpour exécuter un Notebook qui contient des widgets, by default le Notebook spécifié s'exécute avec les valeurs default du widget. Vous pouvez également transmettre des valeurs aux widgets ; consultez Utiliser des widgets Databricks avec %run.
Utilisez dbutils.notebook.run pour start un nouveau Job
Exécuter un Notebook et retourner sa valeur de sortie. La méthode start un job éphémère qui s'exécute immédiatement.
Les méthodes disponibles dans l'API dbutils.notebook sont run et exit. Les paramètres et les valeurs de retour doivent être des chaînes.
run(path: String, timeout_seconds: int, arguments: Map): String
Le paramètre timeout_seconds contrôle le délai d'expiration de l'exécution (0 signifie pas de délai d'expiration). L'appel à
run lève une exception s'il ne se termine pas dans le délai spécifié. Si Databricks est en panne pendant plus de 10 minutes,
l'exécution du Notebook échoue quel que soit timeout_seconds.
Le paramètre arguments définit les valeurs des widgets du Notebook cible. Plus précisément, si le Notebook que vous exécutez possède un widget
nommé A, et que vous passez une paire clé-valeur ("A": "B") dans le cadre du paramètre arguments à l'appel run(),
alors la récupération de la valeur du widget A renverra "B". Vous trouverez les instructions pour créer et
travailler avec des widgets sur la page des widgets Databricks.
- Le paramètre
argumentsn'accepte que les caractères latins (jeu de caractères ASCII). L'utilisation de caractères non-ASCII renvoie une erreur. - Les jobs créés à l’aide de l’API
dbutils.notebookdoivent être terminés en 30 jours ou moins.
run utilisation
- Python
- Scala
dbutils.notebook.run("notebook-name", 60, {"argument": "data", "argument2": "data2", ...})
dbutils.notebook.run("notebook-name", 60, Map("argument" -> "data", "argument2" -> "data2", ...))
Transférez des données structurées entre les Notebooks
Cette section illustre comment transmettre des données structurées entre notebooks.
- Python
- Scala
# Example 1 - returning data through temporary views.
# You can only return one string using dbutils.notebook.exit(), but since called notebooks reside in the same JVM, you can
# return a name referencing data stored in a temporary view.
## In callee notebook
spark.range(5).toDF("value").createOrReplaceGlobalTempView("my_data")
dbutils.notebook.exit("my_data")
## In caller notebook
returned_table = dbutils.notebook.run("LOCATION_OF_CALLEE_NOTEBOOK", 60)
global_temp_db = spark.conf.get("spark.sql.globalTempDatabase")
display(table(global_temp_db + "." + returned_table))
# Example 2 - returning data through DBFS.
# For larger datasets, you can write the results to DBFS and then return the DBFS path of the stored data.
## In callee notebook
dbutils.fs.rm("/tmp/results/my_data", recurse=True)
spark.range(5).toDF("value").write.format("parquet").save("dbfs:/tmp/results/my_data")
dbutils.notebook.exit("dbfs:/tmp/results/my_data")
## In caller notebook
returned_table = dbutils.notebook.run("LOCATION_OF_CALLEE_NOTEBOOK", 60)
display(spark.read.format("parquet").load(returned_table))
# Example 3 - returning JSON data.
# To return multiple values, you can use standard JSON libraries to serialize and deserialize results.
## In callee notebook
import json
dbutils.notebook.exit(json.dumps({
"status": "OK",
"table": "my_data"
}))
## In caller notebook
import json
result = dbutils.notebook.run("LOCATION_OF_CALLEE_NOTEBOOK", 60)
print(json.loads(result))
Serverless compatibility Databricks vous recommande d'abandonner les APIs RDD, car elles ne sont pas compatibles avec l'architecture de compute serverless de Databricks. Utilisez l'API DataFrame à la place.
// Example 1 - returning data through temporary views.
// You can only return one string using dbutils.notebook.exit(), but since called notebooks reside in the same JVM, you can
// return a name referencing data stored in a temporary view.
/** In callee notebook */
sc.parallelize(1 to 5).toDF().createOrReplaceGlobalTempView("my_data")
dbutils.notebook.exit("my_data")
/** In caller notebook */
val returned_table = dbutils.notebook.run("LOCATION_OF_CALLEE_NOTEBOOK", 60)
val global_temp_db = spark.conf.get("spark.sql.globalTempDatabase")
display(table(global_temp_db + "." + returned_table))
// Example 2 - returning data through DBFS.
// For larger datasets, you can write the results to DBFS and then return the DBFS path of the stored data.
/** In callee notebook */
dbutils.fs.rm("/tmp/results/my_data", recurse=true)
sc.parallelize(1 to 5).toDF().write.format("parquet").save("dbfs:/tmp/results/my_data")
dbutils.notebook.exit("dbfs:/tmp/results/my_data")
/** In caller notebook */
val returned_table = dbutils.notebook.run("LOCATION_OF_CALLEE_NOTEBOOK", 60)
display(sqlContext.read.format("parquet").load(returned_table))
// Example 3 - returning JSON data.
// To return multiple values, use standard JSON libraries to serialize and deserialize results.
/** In callee notebook */
// Import jackson json libraries
import com.fasterxml.jackson.module.scala.DefaultScalaModule
import com.fasterxml.jackson.module.scala.experimental.ScalaObjectMapper
import com.fasterxml.jackson.databind.ObjectMapper
// Create a json serializer
val jsonMapper = new ObjectMapper with ScalaObjectMapper
jsonMapper.registerModule(DefaultScalaModule)
// Exit with json
dbutils.notebook.exit(jsonMapper.writeValueAsString(Map("status" -> "OK", "table" -> "my_data")))
/** In caller notebook */
// Import jackson json libraries
import com.fasterxml.jackson.module.scala.DefaultScalaModule
import com.fasterxml.jackson.module.scala.experimental.ScalaObjectMapper
import com.fasterxml.jackson.databind.ObjectMapper
// Create a json serializer
val jsonMapper = new ObjectMapper with ScalaObjectMapper
jsonMapper.registerModule(DefaultScalaModule)
val result = dbutils.notebook.run("LOCATION_OF_CALLEE_NOTEBOOK", 60)
println(jsonMapper.readValue[Map[String, String]](result))
Gérer les erreurs
Cette section illustre comment gérer les erreurs.
- Python
- Scala
# Errors throw a WorkflowException.
def run_with_retry(notebook, timeout, args = {}, max_retries = 3):
num_retries = 0
while True:
try:
return dbutils.notebook.run(notebook, timeout, args)
except Exception as e:
if num_retries > max_retries:
raise e
else:
print("Retrying error", e)
num_retries += 1
run_with_retry("LOCATION_OF_CALLEE_NOTEBOOK", 60, max_retries = 5)
// Errors throw a WorkflowException.
import com.databricks.WorkflowException
// Since dbutils.notebook.run() is just a function call, you can retry failures using standard Scala try-catch
// control flow. Here, we show an example of retrying a notebook a number of times.
def runRetry(notebook: String, timeout: Int, args: Map[String, String] = Map.empty, maxTries: Int = 3): String = {
var numTries = 0
while (true) {
try {
return dbutils.notebook.run(notebook, timeout, args)
} catch {
case e: WorkflowException if numTries < maxTries =>
println("Error, retrying: " + e)
}
numTries += 1
}
"" // not reached
}
runRetry("LOCATION_OF_CALLEE_NOTEBOOK", timeout = 60, maxTries = 5)
Exécutez plusieurs Notebooks simultanément
Vous pouvez exécuter plusieurs notebooks simultanément en utilisant des constructions Scala et Python standard telles que les Threads (Scala, Python) et les Futures (Scala, Python). Les notebooks d’exemple montrent comment utiliser ces constructions.
- download les quatre notebooks suivants. Les notebooks sont écrits en Scala.
- Importez les notebooks dans un dossier unique du Workspace.
- Exécutez le Notebook Exécuter simultanément .