Cheat sheetsSheet 7 of 10
Cheat sheet · Sheet 7 of 10
PySpark Cheat Sheet
A quick PySpark reference: reading and writing data, column expressions, joins, aggregations, window functions and the settings that matter for performance.
On this page
Setup
from pyspark.sql import SparkSession, Window
from pyspark.sql import functions as F
spark = SparkSession.builder.appName("job").getOrCreate()
Read and write
df = spark.read.parquet("/data/orders")
df = spark.read.option("header", True).csv("/data/raw.csv")
df.write.mode("overwrite").partitionBy("dt").parquet("/data/out")
Columns
df.select("id", F.col("amount") * 1.1)
df.withColumn("is_big", F.col("amount") > 100)
df.filter((F.col("country") == "IN") & F.col("amount").isNotNull())
df.withColumnRenamed("amt", "amount")
df.dropDuplicates(["order_id"])
Aggregate
df.groupBy("customer_id").agg(
F.count("*").alias("orders"),
F.sum("amount").alias("revenue"),
)
Join
orders.join(customers, "customer_id", "left")
orders.join(F.broadcast(countries), "country_code") # small side only
Window
w = Window.partitionBy("region").orderBy(F.desc("amount"), "rep")
df.withColumn("rn", F.row_number().over(w)).filter("rn = 1")
running = Window.partitionBy("rep").orderBy("dt").rowsBetween(Window.unboundedPreceding, Window.currentRow)
df.withColumn("running_total", F.sum("amount").over(running))
Inspect
df.printSchema()
df.explain() # look for Exchange (shuffle) and join type
df.rdd.getNumPartitions()
Performance settings (defaults checked on Spark 4.2)
| Setting | Default | Meaning |
|---|---|---|
spark.sql.shuffle.partitions |
200 | Partitions after a shuffle |
spark.sql.adaptive.enabled |
true | Adaptive Query Execution |
spark.sql.adaptive.skewJoin.enabled |
true | Split skewed join partitions |
spark.sql.autoBroadcastJoinThreshold |
10 MB | Max size to auto-broadcast |
Actions trigger jobs
count, collect, show, take, write. Everything else is lazy. Avoid collect() on large data.