Didacticiel : Créez un pipeline ETL avec Apache Spark sur la plateforme Databricks
Ce tutoriel vous montre comment développer et déployer votre premier pipeline ETL (extraction, transformation et chargement) pour l'orchestration de données avec Apache Spark. Bien que ce tutoriel utilise le compute universel Databricks, vous pouvez également utiliser le compute serverless s'il est activé pour votre workspace.
Vous pouvez également utiliser les LakeFlow Pipelines pour créer des pipelines ETL. Les LakeFlow Pipelines réduisent la complexité de la création, du déploiement et de la maintenance des pipelines ETL de production. Voir Tutoriel : créer un pipeline ETL avec les LakeFlow Pipelines.
À la fin de cet article, vous saurez comment :
- Lancer une Ressource de compute tout usage Databricks.
- Créer un Notebook Databricks.
- Configurez l’ingestion incrémentielle de données dans Delta Lake avec Auto Loader.
- Traiter et interagir avec les données.
- Planifier un Notebook comme Job Databricks.
Ce tutoriel utilise des notebooks interactifs pour effectuer des tâches ETL courantes en Python ou Scala.
Vous pouvez également utiliser le fournisseur Databricks Terraform pour créer les ressources de cet article. Consultez Créer des clusters, des Notebook et des Job avec Terraform.
Exigences
- Vous êtes connecté(e) à un Databricks Workspace.
- Vous avez l'autorisation de créer une ressource de compute.
Si vous ne disposez pas des privilèges de contrôle compute, vous pouvez toujours effectuer la plupart des étapes ci-dessous tant que vous avez accès à une ressource compute.
Étape 1 : Créer une ressource de compute
Pour effectuer une analyse exploratoire des données et de l'ingénierie des données, créez une ressource de compute pour exécuter les commandes.
Si votre Workspace est activé pour le compute serverless, vous pouvez ignorer cette étape. Les Notebooks se rattachent automatiquement au compute serverless à l’exécution du code, ou vous pouvez sélectionner Serverless dans la liste déroulante du compute. Consultez Compute serverless pour les notebooks.
- Cliquez sur
Compute dans la barre latérale.
- Sur la page Compute, cliquez sur Créer un compute .
- Spécifiez un nom unique pour la ressource de compute, laissez les valeurs restantes dans leur état par default et cliquez sur Créer un compute .
Pour en savoir plus sur le compute Databricks, consultez Compute.
Étape 2 : Créer un notebook Databricks
Pour créer un Notebook dans votre Workspace, cliquez sur Nouveau dans la barre latérale, puis cliquez sur Notebook . Un notebook vierge s'ouvre dans le Workspace.
Pour en savoir plus sur la création et la gestion des Notebooks, veuillez consulter Gérer les Notebooks Databricks.
Étape 3 : configurez Auto Loader pour ingérer des données dans Delta Lake
Databricks vous recommande d'utiliser Auto Loader pour l'ingestion incrémentielle de données. Auto Loader détecte et traite automatiquement les nouveaux fichiers au fur et à mesure qu'ils arrivent dans le stockage d'objets cloud.
Databricks recommande de stocker les données avec Delta Lake. Delta Lake est une couche de stockage open source qui fournit des transactions ACID et permet le data lakehouse. Delta Lake est le format default pour les tables créées dans Databricks.
Pour configurer Auto Loader pour ingérer des données dans une table Delta Lake, copiez et collez le code suivant dans la cellule vide de votre Notebook :
- Python
- Scala
# Import functions
from pyspark.sql.functions import col, current_timestamp
# Define variables used in code below
file_path = "/databricks-datasets/structured-streaming/events"
username = spark.sql("SELECT regexp_replace(session_user(), '[^a-zA-Z0-9]', '_')").first()[0]
table_name = f"{username}_etl_quickstart"
checkpoint_path = f"/tmp/{username}/_checkpoint/etl_quickstart"
# Clear out data from previous demo execution
spark.sql(f"DROP TABLE IF EXISTS {table_name}")
dbutils.fs.rm(checkpoint_path, True)
# Configure Auto Loader to ingest JSON data to a Delta table
(spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", checkpoint_path)
.load(file_path)
.select("*", col("_metadata.file_path").alias("source_file"), current_timestamp().alias("processing_time"))
.writeStream
.option("checkpointLocation", checkpoint_path)
.trigger(availableNow=True)
.toTable(table_name))
// Imports
import org.apache.spark.sql.functions.current_timestamp
import org.apache.spark.sql.streaming.Trigger
import spark.implicits._
// Define variables used in code below
val file_path = "/databricks-datasets/structured-streaming/events"
val username = spark.sql("SELECT regexp_replace(session_user(), '[^a-zA-Z0-9]', '_')").first.get(0)
val table_name = s"${username}_etl_quickstart"
val checkpoint_path = s"/tmp/${username}/_checkpoint"
// Clear out data from previous demo execution
spark.sql(s"DROP TABLE IF EXISTS ${table_name}")
dbutils.fs.rm(checkpoint_path, true)
// Configure Auto Loader to ingest JSON data to a Delta table
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", checkpoint_path)
.load(file_path)
.select($"*", $"_metadata.file_path".as("source_file"), current_timestamp.as("processing_time"))
.writeStream
.option("checkpointLocation", checkpoint_path)
.trigger(Trigger.AvailableNow)
.toTable(table_name)
Les variables définies dans ce code devraient vous permettre de l'exécuter en toute sécurité sans risque de conflit avec les assets de Workspace existants ou d'autres utilisateurs. Les autorisations réseau ou de stockage restreintes généreront des erreurs lors de l'exécution de ce code ; contactez votre administrateur de workspace pour résoudre ces restrictions.
Pour en savoir plus sur Auto Loader, consultez Qu'est-ce qu'Auto Loader ?.
Étape 4 : Traiter et interagir avec les données
Les notebooks exécutent la logique cellule par cellule. Pour exécuter la logique dans votre cellule :
-
Pour exécuter la cellule que vous avez terminée à l'étape précédente, sélectionnez la cellule et appuyez sur MAJ+ENTRÉE .
-
Pour interroger la table que vous venez de créer, copiez et collez le code suivant dans une cellule vide, puis appuyez sur MAJ+ENTRÉE pour exécuter la cellule.
- Python
- Scala
df = spark.read.table(table_name)
val df = spark.read.table(table_name)
- Pour prévisualiser les données de votre DataFrame, copiez et collez le code suivant dans une cellule vide, puis appuyez sur SHIFT+ENTER pour exécuter la cellule.
- Python
- Scala
display(df)
display(df)
Pour en savoir plus sur les options interactives pour la visualisation de données, voir Visualisations dans les Notebooks et l'éditeur SQL de Databricks.
Étape 5 : planifier un Job
Vous pouvez exécuter des Notebooks Databricks en tant que scripts de production en les ajoutant comme tâche dans un Job Databricks. À cette étape, vous allez créer un nouveau Job que vous pourrez Trigger manuellement.
Si vous utilisez le compute serverless, sélectionnez Serverless dans le menu déroulant Compute au lieu de la ressource de compute de l'étape 1.
Pour planifier votre notebook en tant que tâche :
- Cliquez sur Planifier à droite de la barre d'en-tête.
- Saisissez un nom unique pour le nom du Job .
- Cliquez sur Manuel .
- Dans le menu déroulant compute , sélectionnez la ressource de compute que vous avez créée à l'étape 1.
- Cliquez sur Créer .
- Dans la fenêtre qui apparaît, cliquez sur Exécuter maintenant .
- Pour voir les résultats de l'exécution du job, cliquez
sur l'icône à côté du Timestamp de la dernière exécution.
Pour plus d'informations sur les Jobs, consultez Qu'est-ce qu'un Job ?.
Intégrations supplémentaires
En savoir plus sur les intégrations et les outils pour le Data Engineering avec Databricks :