Aller au contenu principal

Tutoriel : Créez un pipeline ETL avec LakeFlow Pipelines.

Ce tutoriel explique comment créer et déployer un pipeline ETL (extraction, transformation et chargement) pour l'orchestration de données à l'aide des LakeFlow Pipelines et d'Auto Loader. Un pipeline ETL met en œuvre les étapes pour lire les données des systèmes sources, transformer ces données en fonction des exigences, telles que les contrôles de qualité des données et la déduplication des enregistrements, et écrire les données vers un système cible, tel qu'un data warehouse ou un data lake.

Dans ce tutoriel, vous utiliserez des pipelines et Auto Loader pour :

  • Ingérer les données source brutes dans une table cible.
  • Transformez les données sources brutes et écrivez les données transformées dans deux vues matérialisées cibles.
  • Interrogez les données transformées.
  • Automatisez le pipeline ETL avec un Job Databricks.

Pour plus d’information sur les pipelines et Auto Loader, consultez Spark Declarative Pipelines et Qu’est-ce qu’Auto Loader ?

Exigences

Pour suivre ce didacticiel, vous devez satisfaire aux exigences suivantes :

À propos du dataset

Le jeu de données utilisé dans cet exemple est un sous-ensemble du Million Song Dataset, une collection de fonctionnalités et de métadonnées pour des titres de musique contemporaine. Ce dataset est disponible dans les datasets d'exemple inclus dans votre Workspace Databricks.

Étape 1 : Créer un pipeline

Tout d'abord, créez un pipeline en définissant les datasets dans des fichiers (appelés code source ) à l'aide de la syntaxe de pipeline. Chaque fichier de code source ne peut contenir qu'une seule langue, mais vous pouvez ajouter plusieurs fichiers spécifiques à la langue dans le pipeline. Pour en savoir plus, consultez Spark Declarative Pipelines

Ce tutoriel utilise le compute Serverless et Unity Catalog. Pour toutes les options de configuration non spécifiées, utilisez les paramètres default. Si le compute Serverless n'est pas activé ou pris en charge dans votre workspace, vous pouvez suivre le tutoriel tel quel en utilisant les paramètres de compute default.

Pour créer un nouveau pipeline, suivez ces étapes :

  1. Dans votre Workspace, cliquez sur Icône Plus. Nouveau dans la barre latérale, puis sélectionnez Pipeline ETL . Ceci ouvre l'éditeur de pipeline avec un nom de pipeline default comme New Pipeline <date> <time>.
  2. (Facultatif) Sélectionnez le nom et entrez un nom descriptif pour le pipeline.
  3. (Facultatif) À droite du nom, cliquez sur le catalogue et le schéma pour définir différents default.
  4. (Facultatif) Dans le fichier source my_transformation créé pour vous, sélectionnez Python ou SQL dans la liste déroulante des langues pour définir la langue du fichier.
  5. Cliquez Icône de code. sur **Utiliser l'exemple de code**.

Un exemple de code dans votre langue sélectionnée apparaît dans le fichier my_transformation du dossier transformations.

Étape 2 : Développer la logique de votre pipeline

À cette étape, vous utiliserez l’ éditeur LakeFlow Pipelines pour développer et valider le code source du pipeline de manière interactive.

Le code utilise 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. Pour en savoir plus, consultez Qu'est-ce qu'Auto Loader ?

Un exemple de fichier de code source a été créé dans le dossier Transformations de votre pipeline. By default, tous les fichiers *.py et *.sql dans le dossier Transformations font partie de la source de votre pipeline.

  1. Dans le dossier Transformations, ouvrez le fichier de transformation d'exemple et remplacez l'exemple de code par ce qui suit. Veillez à utiliser la langue que vous avez sélectionnée à l'étape 1.
Python
# Import modules
from pyspark import pipelines as dp
from pyspark.sql.functions import *
from pyspark.sql.types import DoubleType, IntegerType, StringType, StructType, StructField

# Define the path to the source data
file_path = f"/databricks-datasets/songs/data-001/"

# Define a streaming table to ingest data from a volume
schema = StructType(
[
StructField("artist_id", StringType(), True),
StructField("artist_lat", DoubleType(), True),
StructField("artist_long", DoubleType(), True),
StructField("artist_location", StringType(), True),
StructField("artist_name", StringType(), True),
StructField("duration", DoubleType(), True),
StructField("end_of_fade_in", DoubleType(), True),
StructField("key", IntegerType(), True),
StructField("key_confidence", DoubleType(), True),
StructField("loudness", DoubleType(), True),
StructField("release", StringType(), True),
StructField("song_hotnes", DoubleType(), True),
StructField("song_id", StringType(), True),
StructField("start_of_fade_out", DoubleType(), True),
StructField("tempo", DoubleType(), True),
StructField("time_signature", DoubleType(), True),
StructField("time_signature_confidence", DoubleType(), True),
StructField("title", StringType(), True),
StructField("year", IntegerType(), True),
StructField("partial_sequence", IntegerType(), True)
]
)

@dp.table(
comment="Raw data from a subset of the Million Song Dataset; a collection of features and metadata for contemporary music tracks."
)
def songs_raw():
return (spark.readStream
.format("cloudFiles")
.schema(schema)
.option("cloudFiles.format", "csv")
.option("sep","\t")
.load(file_path))

# Define a materialized view that validates data and renames a column
@dp.materialized_view(
comment="Million Song Dataset with data cleaned and prepared for analysis."
)
@dp.expect("valid_artist_name", "artist_name IS NOT NULL")
@dp.expect("valid_title", "song_title IS NOT NULL")
@dp.expect("valid_duration", "duration > 0")
def songs_prepared():
return (
spark.read.table("songs_raw")
.withColumnRenamed("title", "song_title")
.select("artist_id", "artist_name", "duration", "release", "tempo", "time_signature", "song_title", "year")
)

# Define a materialized view that has a filtered, aggregated, and sorted view of the data
@dp.materialized_view(
comment="A table summarizing counts of songs released by the artists who released the most songs each year."
)
def top_artists_by_year():
return (
spark.read.table("songs_prepared")
.filter(expr("year > 0"))
.groupBy("artist_name", "year")
.count().withColumnRenamed("count", "total_number_of_songs")
.sort(desc("total_number_of_songs"), desc("year"))
)

Cette source comprend le code pour trois queries. Vous pouvez également placer ces queries dans des fichiers séparés, afin d'organiser les fichiers et le code comme vous le souhaitez.

  1. Cliquez sur Icône de lecture. Exécuter le fichier ou Exécuter le pipeline pour start une mise à jour du pipeline connecté. Avec un seul fichier source dans votre pipeline, ils sont fonctionnellement équivalents.

Une fois la mise à jour terminée, l'éditeur est mis à jour avec les informations concernant votre pipeline.

  • Le graphe de pipeline, également appelé graphe acyclique dirigé (DAG), dans la barre latérale à droite de votre code, affiche trois tables, songs_raw, songs_prepared et top_artists_by_year.
  • Un résumé de la mise à jour est affiché en haut du navigateur d'actifs de pipeline.
  • Les détails des tables générées sont affichés dans le volet inférieur, et vous pouvez parcourir les données des tables en en sélectionnant une.

Cela inclut les données brutes et nettoyées, ainsi qu'une simple analyse pour trouver les meilleurs artistes par année. À l'étape suivante, vous créez des queries ad hoc pour une analyse plus approfondie dans un fichier séparé de votre pipeline.

Étape 3 : Explorez les jeux de données créés par votre pipeline

Dans cette étape, vous exécutez des requêtes ad hoc sur les données traitées dans le pipeline ETL pour analyser les données de chansons dans l'éditeur Databricks SQL. Ces queries utilisent les enregistrements préparés créés à l'étape précédente.

Tout d'abord, exécutez une query qui trouve les artistes ayant sorti le plus de chansons chaque année depuis 1990.

  1. Depuis la barre latérale du navigateur d'actifs de pipeline, cliquez sur Icône Plus. Ajouter , puis Exploration .

  2. Saisissez un **Nom** et sélectionnez **SQL** pour le fichier d'exploration. Un Notebook SQL est créé dans un nouveau dossier explorations. Les fichiers du dossier explorations ne sont pas exécutés par default dans le cadre d'une mise à jour de pipeline. Le Notebook SQL contient des cellules que vous pouvez exécuter ensemble ou séparément.

  3. Pour créer une table d'artistes qui sortent le plus de chansons chaque année après 1990, saisissez le code suivant dans le nouveau fichier SQL (s'il y a un exemple de code dans le fichier, remplacez-le). Puisque ce Notebook ne fait pas partie du pipeline, il n'utilise pas le catalogue et le schéma par default. Remplacez le <catalog>.<schema> par le catalogue et le schéma que vous avez utilisés par default pour le pipeline :

    SQL
    -- Which artists released the most songs each year in 1990 or later?
    SELECT artist_name, total_number_of_songs, year
    -- replace with the catalog/schema you are using:
    FROM <catalog>.<schema>.top_artists_by_year
    WHERE year >= 1990
    ORDER BY total_number_of_songs DESC, year DESC;
  4. Cliquez sur Icône de lecture. ou appuyez sur Shift + Enter pour exécuter cette query.

Exécutez maintenant une autre query qui trouve des chansons avec un rythme 4/4 et un tempo entraînant.

  1. Ajoutez le code suivant à la cellule suivante dans le même fichier. De nouveau, remplacez le <catalog>.<schema> par le catalogue et le schéma que vous avez utilisés par défaut pour le pipeline :

    SQL
    -- Find songs with a 4/4 beat and danceable tempo
    SELECT artist_name, song_title, tempo
    -- replace with the catalog/schema you are using:
    FROM <catalog>.<schema>.songs_prepared
    WHERE time_signature = 4 AND tempo between 100 and 140;
  2. Cliquez sur Icône de lecture. ou appuyez sur Shift + Enter pour exécuter cette query.

Étape 4 : Créez un job pour exécuter le pipeline

Ensuite, créez un workflow pour automatiser les étapes d'ingestion, de traitement et d'analyse des données à l'aide d'un Job Databricks qui s'exécute selon un calendrier.

  1. En haut de l'éditeur, choisissez le bouton Planifier .
  2. Si la boîte de dialogue Planifications apparaît, choisissez Ajouter une planification .
  3. Cela ouvre la boîte de dialogue Nouveau schedule , où vous pouvez créer un Job pour exécuter votre pipeline selon un schedule.
  4. Vous pouvez, si vous le souhaitez, donner un nom au Job.
  5. Par default, la planification est définie pour s'exécuter une fois par jour. Vous pouvez accepter ce réglage par défaut ou définir votre propre planification. Choisir Avancé vous donne la possibilité de définir une heure spécifique à laquelle le Job s'exécutera. La sélection de **Plus d'options** vous permet de créer des notifications lorsque le Job s'exécute.
  6. Sélectionnez Créer pour appliquer les modifications et créer le job.

Désormais, le job s'exécutera quotidiennement pour maintenir votre pipeline à jour. Vous pouvez choisir Planification à nouveau pour afficher la liste des planifications. Vous pouvez gérer les calendriers de votre pipeline à partir de cette boîte de dialogue, notamment l'ajout, la modification ou la suppression de calendriers.

Cliquer sur le nom du planning (ou du job) vous mène à la page du job dans la liste **Jobs et pipelines**. À partir de là, vous pouvez consulter les détails des exécutions de job, y compris l’historique des exécutions, ou exécuter le job immédiatement avec le bouton Run now .

Consultez Monitor Lakeflow Jobs pour plus d'informations sur les exécutions de Jobs.

En savoir plus