๐ Advanced Spark Streaming
๐๏ธ Real-World Streaming Workflow
Follow the Data: A typical production streaming pipeline follows three distinct phases:
- Ingestion: Data flows in continuously from a source (e.g., Apache Kafka, IoT sensors, or TCP sockets).
- Transformations: Spark's
StreamingContextchops the data into micro-batches and applies transformations (likemap,filter, or Machine Learning models) in real-time. - Sinks: The processed insights are immediately pushed out to live dashboards, databases, or alert systems.
Uber uses this exact workflow for dynamic pricing. The ingestion layer receives live GPS pings from drivers and riders. The transformation layer continuously calculates demand and traffic. The sink instantly pushes price surges back to the app.
๐ช Windowed Operations
What Existed Previously:
Standard DStream operations (like map or reduceByKey) only look at the data that arrived within the current micro-batch (e.g., the last 5 seconds).
Problems Faced: What if you want to calculate a moving average? For example, "Show me the top trending hashtags over the last 60 minutes, updated every 10 seconds." A single 10-second micro-batch doesn't contain 60 minutes of history.
How Present Technology Solves It:
Spark Streaming introduces Windowed Operators (like reduceByKeyAndWindow). These operators look across multiple consecutive micro-batches by defining two parameters:
- Window Length: How far back in time to look (e.g., 60 minutes).
- Slide Interval: How often to calculate the result (e.g., every 10 seconds).
๐ป Practical: Sliding Window Word Count
The Smart Way: Let's compute a word count over a 10-second window, updating every 5 seconds.
from pyspark import SparkContext
from pyspark.streaming import StreamingContext
sc = SparkContext(appName="WindowedWordCount")
ssc = StreamingContext(sc, 5) # Batch interval = 5s
# CRITICAL: Stateful operations require a checkpoint directory to save state!
ssc.checkpoint("file:///tmp/ss_checkpoint")
lines = ssc.socketTextStream("localhost", 9999)
pairs = lines.flatMap(lambda l: l.split()).map(lambda w: (w, 1))
# Anatomy Breakdown of reduceByKeyAndWindow:
# 1st func: How to add new data entering the window (x + y)
# 2nd func: How to remove old data leaving the window (x - y)
# windowDuration: 10 seconds
# slideDuration: 5 seconds
window_counts = pairs.reduceByKeyAndWindow(
lambda x, y: x + y,
lambda x, y: x - y,
windowDuration=10,
slideDuration=5
)
window_counts.pprint()
ssc.start()
ssc.awaitTermination()
-------------------------------------------
Window [5s - 15s]
-------------------------------------------
('two', 3)
('three', 3)
๐ป Practical: Hashtag Top-5 Extraction
Let's extract trending hashtags from a continuous stream of social media posts.
from pyspark import SparkContext
from pyspark.streaming import StreamingContext
sc = SparkContext(appName="HashtagTop5")
ssc = StreamingContext(sc, 5)
lines = ssc.socketTextStream("localhost", 9999)
# 1. Extract hashtags and count them
hashtags = lines.flatMap(lambda l: [t for t in l.split() if t.startswith('#')])
counts = hashtags.map(lambda h: (h.lower(), 1)).reduceByKey(lambda a, b: a + b)
# 2. Sort and extract the Top 5
def top5(rdd):
return rdd.sortBy(lambda kv: kv[1], ascending=False).take(5)
# 3. Process each RDD in the DStream using foreachRDD
def print_top5(time, rdd):
items = top5(rdd)
print(f"==== Batch time: {time} ====")
for item in items:
print(item)
counts.foreachRDD(print_top5)
ssc.start()
ssc.awaitTermination()
๐งช Practice Drill
// Try answering these:
Q1. What is the purpose of the "Sink" phase in a streaming workflow?
Q2. If you set a `windowDuration` of 60 seconds and a `slideDuration` of 10 seconds, how often does Spark compute the result?
Q3. Why do stateful operations (like windowing) require you to set a checkpoint directory (`ssc.checkpoint`)?
๐ก Click for Solutions
A1. It is the final destination for processed streaming data, such as a live dashboard, a database, or an alerting system.
A2. Every 10 seconds (the slideDuration).
A3. Because Spark needs a durable place to save the accumulated "state" of the window across batches. If a node crashes, Spark uses the checkpoint to restore the running totals without losing data.
โ ๐ Foundations of Spark Streaming | Next Topic โ ๐ฌ Introduction to Apache Kafka