Install any skill in seconds. Free to start, no credit card required.
Get Started Free →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.
.claude/skills/dicklesworthstone-spark-optimization/SKILL.md| Test case | Without → With | Effect | Δ tokens | Δ turns |
|---|---|---|---|---|
| case-12 | ✗→✓ | ▲ Improved | 144% | 0% |
| case-13 | ✗→✓ | ▲ Improved | 139% | 0% |
| case-16 | ✗→✓ | ▲ Improved | 92% | 0% |
| case-10 | ✓→✓ | = Same ✓ | 140% | 0% |
| case-11 | ✓→✓ | = Same ✓ | 181% | 0% |
Production patterns for optimizing Apache Spark jobs including partitioning strategies, memory management, shuffle optimization, and performance tuning.
Driver Program
↓
Job (triggered by action)
↓
Stages (separated by shuffles)
↓
Tasks (one per partition)| Factor | Impact | Solution | | ----------------- | --------------------- | ----------------------------- | | Shuffle | Network I/O, disk I/O | Minimize wide transformations | | Data Skew | Uneven task duration | Salting, broadcast joins | | Serialization | CPU overhead | Use Kryo, columnar formats | | Memory | GC pressure, spills | Tune executor memory | | Partitions | Parallelism | Right-size partitions |
pythonfrom pyspark.sql import SparkSession from pyspark.sql import functions as F # Create optimized Spark session spark = (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 settings df = (spark.read .format("parquet") .option("mergeSchema", "false") .load("s3://bucket/data/")) # Efficient transformations result = (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/")
python# Calculate optimal partition count def calculate_partitions(data_size_gb: float, partition_size_mb: int = 128) -> int: """ Optimal partition size: 128MB - 256MB Too few: Under-utilization, memory pressure Too many: Task scheduling overhead """ return max(int(data_size_gb * 1024 / partition_size_mb), 1) # Repartition for even distribution df_repartitioned = df.repartition(200, "partition_key") # Coalesce to reduce partitions (no shuffle) df_coalesced = df.coalesce(100) # Partition pruning with predicate pushdown df = (spark.read.parquet("s3://bucket/data/") .filter(F.col("date") == "2024-01-01")) # Spark pushes this down # Write with partitioning for future queries (df.write .partitionBy("year", "month", "day") .mode("overwrite") .parquet("s3://bucket/partitioned_output/"))
pythonfrom pyspark.sql import functions as F from pyspark.sql.types import * # 1. Broadcast Join - Small table joins # Best when: One side < 10MB (configurable) small_df = spark.read.parquet("s3://bucket/small_table/") # < 10MB large_df = spark.read.parquet("s3://bucket/large_table/") # TBs # Explicit broadcast hint result = large_df.join( F.broadcast(small_df), on="key", how="left" ) # 2. Sort-Merge Join - Default for large tables # Requires shuffle, but handles any size result = large_df1.join(large_df2, on="key", how="inner") # 3. Bucket Join - Pre-sorted, no shuffle at join time # Write bucketed tables (df.write .bucketBy(200, "customer_id") .sortBy("customer_id") .mode("overwrite") .saveAsTable("bucketed_orders")) # Join bucketed tables (no shuffle!) orders = spark.table("bucketed_orders") customers = spark.table("bucketed_customers") # Same bucket count result = orders.join(customers, on="customer_id") # 4. Skew Join Handling # Enable AQE skew join optimization 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") # Manual salting for severe skew def salt_join(df_skewed, df_other, key_col, num_salts=10): """Add salt to distribute skewed keys""" # Add salt to skewed side df_salted = df_skewed.withColumn( "salt", (F.rand() * num_salts).cast("int") ).withColumn( "salted_key", F.concat(F.col(key_col), F.lit("_"), F.col("salt")) ) # Explode other side with all salts df_exploded = df_other.crossJoin( spark.range(num_salts).withColumnRenamed("id", "salt") ).withColumn( "salted_key", F.concat(F.col(key_col), F.lit("_"), F.col("salt")) ) # Join on salted key return df_salted.join(df_exploded, on="salted_key", how="inner")
pythonfrom pyspark import StorageLevel # Cache when reusing DataFrame multiple times df = spark.read.parquet("s3://bucket/data/") df_filtered = df.filter(F.col("status") == "active") # Cache in memory (MEMORY_AND_DISK is default) df_filtered.cache() # Or with specific storage level df_filtered.persist(StorageLevel.MEMORY_AND_DISK_SER) # Force materialization df_filtered.count() # Use in multiple actions agg1 = df_filtered.groupBy("category").count() agg2 = df_filtered.groupBy("region").sum("amount") # Unpersist when done df_filtered.unpersist() # Storage levels explained: # MEMORY_ONLY - Fast, but may not fit # MEMORY_AND_DISK - Spills to disk if needed (recommended) # MEMORY_ONLY_SER - Serialized, less memory, more CPU # DISK_ONLY - When memory is tight # OFF_HEAP - Tungsten off-heap memory # Checkpoint for complex lineage spark.sparkContext.setCheckpointDir("s3://bucket/checkpoints/") df_complex = (df .join(other_df, "key") .groupBy("category") .agg(F.sum("amount"))) df_complex.checkpoint() # Breaks lineage, materializes
python# Executor memory configuration # spark-submit --executor-memory 8g --executor-cores 4 # Memory breakdown (8GB executor): # - spark.memory.fraction = 0.6 (60% = 4.8GB for execution + storage) # - spark.memory.storageFraction = 0.5 (50% of 4.8GB = 2.4GB for cache) # - Remaining 2.4GB for execution (shuffles, joins, sorts) # - 40% = 3.2GB for user data structures and internal metadata spark = (SparkSession.builder .config("spark.executor.memory", "8g") .config("spark.executor.memoryOverhead", "2g") # For non-JVM memory .config("spark.memory.fraction", "0.6") .config("spark.memory.storageFraction", "0.5") .config("spark.sql.shuffle.partitions", "200") # For memory-intensive operations .config("spark.sql.autoBroadcastJoinThreshold", "50MB") # Prevent OOM on large shuffles .config("spark.sql.files.maxPartitionBytes", "128MB") .getOrCreate()) # Monitor memory usage def print_memory_usage(spark): """Print current memory usage""" sc = spark.sparkContext for executor in sc._jsc.sc().getExecutorMemoryStatus().keySet().toArray(): mem_status = sc._jsc.sc().getExecutorMemoryStatus().get(executor) total = mem_status._1() / (1024**3) free = mem_status._2() / (1024**3) print(f"{executor}: {total:.2f}GB total, {free:.2f}GB free")
python# Reduce shuffle data size spark.conf.set("spark.sql.shuffle.partitions", "auto") # With AQE spark.conf.set("spark.shuffle.compress", "true") spark.conf.set("spark.shuffle.spill.compress", "true") # Pre-aggregate before shuffle df_optimized = (df # Local aggregation first (combiner) .groupBy("key", "partition_col") .agg(F.sum("value").alias("partial_sum")) # Then global aggregation .groupBy("key") .agg(F.sum("partial_sum").alias("total"))) # Avoid shuffle with map-side operations # BAD: Shuffle for each distinct distinct_count = df.select("category").distinct().count() # GOOD: Approximate distinct (no shuffle) approx_count = df.select(F.approx_count_distinct("category")).collect()[0][0] # Use coalesce instead of repartition when reducing partitions df_reduced = df.coalesce(10) # No shuffle # Optimize shuffle with compression spark.conf.set("spark.io.compression.codec", "lz4") # Fast compression
python# Parquet optimizations (df.write .option("compression", "snappy") # Fast compression .option("parquet.block.size", 128 * 1024 * 1024) # 128MB row groups .parquet("s3://bucket/output/")) # Column pruning - only read needed columns df = (spark.read.parquet("s3://bucket/data/") .select("id", "amount", "date")) # Spark only reads these columns # Predicate pushdown - filter at storage level df = (spark.read.parquet("s3://bucket/partitioned/year=2024/") .filter(F.col("status") == "active")) # Pushed to Parquet reader # Delta Lake optimizations (df.write .format("delta") .option("optimizeWrite", "true") # Bin-packing .option("autoCompact", "true") # Compact small files .mode("overwrite") .save("s3://bucket/delta_table/")) # Z-ordering for multi-dimensional queries spark.sql(""" OPTIMIZE delta.`s3://bucket/delta_table/` ZORDER BY (customer_id, date) """)
python# Enable detailed metrics spark.conf.set("spark.sql.codegen.wholeStage", "true") spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true") # Explain query plan df.explain(mode="extended") # Modes: simple, extended, codegen, cost, formatted # Get physical plan statistics df.explain(mode="cost") # Monitor task metrics def analyze_stage_metrics(spark): """Analyze recent stage metrics""" status_tracker = spark.sparkContext.statusTracker() for stage_id in status_tracker.getActiveStageIds(): stage_info = status_tracker.getStageInfo(stage_id) print(f"Stage {stage_id}:") print(f" Tasks: {stage_info.numTasks}") print(f" Completed: {stage_info.numCompletedTasks}") print(f" Failed: {stage_info.numFailedTasks}") # Identify data skew def check_partition_skew(df): """Check for partition skew""" partition_counts = (df .withColumn("partition_id", F.spark_partition_id()) .groupBy("partition_id") .count() .orderBy(F.desc("count"))) partition_counts.show(20) stats = partition_counts.select( F.min("count").alias("min"), F.max("count").alias("max"), F.avg("count").alias("avg"), F.stddev("count").alias("stddev") ).collect()[0] skew_ratio = stats["max"] / stats["avg"] print(f"Skew ratio: {skew_ratio:.2f}x (>2x indicates skew)")
python# Production configuration template spark_configs = { # Adaptive Query Execution (AQE) "spark.sql.adaptive.enabled": "true", "spark.sql.adaptive.coalescePartitions.enabled": "true", "spark.sql.adaptive.skewJoin.enabled": "true", # Memory "spark.executor.memory": "8g", "spark.executor.memoryOverhead": "2g", "spark.memory.fraction": "0.6", "spark.memory.storageFraction": "0.5", # Parallelism "spark.sql.shuffle.partitions": "200", "spark.default.parallelism": "200", # Serialization "spark.serializer": "org.apache.spark.serializer.KryoSerializer", "spark.sql.execution.arrow.pyspark.enabled": "true", # Compression "spark.io.compression.codec": "lz4", "spark.shuffle.compress": "true", # Broadcast "spark.sql.autoBroadcastJoinThreshold": "50MB", # File handling "spark.sql.files.maxPartitionBytes": "128MB", "spark.sql.files.openCostInBytes": "4MB", }
.count() for existence - Use .take(1) or .isEmpty()| Case | Status | Duration (ms) | Turns | Tokens | Tool calls | ||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| Without | With | Δ | Without | With | Δ | Without | With | Δ | Without | With | Δ | ||
case-10 | pass→pass | 11,740 | 7,231 | -38% | 1 | 1 | 0% | 2,130 | 5,120 | +140% | 0 | 0 | — |
case-11 | pass→pass | 13,370 | 31,838 | +138% | 1 | 1 | 0% | 2,277 | 6,394 | +181% | 0 | 0 | — |
case-12 | fail→pass | 15,243 | 14,117 | -7% | 1 | 1 | 0% | 2,614 | 6,372 | +144% | 0 | 0 | — |
case-13 | fail→pass | 16,236 | 17,271 | +6% | 1 | 1 | 0% | 2,789 | 6,666 | +139% | 0 | 0 | — |
case-14 | pass→pass | 7,617 | 12,507 | +64% | 1 | 1 | 0% | 1,215 | 5,891 | +385% | 0 | 0 | — |
case-01 | fail→fail | 31,441 | 37,262 | +19% | 1 | 1 | 0% | 6,019 | 8,364 | +39% | 0 | 0 | — |
case-02 | fail→fail | 21,631 | 16,180 | -25% | 1 | 1 | 0% | 3,707 | 7,025 | +90% | 0 | 0 | — |
case-03 | pass→pass | 18,759 | 12,737 | -32% | 1 | 1 | 0% | 3,570 | 6,121 | +71% | 0 | 0 | — |
case-04 | pass→pass | 6,526 | 7,956 | +22% | 1 | 1 | 0% | 1,084 | 5,043 | +365% | 0 | 0 | — |
case-05 | pass→pass | 22,931 | 37,860 | +65% | 1 | 1 | 0% | 4,052 | 11,320 | +179% | 0 | 0 | — |
case-06 | pass→pass | 15,566 | 11,914 | -23% | 1 | 1 | 0% | 2,395 | 5,986 | +150% | 0 | 0 | — |
case-07 | pass→pass | 9,480 | 8,138 | -14% | 1 | 1 | 0% | 1,727 | 5,160 | +199% | 0 | 0 | — |
case-08 | pass→pass | 5,606 | 6,273 | +12% | 1 | 1 | 0% | 858 | 4,948 | +477% | 0 | 0 | — |
case-09 | pass→pass | 8,180 | 8,381 | +2% | 1 | 1 | 0% | 1,611 | 5,536 | +244% | 0 | 0 | — |
case-15 | pass→pass | 32,580 | 17,997 | -45% | 1 | 1 | 0% | 3,573 | 7,796 | +118% | 0 | 0 | — |
case-16 | fail→pass | 16,221 | 7,364 | -55% | 1 | 1 | 0% | 2,682 | 5,145 | +92% | 0 | 0 | — |
case-17 | pass→pass | 10,301 | 7,001 | -32% | 1 | 1 | 0% | 1,794 | 5,061 | +182% | 0 | 0 | — |
case-18 | pass→pass | 3,639 | 4,020 | +10% | 1 | 1 | 0% | 637 | 4,476 | +603% | 0 | 0 | — |
case-19 | pass→pass | 11,300 | 24,400 | +116% | 1 | 1 | 0% | 2,128 | 5,834 | +174% | 0 | 0 | — |
case-20 | pass→pass | 27,546 | 29,012 | +5% | 1 | 1 | 0% | 4,490 | 9,567 | +113% | 0 | 0 | — |
case-21 | pass→pass | 13,459 | 13,312 | -1% | 1 | 1 | 0% | 2,484 | 6,289 | +153% | 0 | 0 | — |
case-22 | pass→pass | 8,438 | 7,489 | -11% | 1 | 1 | 0% | 1,400 | 5,206 | +272% | 0 | 0 | — |
case-23 | pass→pass | 5,117 | 3,304 | -35% | 1 | 1 | 0% | 939 | 4,449 | +374% | 0 | 0 | — |
DecimalAI ran this skill against gemini-3.6-flash twice over the same eval suite — once with the skill loaded and once without — and compared the two runs case by case. 23 cases were attempted. The headline lift of +13 percentage points is the difference between those two pass rates over the 23 comparable cases.
Without the skill loaded, the model failed this case. With it loaded, the same prompt on the same model passed. This is one improved case from the latest verified run; every case, including any that regressed, is in the table above.
Other measured skills in the registry, with their headline benchmark lift.