Standardmäßig ist der Prompt ausgewählt, der zuerst die Quelle prüft. Sie können zu einem direkten Befehl wechseln oder eine lokale Kopie herunterladen.
Quelldateien prüfen
Lesen Sie SKILL.md und alle von SkillsMP angezeigten Begleitdateien, bevor Sie sich für eine Installation entscheiden.
Mit Codex oder Claude installieren Kopieren Sie diesen Prompt, fügen Sie ihn in Codex, Claude oder einen anderen Assistant ein und lassen Sie die Skill-Seite prüfen und installieren.
Ein direkter Befehl überspringt den Prüf-Prompt. Prüfen Sie die Quelle, bevor Sie ihn ausführen.
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.