PySpark Skill Sharpening Guide

Skill: PySpark

1. Setup & Essentials

Environment Setup

  • Local Setup:
    pip install pyspark

  • Cloud 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?