Structured Streaming is Spark's high-level streaming engine that treats a live data stream as an unbounded table. You write almost the same DataFrame/SQL transformations as batch; Spark runs them incrementally.
Model
Source (Kafka/files/socket)
-> continuous Input Table
-> query (filter/join/agg)
-> Result Table
-> Sink (Kafka/Parquet/foreachBatch)
with trigger (micro-batch or availableNow)Example
stream = (
spark.readStream
.format("kafka")
.option("subscribe", "orders")
.load()
)
parsed = parse_orders(stream) # your parsing
query = (
parsed.writeStream
.format("parquet")
.option("path", "/lake/bronze/orders")
.option("checkpointLocation", "/chk/orders")
.outputMode("append")
.start()
)
query.awaitTermination()Output modes
- append: new rows only (common for raw ingest)
- complete: full result rewritten (some aggs)
- update: changed rows (supported sinks)
Interview tip
Emphasize checkpointing for fault tolerance, exactly-once sinks when supported, and that micro-batch is the common execution mode.