Écrire en Spark : Parquet, Delta Lake, petits fichiers

Comment écrire en Parquet ou Delta Lake sans noyer le stockage : partitionnement, taille des fichiers, et le moment où Delta devient indispensable.

6 min de lecturesparkscalaparquetdelta-lakebig-data

Toute la performance patiemment gagnée sur les partitions et les joins peut être perdue en une seule ligne : un write.parquet(...) mal calibré crée trois millions de fichiers de 200 Ko, votre stockage objet part en tarif dégradé, et le prochain job qui lit ces données passe la moitié de son temps à ouvrir des connexions. C’est un problème très courant, très visible dans les factures, et étonnamment facile à corriger.

Cet article couvre les deux formats qui comptent aujourd’hui — Parquet et Delta Lake — les erreurs classiques à l’écriture, et la règle simple pour dimensionner les fichiers correctement.

Pourquoi Parquet a gagné

Parquet est un format colonnaire, compressé et splittable. Trois propriétés qui changent tout côté lecture :

  • Colonnaire — lire uniquement user_id et ts sans charger les 40 autres colonnes du fichier. Sur des tables larges, c’est un facteur 10.
  • Compressé par colonne (Snappy, ZSTD, GZIP). Les répétitions et les faibles cardinalités se compressent extrêmement bien.
  • Splittable — un fichier peut être lu en parallèle par plusieurs tâches Spark.

À cela s’ajoutent les statistiques stockées dans le footer (min/max, count, null count par row group) qui permettent le predicate pushdown : un WHERE date = '2026-08-10' fait sauter les row groups entiers sans lire les colonnes.

CSV et JSON n’ont rien de tout cela. Ils restent utiles pour l’export vers un humain ou un système ancien — jamais pour une couche analytique.

Le partitionnement, ce qu’il faut faire et pas faire

Un mot piégeux : « partition » désigne deux choses en Spark.

SensRôle
Partition en mémoireRuntimeUnité de parallélisme des tâches
Partition sur disqueLayout physiqueSous-dossier col=valeur/ créé par l’écriture

Le partitionBy de .write fait la seconde :

scala
df.write
  .partitionBy("dt")
  .parquet("s3://bucket/events/")

// Produit :
// s3://bucket/events/dt=2026-08-08/
// s3://bucket/events/dt=2026-08-09/
// s3://bucket/events/dt=2026-08-10/

Le pushdown de partition est le plus efficace des filtres : un WHERE dt = '2026-08-10' fait littéralement ignorer les autres dossiers, sans même les lister au niveau des fichiers.

Les règles à respecter :

  • Partitionner par une colonne de faible cardinalité (jour, région, type d’événement). Pas par user_id ni par session_id — vous créez des millions de dossiers.
  • Une valeur de partition doit peser au moins 1 Go dans la plupart des cas. Sinon, le coût de gestion (listing, ouverture, metadata) dépasse le gain.
  • Maximum une ou deux colonnes de partition. partitionBy("year", "month", "day") est presque toujours pire que partitionBy("dt").

Le problème des petits fichiers

Un DataFrame de 200 partitions en mémoire écrit avec partitionBy("dt") sur 30 jours de données produit potentiellement 6 000 fichiers. Beaucoup sont minuscules — Spark en produit un par (partition en mémoire) × (valeur de partition sur disque) qui contient au moins une ligne.

C’est le small files problem, et il coûte cher :

  • Le stockage objet (S3, GCS, Azure Blob) facture par requête. Un million de petits fichiers = un million de GET par lecture.
  • Le driver Spark liste tous les fichiers avant de démarrer un job. Sur 100 000 fichiers, le listing seul prend plusieurs minutes.
  • Chaque fichier a un footer Parquet à décoder. À 200 Ko par fichier, le coût par octet lu est absurde.

La cible réaliste : des fichiers entre 128 Mo et 1 Go.

Le geste qui répare : contrôler les partitions en mémoire avant l’écriture

scala
// Mauvais : 200 partitions en mémoire × 30 jours = beaucoup de petits fichiers
df.write.partitionBy("dt").parquet("s3://bucket/events/")

// Bon : on redistribue en amont pour viser 1 fichier par jour et par partition
df.repartition($"dt")
  .write
  .partitionBy("dt")
  .parquet("s3://bucket/events/")

repartition($"dt") redistribue les données par jour avant l’écriture : toutes les lignes d’un même jour arrivent dans la même partition en mémoire, qui produit un seul fichier par valeur de dt.

Si un jour est trop volumineux (par exemple 5 Go), on force une sous-division :

scala
df.repartition($"dt", floor(rand() * 5))
  .write
  .partitionBy("dt")
  .parquet(...)
// 5 fichiers de ~1 Go par jour
Et `coalesce` avant l’écriture ?

coalesce(n) réduit sans shuffle et fonctionne quand vous êtes déjà partitionné correctement — typiquement à la fin d’un pipeline qui a fait un groupBy et où vous voulez limiter le nombre de fichiers de sortie. Mais coalesce ne rééquilibre pas : si les données sont réparties inégalement, vos fichiers de sortie le seront aussi. Pour partir d’un DataFrame qui a subi plusieurs stages, repartition est le choix par défaut.

Delta Lake : quand Parquet ne suffit plus

Delta Lake est Parquet + un log de transactions JSON. Ce petit ajout débloque tout ce qui manque à Parquet nu :

CapacitéParquetDelta Lake
Écriture atomique (tout ou rien)NonOui
Lecteurs concurrents pendant une écritureCasseIsolation snapshot
UPDATE, DELETE, MERGENonOui
Time travel (lire une ancienne version)NonOui
Évolution du schémaManuelleOui
Compaction / vacuumManuelleOPTIMIZE + VACUUM

Concrètement, ce qui change au quotidien :

scala
// Un upsert propre — impossible en Parquet nu
DeltaTable.forPath(spark, path)
  .as("t")
  .merge(updates.as("u"), "t.id = u.id")
  .whenMatched.updateAll()
  .whenNotMatched.insertAll()
  .execute()

// Reprise en arrière d'une release ratée
spark.read
  .format("delta")
  .option("versionAsOf", 42)
  .load(path)

// Compaction automatique en une commande
spark.sql(s"OPTIMIZE delta.`$path`")

Delta Lake devient indispensable dès que :

  • Plusieurs jobs écrivent la même table (concurrence).
  • Vous devez corriger un ligne sur un milliard sans réécrire tout.
  • Vous devez pouvoir revenir en arrière (audit, débogage, RGPD).
  • Vous êtes en streaming (Delta est source et sink stables).

Pour un lot batch simple qui écrit une fois par jour et n’est jamais modifié, Parquet nu reste très bien — pas la peine de payer la complexité du log de transactions.

L’ordre correct des opérations en écriture

Un pipeline de sortie robuste ressemble à ceci :

scala
val output = pipeline(input)      // votre logique métier

output
  .repartition($"dt")             // 1. rééquilibrer par la colonne de partition disque
  .sortWithinPartitions("user_id") // 2. tri intra-fichier → meilleure compression + pushdown
  .write
  .mode("overwrite")
  .partitionBy("dt")               // 3. layout physique
  .option("compression", "zstd")   // 4. ZSTD > Snappy sur presque tous les cas
  .format("delta")                 // 5. Delta plutôt que Parquet, quand pertinent
  .save("s3://bucket/events/")

Chacune de ces cinq lignes change quelque chose de mesurable : la répartition règle le nombre de fichiers, le tri triple souvent la compression, la compression ZSTD gagne 20-40 % sur Snappy pour un coût CPU faible, et Delta ajoute l’ACID sans coût mesurable à la lecture.

Ce qu’il faut retenir

Écrire, c’est un contrat avec les prochains lecteurs — ceux d’autres jobs, ceux d’autres équipes, ceux dans six mois. Ce contrat tient dans trois règles : un fichier de 128 Mo à 1 Go, une ou deux colonnes de partitionnement de faible cardinalité, Delta dès qu’on a besoin d’écrire deux fois au même endroit. Les respecter transforme radicalement la qualité de service d’un lac de données ; les ignorer produit ces plateformes où plus personne ne comprend pourquoi le moindre SELECT prend cinq minutes.

Cet article clôt le cluster Spark. Les cinq pièces ensemble — word count RDD, RDD vs DataFrame, transformations vs actions, partitions et shuffle, joins — couvrent l’essentiel de ce qu’un ingénieur Spark rencontre les six premiers mois. Le reste, c’est du réglage fin sur ces mêmes fondations.

Ce sujet fait partie d’un cours complet

Développement et déploiement de solutions de données — les premiers modules sont en accès libre.

Voir le plan du cours

Continuer sur le même sujet