๐ค Spark Streaming & Kafka Integration
๐ Connecting Spark to Kafka
What Existed Previously: In the early days of Spark Streaming, Kafka data was read using a Receiver-based approach. A dedicated Spark receiver would constantly pull data and store it in Spark's memory (or Write-Ahead Logs) before processing.
Problems Faced: This caused data duplication issues and performance bottlenecks. The receiver was a single point of failure and keeping Spark's Write-Ahead Logs perfectly synchronized with Kafka was highly complex.
How Present Technology Solves It: Spark introduced the Direct Stream approach. Spark no longer uses receivers. Instead, Spark periodically asks Kafka for the latest offsets and reads the data directly from Kafka partitions in parallel. This guarantees exactly-once processing with zero data duplication.
๐ป Practical: Direct Integration & Delta Lake
The Smart Way: We can use Spark Structured Streaming to seamlessly read JSON from Kafka and continuously append the results directly into a Delta Lake table with full ACID guarantees.
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("KafkaToDelta").getOrCreate()
# 1. Read directly from Kafka (Direct Stream Approach)
df = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "live-transactions") \
.load()
# 2. Extract the raw bytes and parse the JSON string
from pyspark.sql.types import StructType, StringType
from pyspark.sql.functions import from_json, col
schema = StructType().add("transaction_id", StringType()).add("amount", StringType())
text_df = df.selectExpr("CAST(value AS STRING) as json_payload")
parsed_df = text_df.select(from_json(col("json_payload"), schema).alias("data")).select("data.*")
# 3. Write continuously to Delta Lake with ACID guarantees
query = parsed_df.writeStream \
.format("delta") \
.outputMode("append") \
.option("checkpointLocation", "/mnt/checkpoints/transactions") \
.start("/mnt/delta/transactions")
query.awaitTermination()
Notice the checkpointLocation. Spark saves the Kafka offsets it has processed inside this checkpoint directory. If the cluster crashes, Spark reads this directory, sees exactly where it left off in Kafka, and resumes perfectly without processing the same message twice!
๐งช Practice Drill
// Try answering these:
Q1. Why is the Direct Stream approach better than the old Receiver-based approach for connecting Spark to Kafka?
Q2. When writing streaming data to a Delta Lake table, which `outputMode` ensures you do not accidentally overwrite your entire historical dataset?
Q3. How does Spark Streaming know exactly where it left off in a Kafka topic if the Spark cluster suddenly crashes?
๐ก Click for Solutions
A1. The Direct Stream approach eliminates the need for receivers and Write-Ahead logs, allowing Spark to read in parallel directly from Kafka partitions. This guarantees exactly-once semantics.
A2. append mode.
A3. By reading the saved offsets from the Checkpoint Directory (checkpointLocation).
โ ๐ค Kafka Producers & Advanced Features | Next Topic โ ๐ฅ Kafka Consumer API & Monitoring