Aller au contenu principal

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 format samples.<schema-name>.<table-name>. Cet article utilise des tables dans le schéma samples.tpch, qui contient des données d’une entreprise fictive. La table customer contient des information sur les clients, et la table orders contient des information sur les commandes passées par ces clients.
  • Utilisez dbutils.fs.ls pour 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.

Python
# 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 :

Python
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.

Python
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.

Python
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 Icône de données. 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.

Python
# 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.

Python
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 :

Python
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 :

Python
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:

Python
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 :

Python
%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 :

Python
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 de colonne

Spark fournit de nombreuses opérations de colonne de base :

astuce

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.

Python
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 :

Python
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 :

Python
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 :

Python
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.

Python
df_customer.select(
df_customer["c_custkey"],
df_customer["c_acctbal"]
)
Python
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:

Python
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 :

Python
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 :

Python
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 :

Python
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 :

Python
df_customer_flag_renamed.drop("balance_flag_renamed")

Vous pouvez également supprimer plusieurs colonnes à la fois :

Python
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

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.

Python
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.

Python
df_customer.filter((col("c_nationkey") == 20) & (col("c_acctbal") > 1000))
Python
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.

Python
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.

Python
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 :

Python
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 :

Python
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.

Python
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:

Python
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 :

Python
df_appended_rows = df_that_one_customer.union(df_filtered_customer)

display(df_appended_rows)
remarque

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

important

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 :

Python
df_customer.orderBy(col("c_acctbal"))

Pour filtrer par ordre décroissant, utilisez desc:

Python
df_customer.sort(col("c_custkey").desc())

L'exemple suivant montre comment trier sur deux colonnes :

Python
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 :

Python
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 parameter on à 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.

Python
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:

Python
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é :

Python
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)
Python
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 :

Python
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 :

Python
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.

Python
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>.

Python
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 :

  • overwrite remplace toutes les données existantes dans le chemin cible avec le contenu du DataFrame.
  • append ajoute 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 :

Python
# 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 :