Tutoriel : Créez votre premier pipeline à l'aide de l'éditeur Lakeflow Pipelines
Vous créez un nouveau LakeFlow pipeline pour l’orchestration des données avec Auto Loader, puis étendez le pipeline d’exemple en nettoyant les données et en créant une query pour trouver les 100 principaux utilisateurs.
Dans ce tutoriel, vous apprendrez à utiliser l'Éditeur de Lakeflow Pipelines pour :
- Créez un nouveau pipeline avec la structure de dossiers par default et start avec un ensemble de fichiers d'exemple.
- Définir les contraintes de qualité des données à l'aide d'attentes.
- Utilisez les fonctionnalités de l'éditeur pour étendre le pipeline avec une nouvelle transformation afin d'effectuer des analyses sur vos données.
Exigences
Avant de start ce tutoriel, vous devez :
- Connectez-vous à un Workspace Databricks.
- Faites en sorte que Unity Catalog soit activé pour votre Workspace.
- Disposer de l’autorisation de créer une ressource de compute ou d’accéder à une ressource de compute.
- Avoir les autorisations pour créer un nouveau schéma dans un catalogue. Les autorisations requises sont
ALL PRIVILEGESouUSE CATALOGetCREATE SCHEMA. - Pour l'ensemble complet des privilèges requis pour créer, exécuter, refresh et consulter les pipelines et leur sortie, consultez Gérer les identités, les autorisations et les privilèges pour les pipelines.
Étape 1 : Créez une pipeline
Dans cette étape, vous créez un pipeline en utilisant la structure de dossier default et des exemples de code. Les exemples de code font référence à la table users dans la source de données d'échantillon wanderbricks.
-
Dans votre Workspace Databricks, cliquez
sur **Nouveau**, puis sur
**pipeline ETL**. Cela ouvre l'éditeur de pipeline avec un nom de pipeline default tel que
New Pipeline <date> <time>. -
(Facultatif) Sélectionnez le nom et entrez un nom descriptif pour le pipeline.
-
(Facultatif) À droite du nom, cliquez sur le catalogue et le schéma pour définir différents default.
-
(Facultatif) Dans le fichier source
my_transformationcréé pour vous, sélectionnez Python ou SQL dans la liste déroulante des langues pour définir la langue du fichier. -
Cliquez
sur **Utiliser l'exemple de code**.
Un exemple de code dans la langue sélectionnée apparaît dans le fichier source
my_transformationdu dossiertransformations. Les datasets de sortie n'ont pas encore été créés, et le Graphe du pipeline sur le côté droit de l'écran est vide. -
Pour exécuter le code du pipeline (le code dans le dossier
transformations), cliquez sur Exécuter le pipeline dans la partie supérieure droite de l'écran.Une fois l'exécution terminée, la partie inférieure du workspace affiche les deux nouvelles tables qui ont été créées,
sample_users_<date_time>etsample_aggregation_<date_time>. Le Graphe de pipeline sur le côté droit du Workspace affiche maintenant les deux tables, y compris quesample_usersest la source desample_aggregation. Notez le nom complet de la tablesample_users_<date_time>; vous y faites référence à l'étape suivante.
Étape 2 : Appliquer les contrôles de qualité des données
Dans cette étape, vous ajoutez un contrôle de qualité des données à la table sample_users. Vous utilisez les attentes de pipeline pour contraindre les données. Dans ce cas, vous supprimez tous les enregistrements d'utilisateur qui n'ont pas d'adresse e-mail valide, et affichez la table nettoyée sous la forme users_cleaned.
-
Dans le navigateur d’assets de pipeline à gauche, cliquez sur
, puis sélectionnez Transformation .
-
Dans la boîte de dialogue Créer un nouveau fichier de transformation , effectuez les sélections suivantes :
- Choisissez Python ou SQL pour la langue . Cela ne doit pas correspondre à votre sélection précédente.
- Donnez un nom au fichier. Dans ce cas, choisissez
users_cleaned. - Pour **Chemin de destination**, laissez le default.
- Pour **Type de dataset**, laissez **Aucun sélectionné** ou choisissez **Vue matérialisée**. Si vous sélectionnez Vue matérialisée , il génère un exemple de code pour vous.
-
Cliquez sur Créer pour créer le fichier de code de transformation.
-
Dans votre nouveau fichier de code, modifiez le code pour qu'il corresponde à ce qui suit (utilisez SQL ou Python, selon votre sélection sur l'écran précédent). Remplacez
sample_users_<date_time>par le nom complet de votre tablesample_usersde la section précédente.
- SQL
- Python
-- Drop all rows that do not have an email address
CREATE MATERIALIZED VIEW users_cleaned
(
CONSTRAINT non_null_email EXPECT (email IS NOT NULL) ON VIOLATION DROP ROW
) AS
SELECT *
FROM sample_users_<date_time>;
from pyspark import pipelines as dp
# Drop all rows that do not have an email address
@dp.materialized_view
@dp.expect_or_drop("no null emails", "email IS NOT NULL")
def users_cleaned():
return (
spark.read.table("sample_users_<date_time>")
)
- Cliquez sur Exécuter le pipeline pour mettre à jour le pipeline. Il devrait maintenant y avoir trois tables.
Étape 3 : Analyser les utilisateurs principaux
Ensuite, obtenez les 100 premiers utilisateurs par le nombre de réservations qu’ils ont créées. Joignez la table wanderbricks.bookings à la vue matérialisée users_cleaned.
-
Dans le navigateur
d'actifs de pipeline à gauche, cliquez sur l'icône, et sélectionnez **Transformation**.
-
Dans la boîte de dialogue Créer un nouveau fichier de transformation , effectuez les sélections suivantes :
- Choisissez Python ou SQL pour la langue . Cela ne doit pas correspondre à vos sélections précédentes.
- Donnez un nom au fichier. Dans ce cas, choisissez
users_and_bookings. - Pour **Chemin de destination**, laissez le default.
- Pour Dataset type , laissez-le sur Aucune sélection .
-
Cliquez sur Créer pour créer le fichier de code de transformation.
-
Dans votre nouveau fichier de code, modifiez le code pour qu'il corresponde à ce qui suit (utilisez SQL ou Python, en fonction de votre sélection sur l'écran précédent).
- SQL
- Python
-- Get the top 100 users by number of bookings
CREATE OR REFRESH MATERIALIZED VIEW users_and_bookings AS
SELECT u.name AS name, COUNT(b.booking_id) AS booking_count
FROM users_cleaned u
JOIN samples.wanderbricks.bookings b ON u.user_id = b.user_id
GROUP BY u.name
ORDER BY booking_count DESC
LIMIT 100;
from pyspark import pipelines as dp
from pyspark.sql.functions import col, count, desc
# Get the top 100 users by number of bookings
@dp.materialized_view
def users_and_bookings():
return (
spark.read.table("users_cleaned")
.join(spark.read.table("samples.wanderbricks.bookings"), "user_id")
.groupBy(col("name"))
.agg(count("booking_id").alias("booking_count"))
.orderBy(desc("booking_count"))
.limit(100)
)
-
Cliquez sur **Exécuter le pipeline** pour mettre à jour les datasets. Une fois l'exécution terminée, vous pouvez voir dans le **Graphe de pipeline** qu'il y a quatre tables, y compris la nouvelle
users_and_bookingstable.
Ressources supplémentaires
Maintenant que vous avez appris à utiliser certaines des fonctionnalités de l'éditeur de Lakeflow pipelines et créé un pipeline, voici d'autres fonctionnalités à découvrir :
-
Outils pour travailler et debugging les Transformations lors de la création de pipelines :
- Exécution sélective
- Aperçus des données
- Graphe de pipeline interactif (graphe des datasets dans votre pipeline)
-
Intégré Declarative Automation Bundles intégration pour une collaboration efficace, le contrôle de version et l'intégration CI/CD directement depuis l'éditeur :