Principes de base de PySpark
Cet article présente des exemples simples pour illustrer l'utilisation de PySpark. Il suppose que vous comprenez les concepts fondamentaux d'Apache Spark et que vous exécutez des commandes dans un Notebook Databricks connecté à compute. Vous créez des DataFrames à l'aide d'échantillons de données, effectuez des transformations de base, y compris des opérations sur les lignes et les colonnes de ces données, combinez plusieurs DataFrames et agrégez ces données, visualisez ces données, puis les enregistrez dans une table ou un fichier.
upload des données
Certains exemples de cet article utilisent des données d’exemple fournies par Databricks pour démontrer l’utilisation des DataFrames pour charger, transformer et enregistrer des données. Si vous souhaitez utiliser vos propres données qui ne sont pas encore dans Databricks, vous pouvez d’abord les upload et créer un DataFrame à partir d’elles. Voir Créer ou modifier un tableau à l’aide du fichier upload et Travailler avec des fichiers dans des volumes Unity Catalog.
À propos des exemples de données Databricks
Databricks fournit des exemples de données dans le catalogue samples et dans le répertoire /databricks-datasets.
- Pour accéder aux exemples de données dans le catalogue
samples, utilisez le formatsamples.<schema-name>.<table-name>. Cet article utilise des tables dans le schémasamples.tpch, qui contient des données d’une entreprise fictive. La tablecustomercontient des information sur les clients, et la tableorderscontient des information sur les commandes passées par ces clients. - Utilisez
dbutils.fs.lspour explorer les données dans/databricks-datasets. Utilisez Spark SQL ou DataFrames pour query les données à cet emplacement à l’aide de chemins de fichier. Pour en savoir plus sur les exemples de données fournis par Databricks, consultez Exemples de datasets.
Importer des types de données
De nombreuses opérations PySpark exigent l'utilisation de fonctions SQL ou l'interaction avec des types Spark natifs. Soit vous importez directement uniquement les fonctions et les types dont vous avez besoin, soit, pour éviter de remplacer les fonctions intégrées de Python, vous importez ces modules en utilisant un alias commun.
# import select functions and types
from pyspark.sql.types import IntegerType, StringType
from pyspark.sql.functions import floor, round
# import modules using an alias
import pyspark.sql.types as T
import pyspark.sql.functions as F
Pour une liste complète des types de données, consultez PySpark Data Types.
Pour une liste complète des fonctions PySpark SQL, consultez les fonctions PySpark.
Créer un DataFrame
Il existe plusieurs façons de créer un DataFrame. Habituellement, vous définissez un DataFrame à partir d'une source de données, telle qu'une table ou une collection de fichiers. Ensuite, comme décrit dans la section des concepts fondamentaux d'Apache Spark, utilisez une action, telle que display, pour déclencher l'exécution des transformations. La méthode display génère des DataFrames.
Créer un DataFrame avec des valeurs spécifiées
Pour créer un DataFrame avec des valeurs spécifiées, utilisez la méthode createDataFrame, où les lignes sont exprimées sous forme de liste de tuples :
df_children = spark.createDataFrame(
data = [("Mikhail", 15), ("Zaky", 13), ("Zoya", 8)],
schema = ['name', 'age'])
display(df_children)
Vous remarquerez dans la sortie que les types de données des colonnes de df_children sont automatiquement inférés. Vous pouvez également spécifier les types en ajoutant un schéma. Les schémas sont définis à l'aide de StructType qui est composé de StructFields qui spécifient le nom, le type de données et un indicateur booléen indiquant s'ils contiennent ou non une valeur nulle. Vous devez importer les types de données de pyspark.sql.types.
from pyspark.sql.types import StructType, StructField, StringType, IntegerType
df_children_with_schema = spark.createDataFrame(
data = [("Mikhail", 15), ("Zaky", 13), ("Zoya", 8)],
schema = StructType([
StructField('name', StringType(), True),
StructField('age', IntegerType(), True)
])
)
display(df_children_with_schema)
Créer un DataFrame à partir d'une table dans Unity Catalog
Pour créer un DataFrame à partir d'une table dans Unity Catalog, utilisez la méthode table en identifiant la table au format <catalog-name>.<schema-name>.<table-name>. Cliquez sur Catalogue dans la barre de navigation de gauche pour utiliser Catalog Explorer afin d'accéder à votre table. Cliquez dessus, puis sélectionnez Copier le chemin de la table pour insérer le chemin de la table dans le notebook.
L'exemple suivant charge la table samples.tpch.customer, mais vous pouvez également fournir le chemin d'accès à votre propre table.
df_customer = spark.table('samples.tpch.customer')
display(df_customer)
Créez un DataFrame à partir d'un fichier upload
Pour créer un DataFrame à partir d'un fichier que vous avez upload vers des volumes Unity Catalog, utilisez la propriété read. Cette méthode renvoie un DataFrameReader, que vous pouvez ensuite utiliser pour lire le format approprié. Parce que read est une transformation, Spark ne charge pas les données tant que vous n'appelez pas une action telle que display. Cliquez sur Catalogue dans la barre latérale et utilisez l'explorateur de catalogue pour localiser votre fichier. Sélectionnez-le, puis cliquez sur Copier le chemin .
L'exemple ci-dessous lit un fichier *.csv, mais DataFrameReader prend en charge l'upload de fichiers dans de nombreux autres formats. Voir méthodes DataFrameReader.
# Assign this variable your full volume file path
volume_file_path = ""
df_csv = (spark.read
.format("csv")
.option("header", True)
.option("inferSchema", True)
.load(volume_file_path)
)
display(df_csv)
Pour plus d'information sur les volumes Unity Catalog, consultez Que sont les volumes Unity Catalog ?.
Créer un DataFrame à partir d'une réponse JSON
Pour créer un DataFrame à partir d’une charge utile de réponse JSON renvoyée par une API REST, utilisez le package Python requests pour query et analyser la réponse. Vous devez importer le package pour l’utiliser. Cet exemple utilise des données de la base de données des demandes de médicaments de la Food and Drug Administration des États-Unis.
import requests
# Download data from URL
url = "https://api.fda.gov/drug/drugsfda.json?limit=100"
response = requests.get(url)
# Create the DataFrame
df_drugs = spark.createDataFrame(response.json()["results"])
display(df_drugs)
Pour plus d’informations sur l’utilisation de JSON et d’autres données semi-structurées sur Databricks, consultez Modéliser les données semi-structurées.
Sélectionnez un champ ou un objet JSON
Pour sélectionner un champ ou un objet spécifique à partir du JSON converti, utilisez la notation []. Par exemple, pour sélectionner le champ products qui est lui-même un tableau de produits :
display(df_drugs.select(df_drugs["products"]))
Vous pouvez également chaîner des appels de méthode pour parcourir plusieurs champs. Par exemple, pour afficher le nom de la marque du premier produit dans une application de médicament :
display(df_drugs.select(df_drugs["products"][0]["brand_name"]))
Créer un DataFrame à partir d'un fichier
Pour illustrer la création d’un DataFrame à partir d’un fichier, cet exemple charge les données CSV dans le répertoire /databricks-datasets.
Pour accéder aux datasets d'exemple, vous pouvez utiliser les commandes du système de fichiers Databricks Utilities. L'exemple suivant utilise dbutils pour répertorier les datasets disponibles dans /databricks-datasets:
display(dbutils.fs.ls('/databricks-datasets'))
Alternativement, vous pouvez utiliser %fs pour accéder aux commandes du système de fichiers Databricks CLI, comme le montre l'exemple suivant :
%fs ls '/databricks-datasets'
Pour créer un DataFrame à partir d'un fichier ou d'un répertoire de fichiers, spécifiez le chemin dans la méthode load :
df_population = (spark.read
.format("csv")
.option("header", True)
.option("inferSchema", True)
.load("/databricks-datasets/samples/population-vs-price/data_geo.csv")
)
display(df_population)
Transformer des données avec les DataFrames
Les DataFrames facilitent la transformation des données à l'aide de méthodes intégrées pour trier, filtrer et agréger les données. De nombreuses Transformations ne sont pas spécifiées comme des méthodes sur les DataFrames, mais sont plutôt fournies dans le package pyspark.sql.functions. Voir les fonctions SQL PySpark de Databricks.
- Opérations sur les colonnes
- Opérations de ligne
- Joindre des DataFrames
- Agréger les données
- Chaînage d'appels
Opérations de colonne
Spark fournit de nombreuses opérations de colonne de base :
- Sélectionner les colonnes
- Créer des colonnes
- Renommer les colonnes
- Convertir les types de colonnes
- Supprimer des colonnes
Pour sortir toutes les colonnes d'un DataFrame, utilisez columns, par exemple df_customer.columns.
Sélectionner des colonnes
Vous pouvez sélectionner des colonnes spécifiques en utilisant select et col. La fonction col se trouve dans le sous-module pyspark.sql.functions.
from pyspark.sql.functions import col
df_customer.select(
col("c_custkey"),
col("c_acctbal")
)
Vous pouvez également faire référence à une colonne en utilisant expr qui prend une expression définie comme une chaîne de caractères :
from pyspark.sql.functions import expr
df_customer.select(
expr("c_custkey"),
expr("c_acctbal")
)
Vous pouvez également utiliser selectExpr, qui accepte les expressions SQL :
df_customer.selectExpr(
"c_custkey as key",
"round(c_acctbal) as account_rounded"
)
Pour sélectionner des colonnes à l'aide d'un littéral de chaîne, procédez comme suit :
df_customer.select(
"c_custkey",
"c_acctbal"
)
Pour sélectionner explicitement une colonne à partir d’un DataFrame spécifique, vous pouvez utiliser l’opérateur [] ou l’opérateur .. (L’opérateur . ne peut pas être utilisé pour sélectionner des colonnes commençant par un entier, ou des colonnes contenant un espace ou un caractère spécial.) Cela peut être particulièrement utile lorsque vous joignez des DataFrames dont certaines colonnes ont le même nom.
df_customer.select(
df_customer["c_custkey"],
df_customer["c_acctbal"]
)
df_customer.select(
df_customer.c_custkey,
df_customer.c_acctbal
)
Créer des colonnes
Pour créer une nouvelle colonne, utilisez la méthode withColumn. L'exemple suivant crée une nouvelle colonne qui contient une valeur booléenne indiquant si le solde du compte client c_acctbal dépasse 1000:
df_customer_flag = df_customer.withColumn("balance_flag", col("c_acctbal") > 1000)
Renommer les colonnes
Pour renommer une colonne, utilisez la méthode withColumnRenamed, qui accepte les noms de colonne existants et nouveaux :
df_customer_flag_renamed = df_customer_flag.withColumnRenamed("balance_flag", "balance_flag_renamed")
La méthode alias est particulièrement utile lorsque vous souhaitez renommer vos colonnes dans le cadre d'agrégations :
from pyspark.sql.functions import avg
df_segment_balance = df_customer.groupBy("c_mktsegment").agg(
avg(df_customer["c_acctbal"]).alias("avg_account_balance")
)
display(df_segment_balance)
Convertir les types de colonnes
Dans certains cas, vous souhaiterez peut-être modifier le type de données d'une ou plusieurs colonnes de votre DataFrame. Pour ce faire, utilisez la méthode cast pour convertir les types de données de colonne. L'exemple suivant montre comment convertir une colonne d'un type entier en un type chaîne de caractères, en utilisant la méthode col pour référencer une colonne :
from pyspark.sql.functions import col
df_casted = df_customer.withColumn("c_custkey", col("c_custkey").cast(StringType()))
print(type(df_casted))
Supprimer les colonnes
Pour supprimer des colonnes, vous pouvez omettre des colonnes lors d'une sélection ou select(*) except ou vous pouvez utiliser la méthode drop :
df_customer_flag_renamed.drop("balance_flag_renamed")
Vous pouvez également supprimer plusieurs colonnes à la fois :
df_customer_flag_renamed.drop("c_phone", "balance_flag_renamed")
Opérations de ligne
Spark fournit de nombreuses Opérations de ligne de base :
- Filtrer les lignes
- Supprimer les lignes en double
- Gérer les valeurs nulles
- Ajouter des lignes
- Trier les lignes
- Filtrer les lignes
Filtrer les lignes
Pour filtrer les lignes, utilisez la méthode filter ou where sur un DataFrame pour ne renvoyer que certaines lignes. Pour identifier une colonne à filtrer, utilisez la méthode col ou une expression qui s'évalue en colonne.
from pyspark.sql.functions import col
df_that_one_customer = df_customer.filter(col("c_custkey") == 412449)
Pour filtrer sur plusieurs conditions, utilisez des opérateurs logiques. Par exemple, & et | vous permettent de AND et OR conditions, respectivement. L'exemple suivant filtre les lignes où le/la c_nationkey est égal(e) à 20 et le/la c_acctbal est supérieur(e) à 1000.
df_customer.filter((col("c_nationkey") == 20) & (col("c_acctbal") > 1000))
df_filtered_customer = df_customer.filter((col("c_custkey") == 412446) | (col("c_custkey") == 412447))
Supprimer les lignes en double
Pour dédoublonner les lignes, utilisez distinct, qui renvoie uniquement les lignes uniques.
df_unique = df_customer.distinct()
Gérer les valeurs nulles
Pour gérer les valeurs nulles, supprimez les lignes qui contiennent des valeurs nulles à l’aide de la méthode na.drop. Cette méthode vous permet de spécifier si vous souhaitez supprimer les lignes contenant any valeurs nulles ou all valeurs nulles.
Pour supprimer les valeurs nulles, utilisez l'un des exemples suivants.
df_customer_no_nulls = df_customer.na.drop()
df_customer_no_nulls = df_customer.na.drop("any")
Si, en revanche, vous souhaitez uniquement filtrer les lignes qui ne contiennent que des valeurs nulles, utilisez ce qui suit :
df_customer_no_nulls = df_customer.na.drop("all")
Vous pouvez appliquer ceci pour un sous-ensemble de colonnes en le spécifiant, comme indiqué ci-dessous :
df_customer_no_nulls = df_customer.na.drop("all", subset=["c_acctbal", "c_custkey"])
Pour remplir les valeurs manquantes, utilisez la méthode fill. Vous pouvez choisir d'appliquer cela à toutes les colonnes ou à un sous-ensemble de colonnes. Dans l'exemple ci-dessous, les soldes de compte qui ont une valeur nulle pour leur solde de compte c_acctbal sont remplis avec 0.
df_customer_filled = df_customer.na.fill("0", subset=["c_acctbal"])
Pour remplacer des chaînes par d'autres valeurs, utilisez la méthode replace. Dans l'exemple ci-dessous, toutes les chaînes d'adresse vides sont remplacées par le mot UNKNOWN:
df_customer_phone_filled = df_customer.na.replace([""], ["UNKNOWN"], subset=["c_phone"])
Ajouter des lignes
Pour ajouter des lignes, vous devez utiliser la méthode union pour créer un nouveau DataFrame. Dans l'exemple suivant, le DataFrame df_that_one_customer créé précédemment et df_filtered_customer sont combinés, ce qui renvoie un DataFrame avec trois clients :
df_appended_rows = df_that_one_customer.union(df_filtered_customer)
display(df_appended_rows)
Vous pouvez également combiner des DataFrames en les écrivant dans une table, puis en ajoutant de nouvelles lignes. Pour les charges de travail de production, le traitement incrémental des sources de données vers une table cible peut réduire considérablement la latence et les coûts de compute à mesure que les données augmentent en taille. Consultez les connecteurs standard dans Lakeflow Connect.
Trier les lignes
Le tri peut être coûteux à l'échelle, et si vous stockez des données triées et rechargez les données avec Spark, l'ordre n'est pas garanti. Assurez-vous d'être intentionnel dans votre utilisation du tri.
Pour trier les lignes par une ou plusieurs colonnes, utilisez la méthode sort ou orderBy. Par default, ces méthodes trient par ordre croissant :
df_customer.orderBy(col("c_acctbal"))
Pour filtrer par ordre décroissant, utilisez desc:
df_customer.sort(col("c_custkey").desc())
L'exemple suivant montre comment trier sur deux colonnes :
df_sorted = df_customer.orderBy(col("c_acctbal").desc(), col("c_custkey").asc())
df_sorted = df_customer.sort(col("c_acctbal").desc(), col("c_custkey").asc())
Pour limiter le nombre de lignes à retourner une fois le DataFrame trié, utilisez la méthode limit. L'exemple suivant affiche uniquement les 10 premiers résultats :
display(df_sorted.limit(10))
Joindre les DataFrames
Pour joindre deux ou plusieurs DataFrames, utilisez la méthode join. Vous pouvez spécifier comment vous souhaitez que les DataFrames soient joints dans les how (le type de jointure) et on (sur quelles colonnes baser la jointure) parameters. Les types de jointures courants comprennent :
inner: Il s'agit du default de type de jointure, qui renvoie un DataFrame qui ne conserve que les lignes où il y a une correspondance pour le parameteronà travers les DataFrames.left: Cela conserve toutes les lignes du premier DataFrame spécifié et uniquement les lignes du second DataFrame spécifié qui correspondent au premier.outer: Une jointure externe conserve toutes les lignes des deux DataFrames, quelle que soit la correspondance.
Pour des informations détaillées sur les jointures, consultez Travailler avec les jointures sur Databricks. Pour une liste des jointures prises en charge dans PySpark, consultez Jointures de DataFrame.
L'exemple suivant renvoie un seul DataFrame où chaque ligne du DataFrame orders est jointe à la ligne correspondante du DataFrame customers. Une jointure interne est utilisée, car l'attente est que chaque commande corresponde exactement à un seul client.
df_customer = spark.table('samples.tpch.customer')
df_order = spark.table('samples.tpch.orders')
df_joined = df_order.join(
df_customer,
on = df_order["o_custkey"] == df_customer["c_custkey"],
how = "inner"
)
display(df_joined)
Pour joindre plusieurs conditions, utilisez des opérateurs booléens tels que & et | pour spécifier AND et OR, respectivement. L'exemple suivant ajoute une condition supplémentaire, en filtrant uniquement les lignes où o_totalprice est supérieur à 500,000:
df_customer = spark.table('samples.tpch.customer')
df_order = spark.table('samples.tpch.orders')
df_complex_joined = df_order.join(
df_customer,
on = ((df_order["o_custkey"] == df_customer["c_custkey"]) & (df_order["o_totalprice"] > 500000)),
how = "inner"
)
display(df_complex_joined)
Agréger les données
Pour agréger des données dans un DataFrame, de manière similaire à un GROUP BY en SQL, utilisez la méthode groupBy pour spécifier les colonnes à regrouper et la méthode agg pour spécifier les agrégations. Importez les agrégations courantes, notamment avg, sum, max et min depuis pyspark.sql.functions. L'exemple suivant montre le solde moyen des clients par segment de marché :
from pyspark.sql.functions import avg
# group by one column
df_segment_balance = df_customer.groupBy("c_mktsegment").agg(
avg(df_customer["c_acctbal"])
)
display(df_segment_balance)
from pyspark.sql.functions import avg
# group by two columns
df_segment_nation_balance = df_customer.groupBy("c_mktsegment", "c_nationkey").agg(
avg(df_customer["c_acctbal"])
)
display(df_segment_nation_balance)
Certaines agrégations sont des actions, ce qui signifie qu'elles Trigger des calculs. Dans ce cas, vous n'avez pas besoin d'utiliser d'autres actions pour produire des résultats.
Pour compter les lignes dans un DataFrame, utilisez la méthode count :
df_customer.count()
Chaînage des appels
Les méthodes qui transforment les DataFrames retournent des DataFrames, et Spark n'agit pas sur les transformations tant que les actions ne sont pas appelées. Cette évaluation paresseuse signifie que vous pouvez enchaîner plusieurs méthodes pour plus de commodité et de lisibilité. L'exemple suivant montre comment enchaîner le filtrage, l'agrégation et le tri :
from pyspark.sql.functions import count
df_chained = (
df_order.filter(col("o_orderstatus") == "F")
.groupBy(col("o_orderpriority"))
.agg(count(col("o_orderkey")).alias("n_orders"))
.sort(col("n_orders").desc())
)
display(df_chained)
Visualisez votre DataFrame
Pour visualiser un DataFrame dans un Notebook, cliquez sur le signe + à côté du tableau, en haut à gauche du DataFrame, puis sélectionnez Visualisation pour ajouter un ou plusieurs graphiques basés sur votre DataFrame. Pour plus de détails sur les visualisations, consultez Visualisations dans les notebooks Databricks et l'éditeur SQL.
display(df_order)
Pour effectuer des visualisations supplémentaires, Databricks recommande d'utiliser l'API pandas pour Spark. Le .pandas_api() vous permet de caster vers l'API pandas correspondante pour un DataFrame Spark. Pour plus d'informations, consultez l'API Pandas sur Spark.
Enregistrer vos données
Une fois que vous avez transformé vos données, vous pouvez les enregistrer en utilisant les méthodes DataFrameWriter. Une liste complète de ces méthodes se trouve dans DataFrameWriter. Les sections suivantes expliquent comment enregistrer votre DataFrame en tant que table et en tant que collection de fichiers de données.
Enregistrez votre DataFrame en tant que table
Pour enregistrer votre DataFrame en tant que table dans Unity Catalog, utilisez la méthode write.saveAsTable et spécifiez le chemin au format <catalog-name>.<schema-name>.<table-name>.
df_joined.write.saveAsTable(f"{catalog_name}.{schema_name}.{table_name}")
Écrivez votre DataFrame au format CSV
Pour écrire votre DataFrame au format *.csv, utilisez la méthode write.csv, en spécifiant le format et les options. default, si des données existent au chemin spécifié, l'opération d'écriture échoue. Vous pouvez spécifier l'un des modes suivants pour prendre une action différente :
overwriteremplace toutes les données existantes dans le chemin cible avec le contenu du DataFrame.appendajoute le contenu du DataFrame aux données du chemin cible.ignoreéchoue silencieusement l'écriture si des données existent dans le chemin d'accès cible.
L’exemple suivant montre comment écraser des données avec le contenu de DataFrame sous forme de fichiers CSV :
# Assign this variable your file path
file_path = ""
(df_joined.write
.format("csv")
.mode("overwrite")
.write(file_path)
)
Ressources supplémentaires
Pour tirer parti de davantage de capacités Spark sur Databricks, consultez :