Développer le code pipeline avec Python
Les LakeFlow Pipelines introduisent plusieurs nouvelles constructions de code Python pour définir des vues matérialisées et des tables de streaming dans les pipelines. Le support Python pour le développement de pipelines s'appuie sur les bases des API PySpark DataFrame et Structured Streaming.
Pour les utilisateurs qui ne sont pas familiers avec Python et les DataFrames, Databricks recommande d'utiliser l'interface SQL. Voir Développer le code des LakeFlow Pipelines avec SQL. Pour vous aider à choisir entre les deux interfaces, consultez Choisir entre SQL et Python.
Pour une référence complète de la syntaxe Python des Lakeflow Pipelines, veuillez consulter la référence du langage Python des Lakeflow Pipelines.
Bases de Python pour le développement de pipeline
Le code Python qui crée les datasets de pipeline doit retourner des DataFrames.
Toutes les APIs Python des Lakeflow pipelines sont implémentées dans le module pyspark.pipelines. Votre code de pipeline implémenté en Python doit explicitement importer le module pipelines en haut de la source Python. Dans nos exemples, nous utilisons la commande d'importation suivante, et utilisons dp dans les exemples pour faire référence à pipelines.
from pyspark import pipelines as dp
Apache Spark™ inclut des pipelines déclaratifs à partir de Spark 4.1, disponibles via le pyspark.pipelines module. Le Databricks Runtime étend ces fonctionnalités open source avec des APIs et des intégrations supplémentaires pour une utilisation en production gérée.
Le code écrit avec le module open source pipelines s'exécute sans modification sur Databricks. Les fonctionnalités suivantes ne font pas partie d'Apache Spark :
dp.create_auto_cdc_flowdp.create_auto_cdc_from_snapshot_flow@dp.expect(...)
default, la lecture et l’écriture du pipeline utilisent le catalogue et le schéma spécifiés lors de la configuration du pipeline. Voir Définir le catalogue et le schéma cibles.
Le code Python spécifique au pipeline diffère des autres types de code Python d'une manière essentielle : le code de pipeline Python n'appelle pas directement les fonctions qui effectuent l'ingestion et la transformation de données pour créer des datasets. Les LakeFlow Pipelines interprètent les fonctions de décorateur du module dp dans tous les fichiers de code source configurés dans un pipeline et créent un graphe de flux de données.
Pour éviter un comportement inattendu lors de l'exécution de votre pipeline, n'incluez pas de code qui pourrait avoir des effets secondaires dans vos fonctions qui définissent des datasets. Pour en savoir plus, consultez la référence Python.
Créer une vue matérialisée ou une table de streaming avec Python
Utilisez @dp.table pour créer une table de streaming à partir des résultats d'une lecture en streaming. Utilisez @dp.materialized_view pour créer une vue matérialisée à partir des résultats d'une lecture par batch.
By default, les noms des vues matérialisées et des tables de streaming sont déduits des noms de fonctions. L'exemple de code suivant montre la syntaxe de base pour créer une vue matérialisée et une table de streaming :
Les deux fonctions font référence à la même table dans le catalogue samples et utilisent la même fonction décoratrice. Ces exemples montrent que la seule différence dans la syntaxe de base des vues matérialisées et des tables de streaming est l'utilisation de spark.read par rapport à spark.readStream.
Toutes les sources de données ne prennent pas en charge les lectures en streaming. Certaines sources de données devraient toujours être traitées avec une sémantique de streaming.
from pyspark import pipelines as dp
@dp.materialized_view()
def basic_mv():
return spark.read.table("samples.nyctaxi.trips")
@dp.table()
def basic_st():
return spark.readStream.table("samples.nyctaxi.trips")
Facultativement, vous pouvez spécifier le nom de la table à l'aide de l'argument name dans le décorateur @dp.table. L'exemple suivant illustre ce modèle pour une vue matérialisée et une table de streaming :
from pyspark import pipelines as dp
@dp.materialized_view(name = "trips_mv")
def basic_mv():
return spark.read.table("samples.nyctaxi.trips")
@dp.table(name = "trips_st")
def basic_st():
return spark.readStream.table("samples.nyctaxi.trips")
Charger des données à partir du stockage d'objets
Les pipelines prennent en charge le chargement de données à partir de tous les formats pris en charge par Databricks. Consultez Options de format de données.
Ces exemples utilisent les données disponibles sous /databricks-datasets et montées automatiquement sur votre workspace. Databricks recommande d'utiliser des chemins de volume ou des URI cloud pour référencer les données stockées dans le stockage d'objets cloud. Consultez Que sont les volumes Unity Catalog ?.
Databricks recommande d'utiliser Auto Loader et les tables de streaming lors de la configuration des charges de travail d'ingestion incrémentielle pour les données stockées dans le stockage d'objets cloud. Consultez Qu'est-ce qu'Auto Loader ?.
L'exemple suivant crée une table de streaming à partir de fichiers JSON à l'aide d'Auto Loader :
from pyspark import pipelines as dp
@dp.table()
def ingestion_st():
return (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("/databricks-datasets/retail-org/sales_orders")
)
L'exemple suivant utilise la sémantique de traitement par batch pour lire un répertoire JSON et créer une vue matérialisée :
from pyspark import pipelines as dp
@dp.materialized_view()
def batch_mv():
return spark.read.format("json").load("/databricks-datasets/retail-org/sales_orders")
Valider les données avec des attentes
Vous pouvez utiliser des attentes pour définir et appliquer des contraintes de qualité des données. Consultez Gérer la qualité des données avec les attentes de pipeline.
Le code suivant utilise @dp.expect_or_drop pour définir une attente nommée valid_data qui supprime les enregistrements nuls pendant l'ingestion des données :
from pyspark import pipelines as dp
@dp.table()
@dp.expect_or_drop("valid_date", "order_datetime IS NOT NULL AND length(order_datetime) > 0")
def orders_valid():
return (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("/databricks-datasets/retail-org/sales_orders")
)
Query les vues matérialisées et les tables de streaming définies dans votre pipeline
L'exemple suivant définit quatre datasets :
- Une table de streaming nommée
ordersqui charge les données JSON. - Une vue matérialisée nommée
customersqui charge les données CSV. - Une vue matérialisée nommée
customer_ordersqui joint les enregistrements des datasetsordersetcustomers, convertit le timestamp de la commande en date, et sélectionne les champscustomer_id,order_number,stateetorder_date. - Une vue matérialisée nommée
daily_orders_by_statequi agrège le décompte quotidien des commandes pour chaque État.
Lorsque vous interrogez des vues ou des tables dans votre pipeline, vous pouvez spécifier le catalogue et le schéma directement, ou vous pouvez utiliser les valeurs default configurées dans votre pipeline. Dans cet exemple, les tables orders, customers et customer_orders sont écrites et lues à partir du catalogue et du schéma default configurés pour votre pipeline.
Le mode de publication hérité utilise le schéma LIVE pour query d'autres vues matérialisées et tables de streaming définies dans votre pipeline. Dans les nouveaux pipelines, la syntaxe de schéma LIVE est ignorée silencieusement. Consultez le schéma EN DIRECT (hérité).
from pyspark import pipelines as dp
from pyspark.sql.functions import col
@dp.table()
@dp.expect_or_drop("valid_date", "order_datetime IS NOT NULL AND length(order_datetime) > 0")
def orders():
return (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("/databricks-datasets/retail-org/sales_orders")
)
@dp.materialized_view()
def customers():
return spark.read.format("csv").option("header", True).load("/databricks-datasets/retail-org/customers")
@dp.materialized_view()
def customer_orders():
return (spark.read.table("orders")
.join(spark.read.table("customers"), "customer_id")
.select("customer_id",
"order_number",
"state",
col("order_datetime").cast("int").cast("timestamp").cast("date").alias("order_date"),
)
)
@dp.materialized_view()
def daily_orders_by_state():
return (spark.read.table("customer_orders")
.groupBy("state", "order_date")
.count().withColumnRenamed("count", "order_count")
)
Créer des tables dans une boucle for
Vous pouvez utiliser des boucles Python for pour créer plusieurs tables par programmation. Cela peut être utile lorsque vous avez de nombreuses source de données ou des dataset cibles qui ne varient que selon quelques parameters, ce qui réduit le volume total de code à maintenir et la redondance du code.
La boucle for évalue la logique en série, mais une fois la planification des datasets terminée, le pipeline exécute la logique en parallèle.
Lorsque vous utilisez ce modèle pour définir des datasets, assurez-vous que la liste des valeurs passées à la boucle for est toujours additive. Si un dataset précédemment défini dans un pipeline est omis d'une future exécution de pipeline, ce dataset est automatiquement supprimé du schéma cible.
L'exemple suivant crée cinq tables qui filtrent les commandes des clients par région. Ici, le nom de la région est utilisé pour définir le nom des vues matérialisées cibles et pour filtrer les données sources. Les vues temporaires sont utilisées pour définir des jointures à partir des tables sources utilisées dans la construction des vues matérialisées finales.
from pyspark import pipelines as dp
from pyspark.sql.functions import collect_list, col
@dp.temporary_view()
def customer_orders():
orders = spark.read.table("samples.tpch.orders")
customer = spark.read.table("samples.tpch.customer")
return (orders.join(customer, orders.o_custkey == customer.c_custkey)
.select(
col("c_custkey").alias("custkey"),
col("c_name").alias("name"),
col("c_nationkey").alias("nationkey"),
col("c_phone").alias("phone"),
col("o_orderkey").alias("orderkey"),
col("o_orderstatus").alias("orderstatus"),
col("o_totalprice").alias("totalprice"),
col("o_orderdate").alias("orderdate"))
)
@dp.temporary_view()
def nation_region():
nation = spark.read.table("samples.tpch.nation")
region = spark.read.table("samples.tpch.region")
return (nation.join(region, nation.n_regionkey == region.r_regionkey)
.select(
col("n_name").alias("nation"),
col("r_name").alias("region"),
col("n_nationkey").alias("nationkey")
)
)
# Extract region names from region table
region_list = spark.read.table("samples.tpch.region").select(collect_list("r_name")).collect()[0][0]
# Iterate through region names to create new region-specific materialized views
for region in region_list:
@dp.materialized_view(name=f"{region.lower().replace(' ', '_')}_customer_orders")
def regional_customer_orders(region_filter=region):
customer_orders = spark.read.table("customer_orders")
nation_region = spark.read.table("nation_region")
return (customer_orders.join(nation_region, customer_orders.nationkey == nation_region.nationkey)
.select(
col("custkey"),
col("name"),
col("phone"),
col("nation"),
col("region"),
col("orderkey"),
col("orderstatus"),
col("totalprice"),
col("orderdate")
).filter(f"region = '{region_filter}'")
)
Voici un exemple du Graphe de flux de données pour ce pipeline :

Dépannage : for boucle crée de nombreuses tables avec les mêmes valeurs
Le modèle d'exécution paresseuse que les pipelines utilisent pour évaluer le code Python exige que votre logique référence directement les valeurs individuelles lorsque la fonction décorée par @dp.materialized_view() est appelée.
L'exemple suivant démontre deux approches correctes pour définir des tables avec une boucle for. Dans les deux exemples, chaque nom de table de la liste tables est explicitement référencé au sein de la fonction décorée par @dp.materialized_view().
from pyspark import pipelines as dp
# Create a parent function to set local variables
def create_table(table_name):
@dp.materialized_view(name=table_name)
def t():
return spark.read.table(table_name)
tables = ["t1", "t2", "t3"]
for t_name in tables:
create_table(t_name)
# Call `@dp.materialized_view()` within a for loop and pass values as variables
tables = ["t1", "t2", "t3"]
for t_name in tables:
@dp.materialized_view(name=t_name)
def create_table(table_name=t_name):
return spark.read.table(table_name)
L’exemple suivant ne référence pas correctement les valeurs. Cet exemple crée des tables avec des noms distincts, mais toutes les tables chargent des données à partir de la dernière valeur dans la boucle for :
from pyspark import pipelines as dp
# Don't do this!
tables = ["t1", "t2", "t3"]
for t_name in tables:
@dp.materialized(name=t_name)
def create_table():
return spark.read.table(t_name)
Supprimer définitivement les enregistrements d'une vue matérialisée ou d'une table de streaming
Pour supprimer définitivement des enregistrements d'une vue matérialisée ou d'une table de streaming avec les vecteurs de suppression activés, par exemple pour la conformité GDPR, des opérations supplémentaires doivent être effectuées sur les tables Delta sous-jacentes de l'objet. Pour assurer la suppression des enregistrements d'une vue matérialisée, consultez Supprimer définitivement des enregistrements d'une vue matérialisée avec les vecteurs de suppression activés. Pour assurer la suppression des enregistrements d'une table de streaming, consultez Supprimer définitivement des enregistrements d'une table de streaming.