ML Engineer MasterClass (October) | 4 seats left

Spark · Streaming DataFrames
AmazonAmazon Analytics
Amazon · Structured Streaming

Streaming DataFrames

Transform a finite file stream and manage the lifecycle of a teaching query.

Step 1 of 6 · Learn

Read a stream with a schema

readStream describes arriving data rather than a completed batch. A file stream needs an explicit schema here. writeStream starts the query; processAllAvailable drains currently available files, not an infinite future source. This example writes a finite product snapshot, uses a uniquely named memory sink, and stops the query in finally. The timeout is a safety limit, not a processing-time guarantee.

Lesson reference: PySpark and Scala

Read a stream with a schema

readStream describes arriving data rather than a completed batch. A file stream needs an explicit schema here. writeStream starts the query; processAllAvailable drains currently available files, not an infinite future source. This example writes a finite product snapshot, uses a uniquely named memory sink, and stops the query in finally. The timeout is a safety limit, not a processing-time guarantee.

PySpark example

from tempfile import TemporaryDirectory
from uuid import uuid4
from threading import Timer
from pyspark.sql.functions import col, lit, array, element_at, to_timestamp, window, count, sum, upper, lower, date_format

seed = products
name = "lesson_" + uuid4().hex
query = None
timer = None
with TemporaryDirectory(prefix="spark-stream-") as directory:
    path = directory + "/input"
    try:
        seed.write.json(path)
        stream = spark.readStream.schema(seed.schema).json(path)
        transformed = stream.filter(col("category") == "Electronics")
        query = (transformed.writeStream.format("memory").queryName(name)
            .outputMode("append").option("checkpointLocation", directory + "/checkpoint").start())
        timer = Timer(45, query.stop)
        timer.daemon = True
        timer.start()
        query.processAllAvailable()

        if not query.isActive:
            raise RuntimeError("Streaming example stopped before completion; retry the run.")
        snapshot = spark.table(name).select(col("product_id"), col("product_name"))
        # Detach these few rows before stopping and removing the memory table.
        result = spark.createDataFrame(snapshot.collect(), snapshot.schema).orderBy("product_id", "product_name")
        result.show(truncate=False)
    finally:
        if timer is not None:
            timer.cancel()
        if query is not None:
            query.stop()
        spark.catalog.dropTempView(name)

Scala example

import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.StreamingQuery
import java.nio.file.Files
import org.apache.hadoop.fs.Path

val seed = products
val name = "lesson_" + java.util.UUID.randomUUID().toString.replace("-", "")
val directory = Files.createTempDirectory("spark-stream-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("input").toString
var query: StreamingQuery = null
val timer = new java.util.Timer(true)
val result = try {
    seed.write.json(path)
    val stream = spark.readStream.schema(seed.schema).json(path)
    val transformed = stream.filter(col("category") === "Electronics")
    query = transformed.writeStream.format("memory").queryName(name)
        .outputMode("append").option("checkpointLocation", directory.resolve("checkpoint").toString).start()
    timer.schedule(new java.util.TimerTask { def run(): Unit = query.stop() }, 45000L)
    query.processAllAvailable()

    require(query.isActive, "Streaming example stopped before completion; retry the run.")
    val snapshot = spark.table(name).select(col("product_id"), col("product_name"))
    val output = spark.createDataFrame(spark.sparkContext.parallelize(snapshot.collect().toSeq), snapshot.schema).orderBy("product_id", "product_name")
    output.show(false)
    output
} finally {
    timer.cancel()
    if (query != null) query.stop()
    spark.catalog.dropTempView(name)
    fs.delete(root, true) // Only this example’s input and checkpoint directory.
}

Apply stateless expressions

Many row-level DataFrame expressions also work on streaming inputs. A lowercase or uppercase expression needs no cross-batch aggregate state. Do not call collect directly on an unstarted streaming DataFrame: start a supported sink, then inspect its output. The memory sink is suitable for small teaching examples, not an unbounded production output store.

PySpark example

from tempfile import TemporaryDirectory
from uuid import uuid4
from threading import Timer
from pyspark.sql.functions import col, lit, array, element_at, to_timestamp, window, count, sum, upper, lower, date_format

seed = customers
name = "lesson_" + uuid4().hex
query = None
timer = None
with TemporaryDirectory(prefix="spark-stream-") as directory:
    path = directory + "/input"
    try:
        seed.write.json(path)
        stream = spark.readStream.schema(seed.schema).json(path)
        transformed = stream.select(col("customer_id"), upper(col("customer_name")).alias("label"))
        query = (transformed.writeStream.format("memory").queryName(name)
            .outputMode("append").option("checkpointLocation", directory + "/checkpoint").start())
        timer = Timer(45, query.stop)
        timer.daemon = True
        timer.start()
        query.processAllAvailable()

        if not query.isActive:
            raise RuntimeError("Streaming example stopped before completion; retry the run.")
        snapshot = spark.table(name).select(col("customer_id"), col("label"))
        # Detach these few rows before stopping and removing the memory table.
        result = spark.createDataFrame(snapshot.collect(), snapshot.schema).orderBy("customer_id", "label")
        result.show(truncate=False)
    finally:
        if timer is not None:
            timer.cancel()
        if query is not None:
            query.stop()
        spark.catalog.dropTempView(name)

Scala example

import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.StreamingQuery
import java.nio.file.Files
import org.apache.hadoop.fs.Path

val seed = customers
val name = "lesson_" + java.util.UUID.randomUUID().toString.replace("-", "")
val directory = Files.createTempDirectory("spark-stream-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("input").toString
var query: StreamingQuery = null
val timer = new java.util.Timer(true)
val result = try {
    seed.write.json(path)
    val stream = spark.readStream.schema(seed.schema).json(path)
    val transformed = stream.select(col("customer_id"), upper(col("customer_name")).alias("label"))
    query = transformed.writeStream.format("memory").queryName(name)
        .outputMode("append").option("checkpointLocation", directory.resolve("checkpoint").toString).start()
    timer.schedule(new java.util.TimerTask { def run(): Unit = query.stop() }, 45000L)
    query.processAllAvailable()

    require(query.isActive, "Streaming example stopped before completion; retry the run.")
    val snapshot = spark.table(name).select(col("customer_id"), col("label"))
    val output = spark.createDataFrame(spark.sparkContext.parallelize(snapshot.collect().toSeq), snapshot.schema).orderBy("customer_id", "label")
    output.show(false)
    output
} finally {
    timer.cancel()
    if (query != null) query.stop()
    spark.catalog.dropTempView(name)
    fs.delete(root, true) // Only this example’s input and checkpoint directory.
}

Join a stream to a static lookup

This inner join enriches incoming order_items with the static products DataFrame. It is not a stream-stream join and does not need to retain both sides as evolving streams. Lookup keys are unique, preventing row multiplication. Do not assume external lookup changes automatically refresh this already-created DataFrame.

PySpark example

from tempfile import TemporaryDirectory
from uuid import uuid4
from threading import Timer
from pyspark.sql.functions import col, lit, array, element_at, to_timestamp, window, count, sum, upper, lower, date_format

seed = order_items
name = "lesson_" + uuid4().hex
query = None
timer = None
with TemporaryDirectory(prefix="spark-stream-") as directory:
    path = directory + "/input"
    try:
        seed.write.json(path)
        stream = spark.readStream.schema(seed.schema).json(path)
        transformed = stream.filter(col("quantity") >= 2).join(products, "product_id").select("order_item_id", "product_name", "quantity")
        query = (transformed.writeStream.format("memory").queryName(name)
            .outputMode("append").option("checkpointLocation", directory + "/checkpoint").start())
        timer = Timer(45, query.stop)
        timer.daemon = True
        timer.start()
        query.processAllAvailable()

        if not query.isActive:
            raise RuntimeError("Streaming example stopped before completion; retry the run.")
        snapshot = spark.table(name).select(col("order_item_id"), col("product_name"), col("quantity"))
        # Detach these few rows before stopping and removing the memory table.
        result = spark.createDataFrame(snapshot.collect(), snapshot.schema).orderBy("order_item_id", "product_name", "quantity")
        result.show(truncate=False)
    finally:
        if timer is not None:
            timer.cancel()
        if query is not None:
            query.stop()
        spark.catalog.dropTempView(name)

Scala example

import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.StreamingQuery
import java.nio.file.Files
import org.apache.hadoop.fs.Path

val seed = order_items
val name = "lesson_" + java.util.UUID.randomUUID().toString.replace("-", "")
val directory = Files.createTempDirectory("spark-stream-")
val root = new Path(directory.toUri)
val fs = root.getFileSystem(spark.sparkContext.hadoopConfiguration)
val path = directory.resolve("input").toString
var query: StreamingQuery = null
val timer = new java.util.Timer(true)
val result = try {
    seed.write.json(path)
    val stream = spark.readStream.schema(seed.schema).json(path)
    val transformed = stream.filter(col("quantity") >= 2).join(products, Seq("product_id")).select("order_item_id", "product_name", "quantity")
    query = transformed.writeStream.format("memory").queryName(name)
        .outputMode("append").option("checkpointLocation", directory.resolve("checkpoint").toString).start()
    timer.schedule(new java.util.TimerTask { def run(): Unit = query.stop() }, 45000L)
    query.processAllAvailable()

    require(query.isActive, "Streaming example stopped before completion; retry the run.")
    val snapshot = spark.table(name).select(col("order_item_id"), col("product_name"), col("quantity"))
    val output = spark.createDataFrame(spark.sparkContext.parallelize(snapshot.collect().toSeq), snapshot.schema).orderBy("order_item_id", "product_name", "quantity")
    output.show(false)
    output
} finally {
    timer.cancel()
    if (query != null) query.stop()
    spark.catalog.dropTempView(name)
    fs.delete(root, true) // Only this example’s input and checkpoint directory.
}
example.pyPySpark
1. Finite input filesSeed a private directory and declare the schema.
2. Start and drainTransform the stream, start a memory sink, and drain available files.
3. Snapshot and stopDetach the small output; stop this query and remove its view and files.
Source products6 rows
product_idproduct_namecategory
201Wireless keyboardElectronics
202Laptop standOffice
203Desk lampHome
204USB-C hubElectronics
205Travel backpackTravel
206Notebook setOffice
Read a stream with a schema
Result2 rows
product_idproduct_name
201Wireless keyboard
204USB-C hub
Two Electronics rows reach the example’s memory sink. The final result is a detached snapshot, not the live stream itself.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.