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.
}