Tutoriel : Créer un pipeline ETL à l'aide de la capture de données modifiées
Vous créez et déployez un pipeline ETL (extraction, transformation et chargement) avec capture des modifications de données (CDC) en utilisant les LakeFlow Pipelines pour l'orchestration des données et 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 les données d'une table customers dans une base de données MySQL pour :
- Extrayez les modifications d'une base de données transactionnelle à l'aide de Debezium ou d'un autre outil et enregistrez-les dans le stockage d'objets cloud (S3, ADLS ou GCS). Dans ce tutoriel, vous ignorez la configuration d'un système CDC externe et générez plutôt des données factices pour simplifier le tutoriel.
- Utilisez Auto Loader pour charger de manière incrémentielle les messages depuis le stockage d'objets cloud, et stockez les messages bruts dans la table
customers_cdc. Auto Loader déduit le schéma et gère l'évolution des schémas. - Créez la table
customers_cdc_cleanpour vérifier la qualité des données à l'aide d'attentes. Par exemple, leidne devrait jamais êtrenullcar il est utilisé pour exécuter des opérations d'upsert. - Effectuez
AUTO CDC ... INTOsur les données CDC nettoyées pour mettre à jour les modifications dans la tablecustomersfinale. - Montrer comment un pipeline peut créer une table de dimension à évolution lente de type 2 (SCD2) pour suivre toutes les modifications.
L'objectif est d'ingérer les données brutes en quasi temps réel et de créer une table pour votre équipe d'analystes tout en garantissant la qualité des données.
Le tutoriel utilise l'architecture Lakehouse en médaillon, où il ingère les données brutes via la couche Bronze, nettoie et valide les données avec la couche Silver, et applique la modélisation dimensionnelle et l'agrégation à l'aide de la couche Gold. Voir Qu'est-ce que l'architecture lakehouse en médaillon ? pour plus d'informations.
Le flux implémenté se présente comme suit :

Pour en savoir plus sur les pipelines, Auto Loader et la CDC, consultez Spark Declarative Pipelines, Qu'est-ce qu'Auto Loader ? et Capture de données modifiées et instantanés
Exigences
Pour suivre ce didacticiel, vous devez satisfaire aux exigences suivantes :
- Unity Catalog est activé pour votre Workspace.
- Le compute Serverless est disponible dans votre Workspace (activé par default dans les workspaces avec Unity Catalog). Les Lakeflow pipelines Serverless ne sont pas disponibles dans toutes les régions de Workspace. Pour connaître les régions disponibles, consultez les Fonctionnalités avec disponibilité régionale limitée. Si le Serverless compute n'est pas disponible, les étapes devraient fonctionner avec le compute par default de votre Workspace.
- Autorisation de créer une ressource de compute ou d’accéder à une ressource de compute.
- Autorisations de créer un nouveau schéma dans un catalogue. Les autorisations requises sont
USE CATALOGetCREATE SCHEMA. - Autorisations pour créer un nouveau volume dans un schéma existant. Les autorisations requises sont
USE SCHEMAetCREATE VOLUME. - 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.
Capture de données modifiées dans un pipeline ETL
La capture de données modifiées (CDC) est le processus qui capture les modifications apportées aux enregistrements d'une base de données transactionnelle (par exemple, MySQL ou PostgreSQL) ou d'un data warehouse. La CDC capture des opérations telles que les suppressions de données, les ajouts et les mises à jour, généralement sous forme de stream pour rematérialiser les tables dans des systèmes externes. La CDC permet le chargement incrémental tout en éliminant le besoin de mises à jour par chargement en masse.
Pour simplifier ce tutoriel, omettez de configurer un système CDC externe. Supposez qu’il s’exécute et enregistre les données CDC sous forme de fichiers JSON dans le stockage d’objets cloud (S3, ADLS ou GCS). Ce tutoriel utilise la bibliothèque Faker pour générer les données utilisées dans le tutoriel.
Capture du CDC
Une variété d’outils CDC sont disponibles. L'une des principales solutions open source est Debezium, mais d'autres implémentations qui simplifient les sources de données existent, telles que Fivetran, Qlik Replicate, StreamSets, Talend, Oracle GoldenGate et AWS DMS.
Dans ce tutoriel, vous utilisez les données CDC provenant d'un système externe tel que Debezium ou DMS. Debezium capture chaque ligne modifiée. Il envoie généralement l'historique des modifications de données aux rubriques Kafka ou les enregistre sous forme de fichiers.
Vous devez ingérer les informations CDC de la table customers (format JSON), vérifier qu'elles sont correctes, puis matérialiser la table des clients dans le Lakehouse.
Entrée CDC de Debezium
Pour chaque modification, vous recevez un message JSON contenant tous les champs de la ligne mise à jour (id, firstname, lastname, email, address). Le message comprend également des métadonnées supplémentaires :
operation: Un code d’opération, généralement (DELETE,APPEND,UPDATE).operation_date: La date et le Timestamp de l'enregistrement pour chaque action d'Opération.
Des outils comme Debezium peuvent produire une sortie plus avancée, telle que la valeur de la ligne avant la modification, mais ce tutoriel les omet par souci de simplicité.
Étape 1 : Créez une pipeline
Créez un nouveau pipeline ETL pour interroger votre source de données CDC et générer des tables dans votre workspace.
-
Dans votre Workspace, cliquez sur
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>. -
Sélectionnez le nom et saisissez un nom descriptif, tel que
Pipelines with CDC tutorial. -
À droite du nom, cliquez sur le catalogue et le schéma pour choisir les valeurs par default pour lesquelles vous disposez de permissions d’écriture.
Ce catalogue et ce schéma sont utilisés par default, si vous ne spécifiez pas de catalogue ou de schéma dans votre code. Votre code peut écrire dans n'importe quel catalogue ou schéma en spécifiant le chemin complet. Ce tutoriel utilise les valeurs par default que vous spécifiez ici.
-
(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**.
L'éditeur de Lakeflow Pipelines s'ouvre avec des exemples de fichiers dans votre pipeline. Vous ajouterez vos propres fichiers de transformations dans les étapes ultérieures. Ensuite, configurez l’échantillon de données à importer dans le tutoriel.
Étape 2 : Créer l’échantillon de données à importer dans ce tutoriel
Cette étape n'est pas nécessaire si vous importez vos propres données à partir d'une source existante. Pour ce didacticiel, générez de fausses données à titre d’exemple pour le didacticiel. Créez un Notebook pour exécuter le script de génération de données Python. Ce code ne doit être exécuté qu'une seule fois pour générer les exemples de données. Créez-le donc dans le dossier explorations du pipeline, qui n'est pas exécuté dans le cadre d'une mise à jour du pipeline.
Ce code utilise Faker pour générer les données CDC d'exemple. Faker est disponible pour s'installer automatiquement, donc le tutoriel utilise %pip install faker. Vous pouvez également définir une dépendance sur Faker pour le Notebook. Voir Ajouter des dépendances au Notebook.
-
Depuis l'éditeur LakeFlow Pipelines, dans la barre latérale du navigateur d'asset à gauche de l'éditeur, cliquez sur
Ajouter , puis choisissez Exploration .
-
Donnez-lui un Nom , tel que
Setup data, sélectionnez Python . Vous pouvez laisser le dossier de destination default, qui est un nouveau dossierexplorations. -
Cliquez sur Créer. Cela crée un Notebook dans le nouveau dossier.
-
Saisissez le code suivant dans la première cellule. Vous devez modifier la définition de
<my_catalog>et<my_schema>pour qu'elle corresponde au catalogue et au schéma par default que vous avez sélectionnés lors de la procédure précédente :Python%pip install faker
# Update these to match the catalog and schema
# that you used for the pipeline in step 1.
catalog = "<my_catalog>"
schema = dbName = db = "<my_schema>"
spark.sql(f'USE CATALOG `{catalog}`')
spark.sql(f'USE SCHEMA `{schema}`')
spark.sql(f'CREATE VOLUME IF NOT EXISTS `{catalog}`.`{db}`.`raw_data`')
volume_folder = f"/Volumes/{catalog}/{db}/raw_data"
try:
dbutils.fs.ls(volume_folder+"/customers")
except:
print(f"folder doesn't exist, generating the data under {volume_folder}...")
from pyspark.sql import functions as F
from faker import Faker
from collections import OrderedDict
import uuid
fake = Faker()
import random
fake_firstname = F.udf(fake.first_name)
fake_lastname = F.udf(fake.last_name)
fake_email = F.udf(fake.ascii_company_email)
fake_date = F.udf(lambda:fake.date_time_this_month().strftime("%m-%d-%Y %H:%M:%S"))
fake_address = F.udf(fake.address)
operations = OrderedDict([("APPEND", 0.5),("DELETE", 0.1),("UPDATE", 0.3),(None, 0.01)])
fake_operation = F.udf(lambda:fake.random_elements(elements=operations, length=1)[0])
fake_id = F.udf(lambda: str(uuid.uuid4()) if random.uniform(0, 1) < 0.98 else None)
df = spark.range(0, 100000).repartition(100)
df = df.withColumn("id", fake_id())
df = df.withColumn("firstname", fake_firstname())
df = df.withColumn("lastname", fake_lastname())
df = df.withColumn("email", fake_email())
df = df.withColumn("address", fake_address())
df = df.withColumn("operation", fake_operation())
df_customers = df.withColumn("operation_date", fake_date())
df_customers.repartition(100).write.format("json").mode("overwrite").save(volume_folder+"/customers") -
Pour générer le dataset utilisé dans le tutoriel, tapez **Maj** + **Entrée** pour exécuter le code :
-
Facultatif. Pour prévisualiser les données utilisées dans ce tutoriel, saisissez le code suivant dans la cellule suivante et exécutez le code. Mettez à jour le catalogue et le schéma pour qu'ils correspondent au chemin du code précédent.
Python# Update these to match the catalog and schema
# that you used for the pipeline in step 1.
catalog = "<my_catalog>"
schema = "<my_schema>"
display(spark.read.json(f"/Volumes/{catalog}/{schema}/raw_data/customers"))
Ceci génère un grand dataset (avec de fausses données CDC) que vous pouvez utiliser dans le reste du tutoriel. À l'étape suivante, ingérez les données à l'aide d'Auto Loader.
Étape 3 : ingérer les données de manière incrémentielle avec Auto Loader
L'étape suivante consiste à ingérer les données brutes depuis le stockage cloud (fictif) dans une couche bronze.
Cela peut être difficile pour plusieurs raisons, car vous devez :
- Fonctionnez à grande échelle, en ingérant potentiellement des millions de petits fichiers.
- Déduire le schéma et le type JSON.
- Gérer les enregistrements incorrects avec un schéma JSON incorrect.
- Gérez l'évolution des schémas (par exemple, une nouvelle colonne dans la table des clients).
Auto Loader simplifie cette ingestion, y compris l'inférence de schéma et l'évolution des schémas, tout en s'adaptant à des millions de fichiers entrants. Auto Loader est disponible en Python à l'aide de cloudFiles et en SQL à l'aide de SELECT * FROM STREAM read_files(...) et peut être utilisé avec une variété de formats (JSON, CSV, Apache Avro, etc.) :
La définition de la table en tant que table de streaming garantit que vous ne consommez que les nouvelles données entrantes. Si vous ne la définissez pas comme une table de streaming, elle scanne et ingère toutes les données disponibles. Consultez les tables de streaming pour plus d’informations.
- Pour ingérer les données CDC entrantes à l'aide d'Auto Loader, copiez et collez le code suivant dans le fichier de code qui a été créé avec votre pipeline (appelé
my_transformation.pyoumy_transformation.sql). Vous pouvez utiliser Python ou SQL, en fonction de la langue que vous avez choisie lors de la création du pipeline. Assurez-vous de remplacer les<catalog>et<schema>par ceux que vous avez configurés par default pour le pipeline.
- Python
- SQL
from pyspark import pipelines as dp
from pyspark.sql.functions import *
# Replace with the catalog and schema name that
# you are using:
path = "/Volumes/<catalog>/<schema>/raw_data/customers"
# Create the target bronze table
dp.create_streaming_table("customers_cdc_bronze", comment="New customer data incrementally ingested from cloud object storage landing zone")
# Create an Append Flow to ingest the raw data into the bronze table
@dp.append_flow(
target = "customers_cdc_bronze",
name = "customers_bronze_ingest_flow"
)
def customers_bronze_ingest_flow():
return (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.inferColumnTypes", "true")
.load(f"{path}")
)
CREATE OR REFRESH STREAMING TABLE customers_cdc_bronze
COMMENT "New customer data incrementally ingested from cloud object storage landing zone";
CREATE FLOW customers_bronze_ingest_flow AS
INSERT INTO customers_cdc_bronze BY NAME
SELECT *
FROM STREAM read_files(
-- replace with the catalog/schema you are using:
"/Volumes/<catalog>/<schema>/raw_data/customers",
format => "json",
inferColumnTypes => "true"
)
- Cliquez sur
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 orienté acyclique (DAG), dans la barre latérale à droite de votre code, affiche une seule table,
customers_cdc_bronze. - Un résumé de la mise à jour est affiché en haut du navigateur d'actifs de pipeline.
- Les détails de la table générée sont affichés dans le volet inférieur, et vous pouvez parcourir les données de la table en la sélectionnant.
Il s'agit des données de la couche bronze brutes importées depuis le stockage cloud. À l'étape suivante, nettoyez les données pour créer une table de la couche argent.
Étape 4 : Nettoyage et attentes pour le suivi de la qualité des données
Après la définition de la couche bronze, créez la couche silver en ajoutant des attentes pour contrôler la qualité des données. Vérifiez les conditions suivantes :
- L'ID ne doit jamais être
null. - Le type d'Opération du CDC doit être valide.
- Le JSON doit être lu correctement par Auto Loader.
Les lignes qui ne remplissent pas ces conditions sont supprimées.
Consultez Gérer la qualité des données avec les attentes du pipeline pour plus d'informations.
-
Dans la barre latérale du navigateur d'asset de pipeline, cliquez sur
Ajouter , puis sur Transformation .
-
Entrez un Nom (par exemple,
customers_silver) et choisissez une langue (Python ou SQL) pour le fichier de code source. L'extension.pyou.sqlest ajoutée en fonction de votre choix de langue. Vous pouvez mélanger les langues au sein d'un pipeline, vous pouvez donc choisir l'une ou l'autre pour cette étape. -
Pour créer une couche Silver avec une table nettoyée et imposer des contraintes, copiez et collez le code suivant dans le nouveau fichier (choisissez Python ou SQL en fonction du langage du fichier).
- Python
- SQL
from pyspark import pipelines as dp
from pyspark.sql.functions import *
dp.create_streaming_table(
name = "customers_cdc_clean",
expect_all_or_drop = {"no_rescued_data": "_rescued_data IS NULL","valid_id": "id IS NOT NULL","valid_operation": "operation IN ('APPEND', 'DELETE', 'UPDATE')"}
)
@dp.append_flow(
target = "customers_cdc_clean",
name = "customers_cdc_clean_flow"
)
def customers_cdc_clean_flow():
return (
spark.readStream.table("customers_cdc_bronze")
.select("address", "email", "id", "firstname", "lastname", "operation", "operation_date", "_rescued_data")
)
CREATE OR REFRESH STREAMING TABLE customers_cdc_clean (
CONSTRAINT no_rescued_data EXPECT (_rescued_data IS NULL) ON VIOLATION DROP ROW,
CONSTRAINT valid_id EXPECT (id IS NOT NULL) ON VIOLATION DROP ROW,
CONSTRAINT valid_operation EXPECT (operation IN ('APPEND', 'DELETE', 'UPDATE')) ON VIOLATION DROP ROW
)
COMMENT "New customer data incrementally ingested from cloud object storage landing zone";
CREATE FLOW customers_cdc_clean_flow AS
INSERT INTO customers_cdc_clean BY NAME
SELECT * FROM STREAM customers_cdc_bronze;
-
Cliquez sur
Exécuter le fichier ou Exécuter le pipeline pour start une mise à jour pour le pipeline connecté.
Étant donné qu’il y a maintenant deux fichiers sources, ceux-ci ne font pas la même chose, mais dans ce cas, la sortie est la même.
- **Exécuter le pipeline** exécute l’intégralité de votre pipeline, y compris le code de l’étape 3. Si vos données d'entrée étaient mises à jour, cela extrairait toutes les modifications de cette source vers votre couche Bronze. Cela n'exécute pas le code de l'étape de configuration des données, car il se trouve dans le dossier d'explorations et ne fait pas partie de la source de votre pipeline.
- Exécuter le fichier exécute uniquement le fichier source actuel. Dans ce cas, sans la mise à jour de vos données d'entrée, cela génère les données silver à partir de la table bronze mise en cache. Il serait utile d'exécuter uniquement ce fichier pour une itération plus rapide lors de la création ou de la modification de votre code de pipeline.
Une fois la mise à jour terminée, vous pouvez voir que le graphe du pipeline affiche désormais deux tables (avec la couche argent dépendant de la couche bronze), et le panneau inférieur affiche les détails des deux tables. Le haut du navigateur d'assets de pipeline affiche désormais les heures de plusieurs exécutions, mais seulement les détails de l'exécution la plus récente.
Ensuite, créez votre version finale de la couche Gold de la table customers.
Étape 5 : Matérialiser la table des clients avec un flux AUTO CDC
Jusqu'à présent, les tables ont simplement transmis les données CDC à chaque étape. Maintenant, créez la table customers pour qu'elle contienne la vue la plus à jour et qu'elle soit une réplique de la table d'origine, et non la liste des opérations CDC qui l'ont créée.
Cela est non trivial à mettre en œuvre manuellement. Vous devez prendre en compte des éléments tels que la déduplication des données pour conserver la ligne la plus récente.
Cependant, les pipelines résolvent ces défis avec l'AUTO CDC Opérations.
-
Dans la barre latérale du navigateur d'asset de pipeline, cliquez
sur **Ajouter** et **Transformation**.
-
Saisissez un **Nom** et choisissez une langue (Python ou SQL) pour le nouveau fichier de code source. Vous pouvez à nouveau choisir l'une ou l'autre langue pour cette étape, mais utilisez le code correct ci-dessous.
-
Pour traiter les données CDC à l’aide de
AUTO CDC, copiez et collez le code suivant dans le nouveau fichier.
- Python
- SQL
from pyspark import pipelines as dp
from pyspark.sql.functions import *
dp.create_streaming_table(name="customers", comment="Clean, materialized customers")
dp.create_auto_cdc_flow(
target="customers", # The customer table being materialized
source="customers_cdc_clean", # the incoming CDC
keys=["id"], # what we'll be using to match the rows to upsert
sequence_by=col("operation_date"), # de-duplicate by operation date, getting the most recent value
ignore_null_updates=False,
apply_as_deletes=expr("operation = 'DELETE'"), # DELETE condition
except_column_list=["operation", "operation_date", "_rescued_data"],
)
CREATE OR REFRESH STREAMING TABLE customers;
CREATE FLOW customers_cdc_flow
AS AUTO CDC INTO customers
FROM stream(customers_cdc_clean)
KEYS (id)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY operation_date
COLUMNS * EXCEPT (operation, operation_date, _rescued_data)
STORED AS SCD TYPE 1;
- Cliquez sur
Exécuter le fichier pour start une mise à jour du pipeline connecté.
Lorsque la mise à jour est terminée, vous pouvez constater que votre Graphe de pipeline affiche 3 tables, progressant du bronze à l'argent puis à Gold.
Étape 6 : Suivi de l'historique des mises à jour avec la dimension à évolution lente de type 2 (SCD2)
Il est souvent nécessaire de créer une table qui suit toutes les modifications résultant de APPEND, UPDATE et DELETE:
- Historique : vous souhaitez conserver un historique de toutes les modifications apportées à votre table.
- Traçabilité : Vous voulez voir quelle opération a eu lieu.
SCD2 avec des LakeFlow Pipelines
Delta prend en charge le flux de données de modification (CDF), et table_change peut interroger les modifications de table en SQL et Python. Cependant, le cas d'utilisation principal de CDF est de capturer les changements dans un pipeline, et non de créer une vue complète des modifications de table depuis le début.
Il devient particulièrement complexe d'implémenter des choses si vous avez des événements désordonnés. Si vous devez séquencer vos modifications par un Timestamp et recevez une modification qui s'est produite dans le passé, vous devez ajouter une nouvelle entrée dans votre table SCD et mettre à jour les entrées précédentes.
Les LakeFlow pipelines suppriment cette complexité. Vous pouvez créer une table distincte qui contient toutes les modifications depuis le début des temps, que le pipeline maintient automatiquement. Cette table prend en charge les optimisations de Layout des données telles que le clustering liquide et gère automatiquement les enregistrements désordonnés en fonction de _sequence_by. Voir Utiliser le clustering liquide pour les tables.
Pour créer une table SCD2, utilisez l'option STORED AS SCD TYPE 2 en SQL ou stored_as_scd_type="2" en Python.
Vous pouvez également limiter les colonnes suivies par la fonctionnalité à l'aide de l'option : TRACK HISTORY ON {columnList | EXCEPT(exceptColumnList)}
-
Dans la barre latérale du navigateur d'asset de pipeline, cliquez
sur **Ajouter** et **Transformation**.
-
Saisissez un Nom et choisissez une langue (Python ou SQL) pour le nouveau fichier de code source.
-
Copiez et collez le code suivant dans le nouveau fichier.
- Python
- SQL
from pyspark import pipelines as dp
from pyspark.sql.functions import *
# create the table
dp.create_streaming_table(
name="customers_history", comment="Slowly Changing Dimension Type 2 for customers"
)
# store all changes as SCD2
dp.create_auto_cdc_flow(
target="customers_history",
source="customers_cdc_clean",
keys=["id"],
sequence_by=col("operation_date"),
ignore_null_updates=False,
apply_as_deletes=expr("operation = 'DELETE'"),
except_column_list=["operation", "operation_date", "_rescued_data"],
stored_as_scd_type="2",
) # Enable SCD2 and store individual updates
CREATE OR REFRESH STREAMING TABLE customers_history;
CREATE FLOW customers_history_cdc
AS AUTO CDC INTO
customers_history
FROM stream(customers_cdc_clean)
KEYS (id)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY operation_date
COLUMNS * EXCEPT (operation, operation_date, _rescued_data)
STORED AS SCD TYPE 2;
- Cliquez sur
Exécuter le fichier pour start une mise à jour du pipeline connecté.
Une fois la mise à jour terminée, le graphe du pipeline inclut la nouvelle table customers_history, également dépendante de la table de la couche silver, et le panneau inférieur affiche les détails pour les 4 tables.
Étape 7 : Créer une vue matérialisée qui suit les personnes ayant le plus modifié leurs informations.
La table customers_history contient toutes les modifications historiques qu'un utilisateur a apportées à ses informations. Créez une vue matérialisée simple dans la couche Gold qui assure le suivi des personnes ayant le plus modifié leurs informations. Ceci pourrait être utilisé pour l'analyse de la détection de fraude ou des recommandations d'utilisateurs dans un scénario réel. De plus, l'application de modifications avec SCD2 a déjà supprimé les doublons, vous pouvez donc directement compter les lignes par ID utilisateur.
-
Dans la barre latérale du navigateur d'asset de pipeline, cliquez
sur **Ajouter** et **Transformation**.
-
Saisissez un Nom et choisissez une langue (Python ou SQL) pour le nouveau fichier de code source.
-
Copiez et collez le code suivant dans le nouveau fichier source.
- Python
- SQL
from pyspark import pipelines as dp
from pyspark.sql.functions import *
@dp.table(
name = "customers_history_agg",
comment = "Aggregated customer history"
)
def customers_history_agg():
return (
spark.read.table("customers_history")
.groupBy("id")
.agg(
count_distinct("address").alias("address_count"),
count_distinct("email").alias("email_count"),
count_distinct("firstname").alias("firstname_count"),
count_distinct("lastname").alias("lastname_count")
)
)
CREATE OR REPLACE MATERIALIZED VIEW customers_history_agg AS
SELECT
id,
count(distinct address) as address_count,
count(distinct email) AS email_count,
count(distinct firstname) AS firstname_count,
count(distinct lastname) AS lastname_count
FROM customers_history
GROUP BY id
- Cliquez sur
Exécuter le fichier pour start une mise à jour du pipeline connecté.
Une fois la mise à jour terminée, il y a une nouvelle table dans le graphe du pipeline qui dépend de la table customers_history, et vous pouvez la consulter dans le panneau inférieur. Votre pipeline est maintenant terminé. Vous pouvez le tester en effectuant une exécution complète du **pipeline**. Les seules étapes restantes consistent à programmer le pipeline pour qu'il se mette à jour régulièrement.
Étape 8 : Créer un Job pour exécuter le pipeline ETL
Ensuite, créez un workflow pour automatiser les étapes d'ingestion, de traitement et d'analyse des données dans votre pipeline à l'aide d'un Databricks job.
- En haut de l'éditeur, choisissez le bouton Planifier .
- Si la boîte de dialogue Planifications apparaît, choisissez Ajouter une planification .
- Cela ouvre la boîte de dialogue Nouveau schedule , où vous pouvez créer un Job pour exécuter votre pipeline selon un schedule.
- Vous pouvez, si vous le souhaitez, donner un nom au Job.
- Par default, la planification est définie pour s'exécuter une fois par jour. Vous pouvez accepter cette default, ou définir votre propre programmation. 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.
- 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.