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), and offset (long).
  • You must deserialize the binary value column (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.

Related Topics