Installer avec Codex ou Claude Copiez ce prompt, collez-le dans Codex, Claude ou un autre assistant, puis laissez-le vérifier la page du skill et l'installer pour vous.
Une commande directe contourne le prompt de vérification. Examinez la source avant de l'exécuter.
Cette compétence couvre les patterns de production : déploiement sur cluster (YARN/Kubernetes), optimisation Catalyst/Tungsten, tuning mémoire, gestion des shuffles, et intégration avec le lac de données (Delta Lake, Iceberg, Hudi).
Quand l'utiliser
Activez cette compétence lorsque l'utilisateur :
Demande de traiter de gros volumes de données (Go → Po) avec PySpark ou Scala.
Veut exécuter des transformations ETL distribuées sur un cluster (Databricks, EMR, Glue).
A besoin de streaming structuré à partir de Kafka ou de fichiers.
Pose des questions sur l'optimisation Spark (shuffle, partitionnement, Catalyst).
Veut écrire/lire en Delta Lake, Iceberg ou Parquet avec Spark.
Prérequis
Installation Locale (développement)
pip install pyspark delta-spark pyarrow pandas
Vérification :
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("EVA-Spark-Test") \
.config("spark.sql.adaptive.enabled", "true") \
.getOrCreate()
print(f"Spark {spark.version} OK — {spark.sparkContext.defaultParallelism} partitions par défaut")
spark.stop()
# Le checkpoint est la clé de l'exactly-once# Stocker sur un système fiable (HDFS, S3, ADLS)
query = stream_df.writeStream \
.format("delta") \
.outputMode("append") \
.option("checkpointLocation", "s3://checkpoints/sensors-job/") \
.trigger(processingTime="30 seconds") \
.toTable("bronze.sensor_readings")
4. Optimisation des Performances (Catalyst + Tungsten)
# Répartition optimale
df_partitionne = df.repartition(200, "machine_id") # hash par machine_id
df_range = df.repartitionByRange(50, "timestamp") # range pour ordre temporel
df_coalesce = df.coalesce(50) # réduire sans shuffle# Bucketing (optimise les futures jointures)
df.write.bucketBy(20, "produit_id") \
.sortBy("date_vente") \
.mode("overwrite") \
.saveAsTable("ventes_bucketed")
4.3 Taille des Fichiers
# Contrôle des fichiers de sortie
df.write \
.option("maxRecordsPerFile", "500000") \
.option("parquet.block.size", "256MB") \
.mode("overwrite") \
.parquet("/data/curated/")
4.4 Cache et Persistance
# Persistance en mémoire/série
df_cached = df_clean.cache() # cache en mémoire
df_persist = df_clean.persist(StorageLevel.MEMORY_AND_DISK_SER)
# Nettoyage
df_cached.unpersist()
4.5 Plan d'Exécution
# Visualiser le plan Catalyst
df.explain("formatted") # plan textuel détaillé
df.explain("cost") # plan avec coûts estimés
df.explain("codegen") # code généré par Tungsten# Voir les étapes physiques
df.explain(True)
# 1. Salted key pour répartir la chargefrom pyspark.sql import functions as F
import random
salting = F.concat(F.col("machine_id"), F.lit("_"),
(F.rand() * 10).cast("int"))
df_salted = df_dim.withColumn("salt",
(F.rand() * 10).cast("int"))
df_faits_salted = df_faits.crossJoin(
spark.range(10).toDF("salt")
)
# 2. Ou laisser AQE gérer automatiquement
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256MB")
8. Monitoring et Debugging
8.1 Spark UI
# Accès à l'interface Web (port 4040 par défaut)# http://localhost:4040/# Voir : Jobs, Stages, Storage, Environment, SQL, Executors
8.2 Metrics Programmatiques
# Événements Spark (pour analyse des goulots)
spark.sparkContext.setLogLevel("WARN")
# Accumulators pour le suivi
accum = spark.sparkContext.accumulator(0)
deftrack_batch(df):
accum.add(df.count())
return df
df_tracked = df.rdd.mapPartitions(track_batch)
Pièges Courants (Pitfalls)
Shuffle explosif.
Erreur :groupBy(), join() ou distinct() sur des colonnes à haute cardinalité sans partitionnement préalable. Le shuffle peut dépasser la mémoire disponible.
Correction : Pré-partitionner avec repartition() sur la clé de jointure, activer AQE, et augmenter spark.sql.shuffle.partitions.
Petits fichiers (small file problem).
Erreur : Spark écrit des milliers de petits fichiers Parquet (128 Ko) qui ralentissent les lectures futures.
Correction : Utiliser coalesce(), repartition(), maxRecordsPerFile, delta.autoCompact dans Databricks.
Skew dans les jointures.
Erreur : Une clé de jointure très fréquente (ex : machine_id pour la machine la plus active) submerge un seul executor.
Correction : Activer spark.sql.adaptive.skewJoin.enabled=true ou utiliser des salted keys.
Serialisation Kryo non configurée.
Erreur : Par défaut Spark utilise Java serialization (lent). Les UDFs et classes personnalisées ralentissent.