Aller au contenu principal

Utiliser Python avec des pipelines autonomes

Vous pouvez créer et refresh des vues matérialisées et des tables de streaming autonomes à partir d’un notebook en utilisant Python. Cela vous permet de gérer des pipeline autonomes aux côtés de vos autres flux de travail à Notebook basés sur Python.

Il existe deux façons de procéder :

  • Définissez la table avec les décorateurs pyspark.pipelines, @dp.materialized_view et @dp.table. Utilisez cette option lorsque la logique est plus facile à exprimer sous forme de code DataFrame. Voir Définir des tables avec les décorateurs pipelines.
  • Soumettez les mêmes instructions SQL qu'un warehouse Databricks SQL exécute en les transmettant à spark.sql(). Cela vous donne accès à l'ensemble des fonctionnalités SQL pour les vues matérialisées et les tables de streaming autonomes, y compris les instructions REFRESH et les plannings de refresh. Voir Soumettre des instructions SQL avec spark.sql().

La source Python pour les pipelines autonomes nécessite un notebook attaché à un compute serverless général . Vous ne pouvez pas utiliser Python pour créer ou refresh des pipelines autonomes à partir d'un warehouse Databricks SQL, car un warehouse exécute des instructions SQL, et non des notebooks Python. Pour utiliser un SQL warehouse à la place, consultez Utiliser des vues matérialisées autonomes et Utiliser des tables de streaming autonomes.

info

Bêta

La création et l'actualisation de vues matérialisées autonomes et de tables de streaming à partir d'un Notebook sur compute général Serverless sont en version Bêta et disponibles dans certaines régions. See Notebook.

Exigences​

Pour créer et refresh des pipelines autonomes avec Python, vous avez besoin d'un notebook attaché à un compute serverless général sur Databricks Runtime 18,1 ou version supérieure. Pour la liste complète des exigences, y compris la disponibilité régionale et les autorisations, consultez Notebooks.

Define tables with the pipelines decorators​

Vous pouvez définir une vue matérialisée ou une table de streaming autonome avec les mêmes décorateurs que ceux utilisés dans un LakeFlow Pipelines. Chaque fonction décorée définit une table. Lorsque vous exécutez la cellule, Databricks crée la table et exécute un pipeline serverless pour la remplir. La cellule renvoie le résultat une fois la mise à jour terminée.

attention

Les décorateurs de pipeline nécessitent la version de l'environnement serverless 5 ou supérieure.

Définir une vue matérialisée​

Use @dp.materialized_view on a function that returns a batch DataFrame. The following example creates the materialized view daily_booking_revenue from the bookings table in the Wanderbricks sample dataset:

Python
from pyspark import pipelines as dp
from pyspark.sql import functions as F

@dp.materialized_view(name="main.default.daily_booking_revenue")
def daily_booking_revenue():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)

To define a table from a streaming read, use @dp.table instead.

Define a streaming table​

Utilisez @dp.table sur une fonction qui renvoie un DataFrame de streaming. L'exemple suivant crée la table de streaming bookings_raw à partir d'une lecture de streaming de la même table bookings :

Python
from pyspark import pipelines as dp

@dp.table(name="main.default.bookings_raw")
def bookings_raw():
return spark.readStream.table("samples.wanderbricks.bookings").select("booking_id", "property_id", "total_amount")

Si la fonction renvoie un DataFrame batch, @dp.table crée à la place une vue matérialisée. La seule exception est replace_where, qui génère toujours une table de streaming. L'exemple suivant maintient à jour le chiffre d'affaires quotidien des enregistrements effectués à compter du 1er juillet 2025, sans recalculer les dates antérieures :

Python
from pyspark import pipelines as dp
from pyspark.sql import functions as F

@dp.table(
name="main.default.booking_revenue_rw",
replace_where=F.col("check_in") >= F.to_date(F.lit("2025-07-01")),
)
def booking_revenue_rw():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)

Chaque exécution supprime les lignes correspondant au prédicat et recalcule uniquement cette plage. Consultez Traitement par batch avec les flux REPLACE WHERE.

refresh a table​

To refresh a table you defined with a decorator, run the code that defines it again, for example by re-running the notebook cell, running the whole notebook, or running the notebook as a job. Each run creates the table if it doesn't exist and refresh it if it does.

Pour retraiter toutes les données disponibles dans la source, transmettez full_refresh=True sur l'un ou l'autre des décorateurs :

Python
@dp.table(name="main.default.bookings_raw", full_refresh=True)
def bookings_raw():
return spark.readStream.table("samples.wanderbricks.bookings").select("booking_id", "property_id", "total_amount")

You can't use a REFRESH statement on a table defined with a decorator, or schedule refreshes with SCHEDULE or TRIGGER ON UPDATE. To refresh on a schedule, define the table in SQL, or schedule the notebook as a job. See Lakeflow Jobs.

Configure the table​

The decorators accept the same common dataset parameters that they accept inside a pipeline, including comment, table_properties, partition_cols, cluster_by, schema, and spark_conf:

Python
@dp.materialized_view(
name="main.default.daily_booking_revenue",
comment="Daily booking revenue.",
table_properties={"quality": "gold"},
cluster_by=["check_in"],
)
def daily_booking_revenue():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)

Pour la liste des parameters, voir materialized_view et table.

private=True n'est pas pris en charge, car une table privée ne peut être lue que par d'autres jeux de données dans le même pipeline.

APIs non prises en charge​

A standalone table is a single dataset with a single flow, so the APIs that describe relationships between datasets aren't available. The following raise an error outside a pipeline:

  • @dp.temporary_view et dp.create_streaming_table
  • @dp.append_flow and other additional flows
  • dp.create_auto_cdc_flow et dp.create_auto_cdc_from_snapshot_flow
  • @dp.replace_flow et le paramètre replace_using, qui définissent des flux REPLACE USING. Consultez Remplacement de snapshot partiel avec les flux REPLACE USING.
  • dp.create_sink
  • Les attentes, telles que @dp.expect et @dp.expect_or_fail

Pour les utiliser, créez plutôt un pipeline Lakeflow. Consultez Develop pipeline code with Python.

Soumettre des instructions SQL avec spark.sql()​

Dans un Notebook Python, transmettez à spark.sql() les mêmes instructions que celles que vous exécuteriez depuis un warehouse Databricks SQL. La syntaxe des vues matérialisées autonomes et des tables de streaming est identique ; seule la façon de soumettre l'instruction diffère. Comme pour un warehouse, chaque instruction CREATE ou REFRESH exécute un pipeline serverless pour traiter l'opération.

La session spark est disponible par default dans les notebooks Databricks, aucune importation n'est donc nécessaire.

Créer une vue matérialisée​

L'exemple suivant crée la vue matérialisée mv1 à partir de la table de base base_table1:

Python
spark.sql("""
CREATE OR REPLACE MATERIALIZED VIEW mv1
AS SELECT
date,
sum(sales) AS sum_of_sales
FROM base_table1
GROUP BY date
""")

Pour obtenir tous les CREATE MATERIALIZED VIEW details, tels que les scheduled et Trigger refreshes, consultez Créer une vue matérialisée.

Créer une table de streaming​

L'exemple suivant crée la table de streaming sales à partir de la table raw_data :

Python
spark.sql("""
CREATE OR REFRESH STREAMING TABLE sales
AS SELECT product, price FROM STREAM raw_data
""")

Pour tous les détails sur CREATE STREAMING TABLE, y compris le chargement de fichiers avec Auto Loader et la planification, consultez Utiliser les tables de streaming autonomes.

refresh a vue matérialisée ou une table de streaming​

Utilisez une instruction REFRESH pour mettre à jour une table autonome avec les dernières données de sa source :

Python
spark.sql("REFRESH MATERIALIZED VIEW mv1")
spark.sql("REFRESH STREAMING TABLE sales")

Sur le compute général serverless, les actualisations sont synchrones. Les actualisations asynchrones (le mot-clé ASYNC) ne sont pas prises en charge. Voir le compute général serverless.

Paramétrer les instructions​

Pour transmettre des valeurs de votre code Python à une instruction au lieu de les coder en dur, utilisez des marqueurs de parameter nommés dans le SQL et fournissez leurs valeurs via l'argument args de spark.sql(). Utilisez directement un marqueur tel que :min_sales pour les valeurs littérales. N'encadrez le marqueur de IDENTIFIER() que lorsque le paramètre est un nom d'objet, tel qu'une table, une vue ou un schéma, car les identificateurs ne peuvent pas être substitués en tant que valeurs de chaîne simples.

L'exemple suivant paramètre à la fois le nom de la vue matérialisée et une valeur de filtre :

Python
mv_name = "main.sales.regional_sales"
min_sales = 1000

spark.sql("""
CREATE OR REPLACE MATERIALIZED VIEW IDENTIFIER(:mv)
AS SELECT
region,
sum(sales) AS sum_of_sales
FROM base_table1
WHERE sales > :min_sales
GROUP BY region
""", args={
"mv": mv_name,
"min_sales": min_sales,
})

Pour plus d'informations, consultez les marqueurs de parameter et la clause IDENTIFIER.

Exécuter d'autres instructions​

Vous pouvez exécuter toute instruction de vue matérialisée ou de table de streaming autonome à partir d'un Notebook Python en la transmettant à spark.sql(), y compris les instructions pour planifier des refreshs, modifier une table ou supprimer une table. Pour comprendre comment utiliser les vues matérialisées et les tables de streaming, y compris la syntaxe SQL, consultez Utiliser les vues matérialisées autonomes et Utiliser les tables de streaming autonomes.

Limitations​

Les vues matérialisées autonomes et les tables streaming créées sur le compute général Serverless ont des limitations supplémentaires, telles que l'absence de prise en charge des refresh asynchrones et l'absence d'attribution des coûts par table. Pour la liste complète, consultez le compute général Serverless.

Comme ces pipelines s’exécutent sur du compute général serverless plutôt que sur un SQL warehouse, ils n’héritent pas des tags personnalisés d’un warehouse englobant. La propagation des tags de warehouse vers system.billing.usage s’applique uniquement aux vues matérialisées et aux tables de streaming dont les instructions sont exécutées à partir d’un SQL warehouse. Voir Attribuer les coûts au SQL warehouse avec des tags personnalisés.

Les tables définies avec les décorateurs de pipeline présentent les limitations supplémentaires suivantes :

  • Vous ne pouvez pas les refresh avec une instruction REFRESH ou planifier des refreshes avec SCHEDULE ou TRIGGER ON UPDATE. Consultez refresh une table.
  • Les attentes, les flux supplémentaires, les flux de capture de données modifiées (CDC), les récepteurs et les vues temporaires ne sont pas pris en charge. Consultez APIs non prises en charge.
  • private=True n’est pas pris en charge.

Ressources supplémentaires​