
PySpark pour les débutants : au-delà des bases
dans cette série, PySpark pour les débutants : maîtriser les bases alors vous comprenez déjà le cœur de Spark : les données distribuées, les DataFrames et l’exécution paresseuse. Vous avez installé PySpark, lancé une SparkSession, lu un CSV et effectué des manipulations simples des données dans un Dataframe. Je laisserai un lien vers cette histoire à la fin de celle-ci.
Une chose qui mérite d’être répétée dans cet article original est que j’utilise souvent les termes PySpark et Spark de manière interchangeable, mais à proprement parler, Spark est le cadre informatique distribué global (écrit en Scala) et PySpark est une API Python dédiée à Spark.
Au-delà des fondamentaux
Maintenant, quelque chose d’intéressant se produit lorsque vous dépassez ce stade de débutant. Vous réalisez rapidement votre deuxième Le projet PySpark nécessite un état d’esprit légèrement différent :
- Vous souhaitez lire/écrire des données de manière plus sûre, plus rapide et plus prévisible.
- Vous souhaitez combiner des ensembles de données sans vous sentir incertain quant aux jointures.
- Tu veux comprendre pourquoi Spark se comporte comme il le fait – et comment le pousser doucement dans la bonne direction.
Cet article vous guide à travers ces prochaines étapes. C’est délibérément lent et pratique. Pas de profondeur interne. Aucun réglage du cluster. Pas d’optimisations Spark compliquées. Exactement ce que les vrais débutants doivent savoir lorsqu’ils passent d’exemples de jouets à de petits travaux du monde réel.
Nous utilisons Spark open source, exécuté localement, comme avant.
1. Passer à l’étape suivante : lire correctement les données
Dans mon premier article, nous avons utilisé le chargeur CSV le plus simple possible :
df = spark.read.csv("sales.csv", header=True, inferSchema=True)
Cela fonctionne – et c’est bien pour les premières expériences – mais cela cache un problème subtil.
Spark devine vos types de données
Lorsque vous utilisez la directive inferSchema=True, Spark examine un petit échantillon de votre fichier et utilise ces informations pour deviner si une colonne est un entier, une chaîne, un booléen ou un double. Cela signifie :
- Si 99 lignes semblent numériques et que la 100ème ligne est vide, Spark peut interpréter la colonne comme un chaîne.
- Si quelqu’un modifie le fichier la semaine prochaine et ajoute accidentellement 23,50 £ au lieu de 23,50, Spark pourrait traiter la colonne entière différemment.
- Si votre fichier est volumineux, l’exemple utilisé par Spark ne représentera pas l’ensemble de données.
Cela peut conduire à un comportement mystérieux plus tard, le genre de bugs que les débutants trouvent les plus difficiles à diagnostiquer.
Une meilleure habitude pour les débutants : définir un schéma pour vos données
Considérez un schéma comme la version Spark d’un modèle de lecture de données. Avant de construire quoi que ce soit, vous dites à Spark des choses telles que :
Les noms des colonnes
De quel type de données ils devraient être
Si une valeur de colonne est facultative ou non.
Voici à quoi cela ressemble pour notre exemple de données de ventes. Rappelez-vous que les données ressemblaient à ceci :
transaction_id,customer_name,net_amount,tax_amount, is_member
101,Alice,250.50,25.05,true
102,Bob,120.00,6.00, false
103,Charlie,450.75,25.07,true
104,David,89.99,5.73,false
Pour spécifier les types des champs ci-dessus dans Spark, nous définissons notre schéma en utilisant un code comme celui-ci.
from pyspark.sql import types as T
schema = T.StructType([
T.StructField("transaction_id", T.IntegerType(), False),
T.StructField("customer_name", T.StringType(), False),
T.StructField("net_amount", T.DoubleType(), True),
T.StructField("tax_amount", T.DoubleType(), True),
T.StructField("is_member", T.BooleanType(), True),
])
Les noms de colonnes et les paramètres de type sont explicites. Le vrai[False] Le paramètre indique qu’il peut y avoir [not] être des valeurs NULL dans la colonne. Notez que l’indicateur de nullabilité Vrai/Faux est principalement constitué de métadonnées de schéma et d’informations d’optimisation. Elle n’est pas toujours strictement appliquée pour chaque source de données, comme l’est une contrainte de base de données NOT NULL.
Options plus utiles lors de la lecture de données CSV
Il existe de nombreuses options de lecture CSV pratiques que vous pouvez combiner avec la directive de schéma qui rendent le chargement des données CSV encore plus fiable.
Les options les plus courantes incluent :
- mode=”PERMISSIVE” : conserve autant que possible les mauvaises lignes
- mode=”DROPMALFORMED” : supprime les lignes mal formées
- mode=”FAILFAST” : erreurs immédiatement
- en-tête = Vrai[False]: Le fichier contient-il [or not] un enregistrement d’en-tête
- nullValue : quel texte doit remplacer les valeurs nulles dans l’entrée
- formatdate/horodatageFormat
Nous pouvons maintenant charger les sales_data dans un Dataframe comme ceci :
df = (
spark.read
.option("header", True)
# Other modes: "PERMISSIVE" and "DROPMALFORMED".
.option("mode", "FAILFAST")
.option("nullValue", "N/A")
.schema(schema)
.csv("sales_data.csv")
)
Pourquoi est-ce important pour les débutants ?
- Toi savoir quels sont les types de données avant de commencer à travailler.
- Si spécifié, Spark rejeter les lignes étranges au lieu de les interpréter silencieusement.
- Vos transformations deviennent plus prévisible.
- Si vous joignez deux ensembles de données plus tard, incompatibilités de types ne vous surprendra pas.
2. Comprendre les transformations de données
Rappelez-vous, dans mon article précédent, lors de nos premières étapes de manipulation des dataframes avec PySpark, nous avons ajouté une colonne dérivée supplémentaire à notre Dataframe en utilisant un code comme celui-ci :
df2 = df.withColumn("gross_amount", df.net_amount + df.tax_amount)
J’ai expliqué que cette ligne ne calcule encore rien. Cela ajoute simplement une étape au plan interne de Spark :
1. Read the CSV
2. Add a new column (gross_amount = net + tax)
Ensuite, vous pouvez ajouter d’autres étapes comme celle-ci :
df3 = df2.withColumn("tax_percentage", df2.tax_amount / df2.gross_amount * 100)
Pourtant, aucun calcul n’a été effectué. Seulement lorsque vous effectuez un action comme …
df3.show()
… Spark dit-il :
« D’accord, maintenant je dois exécuter toutes ces étapes. »
C’est ce que signifie « exécution paresseuse », mais l’élément important pour les débutants n’est pas le nom. C’est l’effet, et ça signifie,
- Vous pouvez enchaîner de nombreuses transformations sans les « payer » jusqu’à ce que vous ayez besoin du résultat.
- Spark peut réorganiser l’ordre en interne pour gérer les choses efficacement.
- Vous ne perdez pas de temps à effectuer des étapes intermédiaires sur des données que vous pourriez filtrer ultérieurement.
Pensez-y comme à une tâche quotidienne, comme préparer un sandwich :
- Vous rassemblez tous les ingrédients.
- Vous l’assemblez dans votre esprit.
- Vous ne commencez réellement à couper et à préparer que lorsque vous savez précisément ce que vous préparez.
3. Nettoyer les données avant qu’elles ne causent des problèmes
Les données réelles sont généralement désordonnées et contiennent souvent des valeurs manquantes, des chaînes vides, des enregistrements en double ou des valeurs d’espace réservé telles que « N/A » et « inconnu ».
Dans PySpark, l’objectif est de détecter et de résoudre rapidement les problèmes évidents afin que le reste de votre flux de travail se comporte de manière prévisible. PySpark dispose d’un certain nombre de fonctions utiles qui vous permettent de le faire.
Suppression de lignes avec des valeurs manquantes
La fonction de nettoyage la plus simple est dropna().
df_clean = df.dropna()
Cela supprime toute ligne contenant une valeur nulle dans une colonne. Cela peut être utile, mais c’est souvent trop agressif.
Le plus souvent, vous supprimez uniquement les lignes dans lesquelles des colonnes importantes sont manquantes :
df_clean = df.dropna(subset=["net_amount", "tax_amount"])
Cela signifie:
Conservez la ligne tant que net_amount et tax_amount sont présents.
D’autres colonnes peuvent toujours contenir des valeurs nulles, et cela peut convenir.
Remplir les valeurs manquantes
Parfois, vous ne souhaitez pas supprimer des lignes. Vous voulez simplement remplacer les valeurs manquantes par quelque chose de sensé.
C’est là que fillna() est utile.
df_clean = df.fillna({"city": "Unknown"})
Vous pouvez également remplir des colonnes numériques :
df_clean = df.fillna({"tax_amount": 0.0})
Ceci est utile lorsqu’une valeur manquante a une signification claire. Par exemple, un montant de remise manquant peut raisonnablement devenir 0,0. Mais soyez prudent. Remplir les valeurs manquantes peut changer la signification de vos données si vous choisissez la mauvaise valeur par défaut.
Changer les types de colonnes avec cast()
Parfois, Spark lit une colonne avec un type incorrect, en particulier lorsque vous travaillez avec des fichiers CSV. Si tel est le cas, vous pouvez convertir une colonne en utilisant l’opérateur cast() :
from pyspark.sql import functions as F
df_clean = df.withColumn("net_amount",F.col("net_amount").cast("double") )
Ceci est particulièrement courant lorsque des dates, des nombres ou des booléens ont été lus sous forme de chaînes.
Suppression des lignes en double
Des lignes en double peuvent apparaître lorsque les fichiers sont exportés plusieurs fois, joints de manière incorrecte ou combinés à partir de plusieurs sources. Vous pouvez supprimer les doublons exacts comme ceci :
df_clean = df.dropDuplicates()
Ou supprimez les doublons en fonction d’une ou plusieurs colonnes sélectionnées.
df_clean = df.dropDuplicates(["transaction_id"])
Cette deuxième version est souvent plus utile car elle dit :
Chaque identifiant de transaction ne doit apparaître qu’une seule fois.
Un petit exemple de nettoyage de données
Rassembler ces idées :
from pyspark.sql import functions as F
df_clean = (
df
# Remove transactions missing required values.
.dropna(subset=["transaction_id", "net_amount"])
# Supply defaults for optional values.
.fillna(
{
"city": "Unknown",
"tax_amount": 0.0,
}
)
# Apply the expected numeric types.
.withColumn(
"net_amount",
F.col("net_amount").cast("double"),
)
.withColumn(
"tax_amount",
F.col("tax_amount").cast("double"),
)
# Keep one row for each transaction.
.dropDuplicates(["transaction_id"])
)
4. Rejoindre des ensembles de données dans PySpark sans se perdre
Si vous avez déjà travaillé avec des bases de données, vous avez probablement écrit des instructions SQL qui joignent deux ou plusieurs tables. Les jointures dans Spark fonctionnent de la même manière, mais sur des Dataframes.
Qu’est-ce qu’une jointure ?
Si le concept de jointure est nouveau pour vous, il s’agit d’un moyen de faire correspondre les lignes d’un DataFrame avec les lignes associées d’un autre DataFrame. En d’autres termes, il répond à une question telle que :
« Quelles lignes de ce DataFrame correspondent aux lignes de ce DataFrame ? »
C’est l’idée principale derrière chaque jointure dans PySpark. Une fois cette partie claire, la syntaxe et les types de jointure deviennent beaucoup plus faciles à comprendre.
Si vous avez deux Dataframes comme ceci :
sales_data.csv
transaction_id, customer_name, net_amount, tax_amount
101, Alice, 250.50, 25.05
102, Bob, 120.00, 6.00
clients.csv
customer_name, city, loyalty_level
Alice, New York, Gold
Bob, London, Silver
Vous pouvez les rejoindre sur leur champ customer_name commun comme ceci :
df_sales = spark.read.csv("sales_data.csv", header=True)
df_customers = spark.read.csv("customers.csv", header=True)
df_joined = df_sales.join(df_customers, on="customer_name", how="inner")
df_joined.show()
# Output
+-------------+--------------+----------+----------+--------+-------------+
|customer_name|transaction_id|net_amount|tax_amount|city |loyalty_level|
+-------------+--------------+----------+----------+--------+-------------+
|Alice |101 |250.50 |25.05 |New York|Gold |
|Bob |102 |120.00 |6.00 |London |Silver |
+-------------+--------------+----------+----------+--------+-------------+
Quelle jointure les débutants doivent-ils utiliser ?
Il existe plusieurs types de jointures disponibles dans Spark. Pour 99 % des cas d’utilisation pour débutants, vous utiliserez l’un des éléments suivants :
- inner — afficher uniquement les lignes correspondantes
- gauche — affiche tout dans le tableau de gauche, plus les correspondances
- external — affiche toutes les lignes des deux tables
Et parmi celles-ci, la jointure interne sera de loin le type de jointure le plus courant que vous utiliserez dans votre travail quotidien.
Ne vous inquiétez pas encore de la « diffusion », du « tri-fusion », du « shuffle-hash » ou de toute autre stratégie de jointure avancée. Au fur et à mesure que votre expérience de Spark grandit, vous pouvez les lire à votre guise.
N’oubliez pas :
Les jointures sont informatiquement plus coûteuses qu’une simple opération de colonnesalors utilisez-les lorsque cela est nécessaire, mais pas avec désinvolture.
5. Lecture et écriture de données à la manière « Spark » : Parquet
La plupart des débutants s’en tiennent au CSV car il est familier. Mais CSV est lent, rigide et ne prend pas en charge les types de données. Dans la vraie vie, Parquet est le format de données natif de Spark. Parquet est un format de données compressé en colonnes, idéal pour l’analyse de données, le reporting de données et les charges de travail lourdes en lecture.
Lorsque Spark lit un ensemble de données Parquet :
- Il ne charge que les colonnes dont vous avez réellement besoin.
- Il comprend tous les types de données.
- Il se charge beaucoup plus rapidement que CSV.
Vous écrivez le contenu du Dataframe dans des fichiers au format Parquet comme ceci :
df_joined.write.mode("overwrite").parquet("output/enriched_sales")
Ensuite, vous pouvez le relire instantanément comme ceci,
df_fast = spark.read.parquet("output/enriched_sales")
df_fast.show()
N.-B.. L’utilisation de Parquet pour l’entrée et la sortie de fichiers constitue la « mise à niveau » de performances la plus simple pour tout débutant Spark.
6. Penser dans les workflows PySpark
Une fois que vous avez compris comment lire les données, les nettoyer, les transformer, les joindre et les réécrire, l’étape suivante consiste à apprendre à organiser ces actions dans un flux de travail simple. Un projet PySpark débutant suit généralement cette séquence :
Read data
-> check and clean it
-> add useful columns
-> combine with other data
-> write the result
Cela peut paraître évident, mais il s’agit d’un changement important. Vous n’expérimentez plus seulement un DataFrame à la fois. Vous construisez un processus reproductible.
Gardez chaque étape simple
Une habitude utile pour les débutants consiste à donner à chaque étape de votre flux de travail un objectif clair. Par exemple:
df_raw = spark.read.schema(schema).csv("sales_data.csv", header=True)
df_clean = df_raw.dropna(subset=["net_amount", "tax_amount"])
df_enriched = df_clean.withColumn(
"gross_amount",
F.col("net_amount") + F.col("tax_amount")
)
df_final = df_enriched.join(df_customers, on="customer_name", how="left")
df_final.write.mode("overwrite").parquet("output/final_dataset")
Ce style est légèrement plus verbeux que de tout enchaîner en une seule longue expression, mais il est beaucoup plus facile à lire lorsque vous apprenez.
Chaque nom de DataFrame vous indique où vous en êtes dans le workflow :
df_raw -> the data as it arrived
df_clean -> the data after basic cleaning
df_enriched -> the data after adding new meaning
df_final -> the dataset ready to save
Pourquoi c’est important
En cas de problème, cette structure facilite grandement le débogage.
Vous pouvez inspecter chaque étape en examinant les données :
df_raw.show()
df_clean.show()
df_enriched.show()
Vous pouvez vérifier le nombre de lignes :
df_raw.count()
df_clean.count()
df_final.count()
Cela permet de répondre à des questions utiles telles que :
Did rows disappear unexpectedly during cleaning?
Did the join create more rows than expected?
Did a calculated column produce nulls?
Le modèle mental simple de : Entrées → préparation → combinaison → sortie vous mènera étonnamment loin dans votre voyage PySpark.
7. Une introduction douce à l’interface utilisateur Spark
Spark a une jolie petite interface utilisateur Web qui s’active lorsque vous exécutez une action comme .count() ou .write(). Lorsque votre tâche Spark est exécutée localement, visitez :
http://localhost:4040
Vous devriez voir quelque chose comme ceci affiché.

Cela semble un peu écrasant, mais vous n’avez pas besoin de comprendre chaque onglet. À ce stade, il vous suffit de savoir que l’interface utilisateur existe et pourquoi elle est utile. Et c’est utile car cela vous aide à voir quelles tâches Spark ont été exécutées ou sont actuellement en cours d’exécution.
Et, à mesure que votre expérience dans Spark se développe, l’interface utilisateur peut vous aider à comprendre pourquoi les tâches ont échoué ou prennent plus de temps à s’exécuter que prévu. Mais cela arrive bien plus tard. Pour l’instant, traitez l’interface utilisateur Spark comme le tableau de bord de votre voiture : vous n’avez pas besoin de comprendre le moteur pour remarquer quand quelque chose semble étrange.
Résumé : Vous êtes maintenant prêt pour votre premier vrai projet PySpark
À ce stade, vous êtes passé de « Je peux exécuter Spark » à « Je peux créer un pipeline Spark propre et simple ».
Vous savez maintenant comment :
- lire les données en toute sécurité,
- nettoyez-le et préparez-le,
- enrichissez-le de nouvelles colonnes,
- combiner plusieurs ensembles de données,
- enregistrer le résultat efficacement,
- et observez Spark juste assez pour rester confiant.
Rien dans cet article ne nécessitait un cluster. Rien ne nécessitait un réglage avancé. C’est exactement ainsi que commencent de nombreux vrais projets PySpark.
Lorsque vous serez plus expérimenté, vous souhaiterez peut-être développer vos connaissances en recherchant certains de ces sujets.
- lire des plans d’exécution
- comprendre les mélanges
- gestion des partitions
- autres types de jointure
- réglage simple des performances
Ce sont quelques-uns des sujets que j’espère aborder dans un prochain article, mais pour l’instant, vous avez franchi votre prochaine étape majeure et vous pouvez créer quelque chose de significatif et d’utile avec PySpark.
BTW, voici ce lien vers le premier article de cette série,
PySpark pour les débutants : maîtriser les basesdont j’ai parlé au début.



