3 min read

    Lab 01: Advanced PySpark & Spark SQL

    #databricks#pyspark#lab

    Before we write code, let's understand our Goal: We need to read messy data from an S3 bucket, clean it up, flatten nested structures (like arrays), and perform complex math like running totals.

    The Tool: PySpark is the Python API for Apache Spark. It allows us to write Python code that is distributed across hundreds of computers at once.


    1. Explicit Schemas & Data Type Casting

    Real-World Analogy Mapping: Imagine hiring a librarian to organize 10,000 books.

    • The Long Way: You tell the librarian, "Look at every single page of every single book and guess what genre it is." This takes forever, and if they guess wrong, the whole library is ruined. This is what happens when you use inferSchema=True in PySpark.
    • The Smart Way: You hand the librarian a rulebook: "If it's on shelf A, it's a String. If it's on shelf B, it's an Integer." This is instant and 100% accurate. This is an Explicit Schema.
    python
    from pyspark.sql.types import StructType, StructField, StringType, IntegerType, TimestampType
    from pyspark.sql.functions import col
    
    # The Smart Way: Explicitly define the rulebook
    schema = StructType([
        StructField("order_id", StringType(), False),
        StructField("user_id", IntegerType(), True),
        StructField("order_timestamp", TimestampType(), True)
    ])
    
    # Read the data instantly using our schema
    df = spark.read.schema(schema).json("s3://data/orders/")
    
    # Casting: Changing a column type on the fly
    df_casted = df.withColumn("user_id", col("user_id").cast("string"))
    

    2. Advanced Functions: when/otherwise

    Goal: We need to create a new column based on conditions, exactly like an IF/ELSE statement in Python or SQL.

    python
    from pyspark.sql.functions import when, current_date, date_add, substring
    
    df_transformed = df.withColumn(
        "status_flag",
        when(col("status") == "COMPLETED", 1)
        .when(col("status") == "PENDING", 0)
        .otherwise(-1)
    )
    

    3. Nested JSON and Explode

    Follow the Data: Sometimes, a single row contains a list of items. For example, Order #1 contains an array of [Apple, Banana]. We cannot run SQL joins on arrays. We must explode the array. This acts like a bomb that blows the array apart, creating a brand new row for every item in the list, duplicating the Order ID.

    • Row 1: Order #1, Apple
    • Row 2: Order #1, Banana
    python
    from pyspark.sql.functions import explode, col
    
    # df contains a column 'items' which is a list.
    df_exploded = df.withColumn("item", explode(col("items")))
    

    4. Window Functions

    Problems Faced: When you use a normal groupBy(), it squishes all the rows together. If you group by user_id to find their total spend, you lose their individual order IDs!

    How Present Technology Solves It: Window Functions allow you to perform calculations (like running totals or rankings) without squishing the rows together. You open a "window" to look at related rows, but keep the original rows completely intact.

    python
    from pyspark.sql.window import Window
    from pyspark.sql.functions import sum, rank, desc
    
    # Define the "Window" (Look at all rows for a specific user, ordered by time)
    windowSpec = Window.partitionBy("user_id").orderBy(desc("order_timestamp"))
    
    # Calculate the rank of the order (1 = newest)
    df_windowed = df.withColumn("latest_order_rank", rank().over(windowSpec)) \
                    .withColumn("running_total", sum("amount").over(windowSpec))
    
    # Keep only the most recent order per user without losing the order_id!
    latest_orders = df_windowed.filter(col("latest_order_rank") == 1)
    

    5. All Join Types

    Goal: Merging two tables together based on a common key.

    • Inner Join: Keep only users that exist in BOTH tables.
    • Left Join: Keep ALL users from Table A, even if they don't exist in Table B.
    • Left Anti Join: Keep only users from Table A that DO NOT exist in Table B (Great for finding missing records).
    python
    # Inner Join (Default)
    df_inner = df1.join(df2, "user_id", "inner")
    
    # Left Join
    df_left = df1.join(df2, "user_id", "left")
    
    # Left Anti Join
    df_anti = df1.join(df2, "user_id", "left_anti")
    

    ← Previous: Databricks Implementation Labs: Syllabus | Next: Lab 02: Spark Internals & Catalyst Optimizer →**