Fonctions pour définir des datasets
Le module pyspark.pipelines (ici nommé dp) implémente une grande partie de ses fonctionnalités principales à l'aide de décorateurs. Ces décorateurs acceptent une fonction qui définit une query en streaming ou en batch et renvoie un DataFrame Apache Spark. La syntaxe suivante présente un exemple simple pour la définition d'un dataset de pipeline :
from pyspark import pipelines as dp
@dp.table()
def function_name(): # This is the function decorated
return (<query>) # This is the query logic that defines the dataset
Cette page fournit un aperçu des fonctions et des query qui définissent des dataset dans les pipeline. Pour obtenir une liste complète des décorateurs disponibles, consultez la référence du développeur de pipeline.
Les fonctions que vous utilisez pour définir des datasets ne devraient pas inclure de logique Python arbitraire sans rapport avec le dataset, y compris les appels aux APIs tierces. Les pipelines exécutent ces fonctions plusieurs fois pendant la planification, la validation et les mises à jour. L'inclusion d'une logique arbitraire peut entraîner des résultats inattendus.
Lire les données pour commencer une définition de dataset
Les fonctions utilisées pour définir les datasets de pipeline commencent généralement par une opération spark.read ou spark.readStream. Ces opérations de lecture renvoient un objet DataFrame statique ou streaming que vous utilisez pour définir des transformations supplémentaires avant de renvoyer le DataFrame. D'autres exemples d'opérations Spark qui renvoient un DataFrame incluent spark.table ou spark.range.
Les fonctions ne doivent jamais référencer des DataFrames définis en dehors de la fonction. Tenter de référencer des DataFrames définis dans une portée différente pourrait entraîner un comportement inattendu. Pour un exemple de modèle de métaprogrammation pour la création de plusieurs tables, consultez Créer des tables dans une boucle for.
Les exemples suivants montrent la syntaxe de base pour la lecture des données à l'aide de la logique batch ou streaming :
from pyspark import pipelines as dp
# Batch read on a table
@dp.materialized_view()
def function_name():
return spark.read.table("catalog_name.schema_name.table_name")
# Batch read on a path
@dp.materialized_view()
def function_name():
return spark.read.format("parquet").load("/Volumes/catalog_name/schema_name/volume_name/data_path")
# Streaming read on a table
@dp.table()
def function_name():
return spark.readStream.table("catalog_name.schema_name.table_name")
# Streaming read on a path
@dp.table()
def function_name():
return (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "parquet")
.load("/Volumes/catalog_name/schema_name/volume_name/data_path")
)
Si vous devez lire des données d'une API REST externe, implémentez cette connexion en utilisant une source de données Python personnalisée. Consultez les sources de données personnalisées PySpark.
Il est possible de créer des DataFrames Apache Spark arbitraires à partir de collections de données Python, y compris des DataFrames pandas, des dictionnaires et des listes. Ces modèles pourraient être utiles pendant le développement et les tests, mais la plupart des définitions de dataset de pipeline de production devraient commencer par le chargement de données à partir de fichiers, d'un système externe, ou d'une table ou vue existante.
Chaînage des transformations
Les pipelines prennent en charge presque toutes les transformations des DataFrame Apache Spark. Vous pouvez inclure n'importe quel nombre de Transformations dans votre fonction de définition de dataset, mais vous devez vous assurer que les méthodes que vous utilisez retournent toujours un objet DataFrame.
Si vous avez une transformation intermédiaire qui alimente plusieurs charges de travail en aval, mais que vous n'avez pas besoin de la matérialiser sous forme de table, utilisez @dp.temporary_view() pour ajouter une vue temporaire à votre pipeline. Vous pouvez ensuite référencer cette vue en utilisant spark.read.table("temp_view_name") dans plusieurs définitions de datasets en aval. La syntaxe suivante illustre ce modèle :
from pyspark import pipelines as dp
@dp.temporary_view()
def a():
return spark.read.table("source").filter(...)
@dp.materialized_view()
def b():
return spark.read.table("a").groupBy(...)
@dp.materialized_view()
def c():
return spark.read.table("a").groupBy(...)
Cela garantit que le pipeline a une connaissance complète des Transformations dans votre vue pendant la planification du pipeline et prévient les problèmes potentiels liés à l'exécution de code Python arbitraire en dehors des définitions de dataset.
Dans votre fonction, vous pouvez enchaîner des DataFrames pour créer de nouveaux DataFrames sans écrire de résultats incrémentiels sous forme de vues, de vues matérialisées ou de tables de streaming, comme dans l'exemple suivant :
from pyspark import pipelines as dp
@dp.table()
def multiple_transformations():
df1 = spark.read.table("source").filter(...)
df2 = df1.groupBy(...)
return df2.filter(...)
Si tous vos DataFrames effectuent leurs lectures initiales en utilisant la logique de traitement par batch, votre résultat renvoyé est un DataFrame statique. Si vous avez des requêtes qui sont en streaming, votre résultat renvoyé est un DataFrame de streaming.
Retourner un DataFrame
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 batch. La plupart des autres décorateurs fonctionnent à la fois sur des DataFrames de streaming et statiques, tandis que quelques-uns nécessitent un DataFrame de streaming.
La fonction utilisée pour définir un dataset doit renvoyer un DataFrame Spark. N'utilisez jamais de méthodes qui enregistrent ou écrivent dans des fichiers ou des tables dans le cadre du code de votre dataset de pipeline.
Exemples d'Opérations Apache Spark qui ne devraient jamais être utilisées dans le code de pipeline :
collect()count()toPandas()save()saveAsTable()start()toTable()
Les pipelines prennent également en charge l'utilisation de Pandas sur Spark pour les fonctions de définition de dataset. Voir l’ API Pandas sur Spark.
Utilisez SQL dans un pipeline Python
PySpark prend en charge l'opérateur spark.sql pour écrire du code DataFrame à l'aide de SQL. Lorsque vous utilisez ce modèle dans le code source d'un pipeline, il est compilé en vues matérialisées ou en tables de streaming.
L'exemple de code suivant équivaut à l'utilisation de spark.read.table("catalog_name.schema_name.table_name") pour la logique de query du dataset :
@dp.materialized_view
def my_table():
return spark.sql("SELECT * FROM catalog_name.schema_name.table_name")
dlt.read et dlt.read_stream (hérité)
L'ancien module dlt comprend les fonctions dlt.read() et dlt.read_stream() qui ont été introduites pour prendre en charge les fonctionnalités dans le mode de publication de pipeline hérité. Ces méthodes sont prises en charge, mais Databricks recommande d'utiliser toujours les fonctions spark.read.table() et spark.readStream.table() pour les raisons suivantes :
- Les fonctions
dltont un support limité pour la lecture des datasets définis en dehors du pipeline actuel. - Les fonctions
sparkprennent en charge la spécification d'options, telles queskipChangeCommits, pour les opérations de lecture. La spécification d'options n'est pas prise en charge par les fonctionsdlt. - Le module
dlta été remplacé par le modulepyspark.pipelines. Databricks recommande d'utiliserfrom pyspark import pipelines as dppour importerpyspark.pipelinesà utiliser lors de l'écriture de code de pipelines en Python.