Spark Optimization

作者 wshobson46891e7e60da無授權條款收錄於 2026年10月8日更新於 2026年10月8日

Optimize Apache Spark jobs with partitioning, caching, shuffle optimization, and memory tuning. Use when improving Spark performance, debugging slow jobs, or scaling data processing pipelines.

AI 產生的概覽

指導透過分割、快取、減少 shuffle 與記憶體調校來最佳化 Apache Spark 作業。

功能
提供用於調校 Apache Spark 作業的正式環境實務模式,涵蓋分割策略、記憶體與執行器設定、shuffle 縮減以及資料傾斜處理。內容包含 PySpark 工作階段設定的快速入門範例、關鍵效能因素表格,以及注意事項清單。更詳細的模式文件在另外的詳情檔案中引用。
適用情境
適用於提升 Spark 效能、排查執行緩慢的作業,或擴充資料處理管線。也適合調校記憶體與執行器設定、實作分割策略,或減少 shuffle 與資料傾斜。
執行需求
不隨附指令碼,僅為說明性內容。執行範例程式碼需要 Apache Spark 與 PySpark,且範例會在 S3 上讀寫 Parquet 資料,需要相應的儲存存取權限。

Apache Spark Optimization

Production patterns for optimizing Apache Spark jobs including partitioning strategies, memory management, shuffle optimization, and performance tuning.

When to Use This Skill

  • Optimizing slow Spark jobs
  • Tuning memory and executor configuration
  • Implementing efficient partitioning strategies
  • Debugging Spark performance issues
  • Scaling Spark pipelines for large datasets
  • Reducing shuffle and data skew

Core Concepts

1. Spark Execution Model

Driver Program    ↓Job (triggered by action)    ↓Stages (separated by shuffles)    ↓Tasks (one per partition)

2. Key Performance Factors

FactorImpactSolution
ShuffleNetwork I/O, disk I/OMinimize wide transformations
Data SkewUneven task durationSalting, broadcast joins
SerializationCPU overheadUse Kryo, columnar formats
MemoryGC pressure, spillsTune executor memory
PartitionsParallelismRight-size partitions

Quick Start

python
from pyspark.sql import SparkSessionfrom pyspark.sql import functions as F
# Create optimized Spark sessionspark = (SparkSession.builder    .appName("OptimizedJob")    .config("spark.sql.adaptive.enabled", "true")    .config("spark.sql.adaptive.coalescePartitions.enabled", "true")    .config("spark.sql.adaptive.skewJoin.enabled", "true")    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")    .config("spark.sql.shuffle.partitions", "200")    .getOrCreate())
# Read with optimized settingsdf = (spark.read    .format("parquet")    .option("mergeSchema", "false")    .load("s3://bucket/data/"))
# Efficient transformationsresult = (df    .filter(F.col("date") >= "2024-01-01")    .select("id", "amount", "category")    .groupBy("category")    .agg(F.sum("amount").alias("total")))
result.write.mode("overwrite").parquet("s3://bucket/output/")

Detailed patterns and worked examples

Detailed pattern documentation lives in references/details.md. Read that file when the navigation tier above is insufficient.

Best Practices

Do's

  • Enable AQE - Adaptive query execution handles many issues
  • Use Parquet/Delta - Columnar formats with compression
  • Broadcast small tables - Avoid shuffle for small joins
  • Monitor Spark UI - Check for skew, spills, GC
  • Right-size partitions - 128MB - 256MB per partition

Don'ts

  • Don't collect large data - Keep data distributed
  • Don't use UDFs unnecessarily - Use built-in functions
  • Don't over-cache - Memory is limited
  • Don't ignore data skew - It dominates job time
  • Don't use .count() for existence - Use .take(1) or .isEmpty()

來源與署名

來源:wshobson/agents位於plugins/data-engineering/skills/spark-optimization提交46891e7

授權條款: 無授權條款

內容歸原作者所有。SourceWeft 從公開儲存庫中收錄這些內容。

檢舉或申請下架