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 à l'aide de Python. Créez votre pipeline dans un Notebook Python et exécutez-les avec spark.sql(). Cela vous permet de gérer des pipelines autonomes avec vos autres workflows de Notebook basés sur Python.

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.

Comment cela fonctionne

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.

Ressources supplémentaires