4 min read

    🌊 Foundations of Spark Streaming

    #spark#streaming#bigdata

    πŸ—οΈ Batch vs Streaming Architectures

    What Existed Previously: Batch Processing systems (like Hadoop MapReduce) would wait for large amounts of data to accumulate over hours or days, and then process it all at once in massive chunks.

    Problems Faced: This caused critical latency. If an e-commerce site needs to detect a fraudulent transaction, waiting 24 hours for a batch job to run is uselessβ€”the money is already gone.

    How Present Technology Solves It: Streaming Processing continuously ingests and processes data the exact moment it arrives (e.g., from IoT sensors, social media feeds, or credit card transactions).

    The Big Win

    Streaming enables instant, real-time decision making for high-stakes use cases like fraud detection, dynamic pricing, and live recommendation engines.


    ⏱️ Spark's Micro-Batch Architecture

    Real-World Analogy Mapping: Imagine you are a worker at a very busy post office.

    • True Streaming is like a worker taking each individual letter exactly the millisecond it arrives and running to the back room to sort it. This is fast, but easily overwhelming.
    • Spark's Micro-Batching is like a worker waiting exactly 5 seconds, scooping up all the letters that arrived in that 5-second window, and sorting them together as a tiny batch.

    Spark Streaming does not process one record at a time. It collects continuous data streams over a tiny time interval (e.g., 1 second) and processes that small chunk as a Micro-Batch.


    🌊 Discretized Streams (DStreams)

    Anatomy Breakdown: A DStream (Discretized Stream) is the core architectural abstraction in classic Spark Streaming.

    Think of a DStream as a continuous, infinite conveyor belt. Spark chops this conveyor belt into small, fixed-time segments. Under the hood, each segment is just a standard RDD.

    Because a DStream is fundamentally just a sequence of RDDs, you can use all the classic RDD transformations you already know (like map, filter, and reduceByKey) directly on the stream!

    Analogy

    If a video is just a rapid sequence of still image frames, a DStream is just a rapid sequence of still RDDs.


    πŸ’» Practical: Socket Word Count

    Let's build a real-time word counter that listens to a network socket and updates the count every 5 seconds.

    The Smart Way:

    python
    from pyspark import SparkContext
    from pyspark.streaming import StreamingContext
    
    # 1. Initialize Contexts
    sc = SparkContext(appName="DStreamWordCount")
    
    # Create a StreamingContext with a 5-second Micro-batch interval
    ssc = StreamingContext(sc, 5) 
    
    # 2. Connect to a stream source (listening to port 9999)
    lines = ssc.socketTextStream("localhost", 9999)
    
    # 3. Apply standard RDD transformations!
    words = lines.flatMap(lambda line: line.split())
    pairs = words.map(lambda w: (w, 1))
    counts = pairs.reduceByKey(lambda a, b: a + b)
    
    # 4. Output the results to the console
    counts.pprint()
    
    # 5. Start the engine and wait for manual termination
    ssc.start()
    ssc.awaitTermination()
    
    Expected Output (Every 5 seconds)
    text
    -------------------------------------------
    Time: 2026-07-06 12:00:05
    -------------------------------------------
    ('spark', 2)
    ('streaming', 1)
    

    πŸ’» Practical: Socket Word Count (Structured Streaming)

    The modern, preferred way to do streaming in Spark is Structured Streaming, which treats the stream as a continuously appending DataFrame rather than raw RDDs.

    The Smart Way:

    python
    from pyspark.sql import SparkSession
    from pyspark.sql.functions import explode, split
    
    spark = SparkSession.builder.appName("StructuredWordCount").getOrCreate()
    
    # 1. Read from the socket stream
    lines = spark.readStream.format("socket").option("host", "localhost").option("port", 9999).load()
    
    # 2. Split lines into words
    words = lines.select(explode(split(lines.value, r'\s+')).alias("word"))
    
    # 3. Generate running word count
    counts = words.groupBy("word").count()
    
    # 4. Output to console with a 5-second trigger
    query = counts.writeStream \
        .outputMode("complete") \
        .format("console") \
        .trigger(processingTime="5 seconds") \
        .start()
    
    query.awaitTermination()
    

    πŸ§ͺ Practice Drill

    text
    // Try answering these:
    Q1. Why is batch processing unsuitable for credit card fraud detection?
    
    Q2. Does Spark Streaming process data one single record at a time?
    
    Q3. Under the hood, what is a DStream actually made of?
    
    πŸ’‘ Click for Solutions

    A1. Because batch processing analyzes data in large, delayed chunks. Fraud detection requires instant (streaming) processing to stop the transaction before it completes.

    A2. No. It uses a Micro-batch architecture, gathering arriving data into small time windows (e.g., 5 seconds) and processing them as a tiny batch.

    A3. A continuous sequence of RDDs.


    ← πŸ“š Kafka & Spark Streaming Syllabus | Next Topic β†’ 🌊 Advanced Spark Streaming