Skip to main content
PYSPARK • LESSON 271

Structured Streaming Source with readStream

How can we initialize a continuous structured streaming DataFrame reading real-time Parquet files from cloud storage?

Production3 Minutes800 XP
🤔 THE QUESTION

How can we initialize a continuous structured streaming DataFrame reading real-time Parquet files from cloud storage?

💡 WHAT IS IT?

spark.readStream.schema(schema).parquet(path) sets up an unbounded DataFrame monitoring a directory for newly arrived files.

🎯 WHAT IS IT USED FOR?

Micro-batch event ingestion pipelines ingesting IoT telemetry and user clickstreams as files land in S3/GCS.

💻 EXAMPLE
streaming_df = spark.readStream.schema(event_schema).parquet("s3://lakehouse/streaming_landing/")

🎯 Mission Objectives

Practice typing production-grade PySpark code for Structured Streaming Source with readStream.

  • readStream initialization
  • Explicit streaming schema enforcement
  • Directory monitoring source