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 从公开仓库中收录这些内容。

举报或申请下架