PySpark Skill Sharpening Guide
Skill: PySpark
1. Setup & Essentials
Environment Setup
Local Setup:
pip install pysparkCloud Setup:
- Databricks Community Edition
- AWS EMR or Google Cloud Dataproc
- Jupyter Notebook with PySpark Kernel
Create a SparkSession
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("PySpark Practice") \
.getOrCreate()
Load Data
df = spark.read.csv("file.csv", header=True, inferSchema=True)
df.show()
2. Core DataFrame Operations (Parallel to Pandas)
Selection & Filtering
df.select("column1", "column2").show()
df.filter(df.column1 > 100).show()
df.where(df.column2.isNotNull()).show()
Sorting
df.orderBy("column1").show()
df.orderBy(df.column1.desc()).show()
Distinct & Duplicates
df.dropDuplicates(["column1"])
df.distinct()
Column Operations
from pyspark.sql.functions import col
df.withColumn("new_col", col("column1") * 2)
df.withColumnRenamed("old_col", "new_col")
df.drop("unwanted_column")
Type Casting
df.withColumn("col_int", col("col").cast("Integer"))
Aggregations & GroupBy
from pyspark.sql.functions import avg, count
df.groupBy("group_col").agg(avg("value"), count("*")).show()
Window Functions
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number
windowSpec = Window.partitionBy("group_col").orderBy("value")
df.withColumn("row_num", row_number().over(windowSpec)).show()
Joins
df1.join(df2, df1.id == df2.id, "inner").show()
3. Intermediate / Real-World Use Cases
Complex GroupBy + Aggregations
Example: Find the top-selling product per category
Deduplication Strategies
Using .dropDuplicates() or .row_number() with Window functions
Null Handling
df.na.fill(0)
df.na.drop(subset=["col1", "col2"])
UDFs & Pandas UDFs (for complex logic)
from pyspark.sql.functions import udf
from pyspark.sql.types import StringType
@udf(StringType())
def custom_logic(value):
return value.lower()
df.withColumn("new_col", custom_logic(col("col")))
4. Performance & Optimization
Partitioning & Repartitioning
df.repartition(4)
df.coalesce(2)
Caching & Persisting
df.cache()
df.persist()
Shuffle Understanding
Know when .groupBy() and .join() cause a shuffle.
Explain Plan
df.explain()
Broadcast Join
from pyspark.sql.functions import broadcast
df1.join(broadcast(df2), "id").show()
5. Pandas → PySpark Cheat Sheet
| Pandas | PySpark |
|---|---|
df[df.col > 5] |
df.filter(df.col > 5) |
df.sort_values('col') |
df.orderBy('col') |
df.groupby('col').agg() |
df.groupBy('col').agg() |
df['new'] = df.col * 2 |
df.withColumn('new', df.col * 2) |
df.dropna() |
df.na.drop() |
df.fillna(0) |
df.na.fill(0) |
df.merge(df2, on='id') |
df.join(df2, 'id') |
Next Steps & Practice
- Prepare practice datasets + exercises
- Review solutions and suggest optimizations
- Build a personalized PySpark Cheatbook
- Assign a mini-project simulating real-world data transformation
Optional: Shall I prepare your first practice set now?