2 min read

    Lab 03: Structured Streaming & Auto Loader

    #databricks#streaming#autoloader#lab

    Before we write code, let's understand our Goal: We want to continuously ingest files that are dropped into a cloud bucket without ever dropping a file or processing the same file twice, even if our computer crashes.

    The Tool: Auto Loader is a Databricks feature that automatically detects new files as they arrive in cloud storage. It is built on top of Spark Structured Streaming.


    1. Auto Loader Configuration

    What Existed Previously: Before Auto Loader, engineers had to write code that constantly scanned the entire S3 bucket to see if any new files were added. If the bucket had 10 million files, this scan took 30 minutes, causing massive delays.

    How Present Technology Solves It: Auto Loader uses "File Notification Mode." It subscribes to cloud events (like AWS SQS). Instead of searching the bucket, the bucket simply sends a text message saying, "Hey, a new file just arrived!"

    python
    df = spark.readStream.format("cloudFiles") \
        .option("cloudFiles.format", "json") \
        .option("cloudFiles.useNotifications", "true") # The Smart Way: File Notifications \
        .option("cloudFiles.schemaLocation", "s3://metadata/schemas/orders") \
        .option("cloudFiles.schemaEvolutionMode", "rescue") \
        .load("s3://raw-data/orders/")
    

    2. Handling Bad Data: Rescued Data

    Problems Faced: In traditional pipelines, if an API sends a string like "N/A" into an Integer column, the entire pipeline instantly crashes. The data team wakes up at 2 AM to fix it.

    How Present Technology Solves It: When schemaEvolutionMode is set to rescue, the pipeline will not crash. Instead, Databricks creates a special rescue column. It puts the bad "N/A" data into _rescued_data and leaves the rest of the valid columns intact. You can investigate the bad data the next morning.


    3. Checkpoint Recovery Strategy

    Goal: How does Spark remember what it has already processed?

    Real-World Analogy Mapping: Imagine reading a 10,000-page book. If you close the book and don't use a bookmark, you have to start from page 1 the next day. A Checkpoint is Spark's bookmark. It writes down exactly which files it has processed into a special folder.

    python
    df.writeStream.format("delta") \
        .outputMode("append") \
        .option("checkpointLocation", "s3://metadata/checkpoints/orders") # The Bookmark \
        .trigger(availableNow=True) \
        .table("bronze.orders")
    

    Warning: If someone accidentally deletes the checkpoint folder, Spark loses its bookmark. It will start from "Page 1" and re-process all historical files!


    4. Watermarking & Late Data

    Goal: Handling data that arrives late due to bad internet connections.

    Real-World Analogy Mapping: Imagine a teacher collecting homework. The deadline is Friday at 5 PM. A student tries to submit homework on Monday. Should the teacher accept it? A Watermark tells Spark exactly how long it should keep its "memory" open for late data before permanently closing the window and finalizing the math.

    python
    from pyspark.sql.functions import window
    
    # Allow data to arrive up to 2 hours late. After 2 hours, reject it.
    df_watermarked = df.withWatermark("event_time", "2 hours") \
        .groupBy(window("event_time", "10 minutes")) \
        .count()
    

    ← Previous: Lab 02: Spark Internals & Catalyst Optimizer | Next: Lab 04: Delta Lake Advanced Operations →**