advanced
Kafka
10 min readLast updated: 2026-07-08
Overview
Connect Spark to Apache Kafka to read and write real-time message streams.
What You Will Learn
In this lesson, you will learn:
- Kafka Reader: Reading batch or streaming data from Kafka topics.
- Schema Deserialization: Converting binary keys and values into structured columns.
- Kafka Writer: Writing data back to Kafka topics.
Detailed Concept Explanation
Spark can connect to Apache Kafka for both batch queries and real-time structured streaming.
When reading from Kafka:
- Spark retrieves records with a default schema containing columns like
key(binary),value(binary),topic(string),partition(int), andoffset(long). - You must deserialize the binary
valuecolumn (typically cast to string or parsed using Avro/JSON schemas) to extract structured fields.
Code Examples
Python (PySpark) Implementation
python
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("KafkaTest") \
.config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.0") \
.getOrCreate()
# Read batch data from a Kafka topic
kafka_df = spark.read.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "my-topic") \
.load()
# Convert binary message value to string
messages_df = kafka_df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
messages_df.show()
Expected Output
text
+----+--------------------+
| key| value|
+----+--------------------+
| 101|{"action": "signup"}|
+----+--------------------+
Execution Plan Diagram (Python & Scala)
Execution Plan Diagram
SparkSession.builder
read.format(kafka)
selectExpr(CAST value AS STRING)
show()
Scala Implementation
scala
import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder().appName("KafkaScala").getOrCreate()
val kafkaDF = spark.read.format("kafka")
.option("kafka.bootstrap.servers", "localhost:9092")
.option("subscribe", "my-topic")
.load()
val messagesDF = kafkaDF.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
messagesDF.show()
Expected Output
text
+----+--------------------+
| key| value|
+----+--------------------+
| 101|{"action": "signup"}|
+----+--------------------+
Common Mistakes
- Forgetting Package Coordinates: Omitting the SQL-Kafka connector jar dependency coordinates (
spark-sql-kafka-0-10) at session startup. This causes Spark to fail to compile the Kafka format.