The Ultimate PySpark Cheatsheet
A quick-reference guide for PySpark covering core DataFrame syntax, Window functions, batch processing, and Structured Streaming.
When you are deep in the trenches of a data pipeline, you don’t have time to dig through the official Apache Spark docs. This is a living cheatsheet of the most common PySpark patterns for both batch and streaming workloads.
[!TIP] Use
Cmd+ForCtrl+Fto quickly find the exact syntax you need.
1. Initializing the Spark Session
Every PySpark application starts here.
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("MyDataPipeline") \
.config("spark.sql.shuffle.partitions", "200") \
.getOrCreate()
2. Core DataFrame Operations (Batch)
The bread and butter of data transformations: Reading, Writing, Selecting, and Grouping.
Reading and Writing Data
# Reading Parquet
df = spark.read.parquet("s3://bucket/path/to/data/")
# Reading CSV with schema inference
df_csv = spark.read.option("header", "true").option("inferSchema", "true").csv("data.csv")
# Writing Partitioned Parquet
df.write.partitionBy("year", "month").mode("overwrite").parquet("s3://bucket/output/")
Select, Filter, and Mutate
import pyspark.sql.functions as F
# Select and Rename
df = df.select(
F.col("user_id").alias("id"),
F.col("event_type")
)
# Filter / Where
active_users = df.filter(F.col("status") == "active")
# Add or Modify a Column
df = df.withColumn("created_date", F.to_date(F.col("created_timestamp")))
GroupBy and Aggregations
# Grouping by single column
summary = df.groupBy("department").agg(
F.count("employee_id").alias("headcount"),
F.avg("salary").alias("avg_salary")
)
3. Window Functions
Essential for running totals, rankings, and time-series logic.
from pyspark.sql.window import Window
# Define the Window
w = Window.partitionBy("user_id").orderBy("event_timestamp")
# 1. Lead / Lag
df = df.withColumn("next_event", F.lead("event_type").over(w))
df = df.withColumn("prev_event", F.lag("event_type").over(w))
# 2. Ranking
df = df.withColumn("row_num", F.row_number().over(w))
df = df.withColumn("rank", F.rank().over(w))
# 3. Running Total
w_running = Window.partitionBy("user_id").orderBy("date").rowsBetween(Window.unboundedPreceding, Window.currentRow)
df = df.withColumn("running_total", F.sum("amount").over(w_running))
4. Structured Streaming
Spark's continuous processing engine for real-time pipelines.
Reading a Stream (Kafka Example)
stream_df = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "host1:port1,host2:port2") \
.option("subscribe", "topic_name") \
.option("startingOffsets", "earliest") \
.load()
Writing a Stream (Console / S3)
# Output to Console for Debugging
query = stream_df.writeStream \
.format("console") \
.outputMode("append") \
.start()
# Output to Cloud Storage with Checkpointing
query = stream_df.writeStream \
.format("parquet") \
.option("path", "s3://bucket/streaming/output/") \
.option("checkpointLocation", "s3://bucket/checkpoints/my_stream/") \
.trigger(processingTime="1 minute") \
.start()
query.awaitTermination()
5. Optimization & Performance Tuning
Never crash an executor again. Learn when to coalesce versus repartition.
| Operation | Best For | Behavior |
|---|---|---|
df.coalesce(n) | Decreasing partitions | Avoids full shuffle. Fast. Use before writing to disk. |
df.repartition(n) | Increasing partitions | Forces full shuffle. Distributes data evenly across cluster. |
df.cache() | Iterative algorithms | Stores dataframe in Memory (and Disk if Memory is full). |
# Coalesce before writing to avoid small file problem
df.coalesce(10).write.parquet("s3://bucket/output/")