# The Structured Streaming Model — Apache Spark: Big Data Processing with DataFrames

Source: https://www.geekswithgeeks.com/en/spark/st-model

> Write streaming queries with the same DataFrame API and understand triggers and output modes.

## Same API, unbounded input

Structured Streaming lets you write a streaming job as if it were a batch job on a table that **new rows keep being appended to**. You read with `spark.readStream` from a source (Kafka, files, a test `rate` source), apply ordinary DataFrame operations and write with `writeStream` to a **sink** (Kafka, files/Delta, a console or memory table for testing). A **trigger** controls how often micro-batches run (`processingTime="10 seconds"`, or `availableNow=True` to process what is there and stop). The **output mode** says what is written each time: `append` (only new final rows), `update` (rows that changed) or `complete` (the whole result table, for small aggregations). A **checkpoint location** stores progress and state so a restarted query resumes without losing or duplicating data (with replayable sources and idempotent or transactional sinks).

## A table that keeps growing

Structured Streaming treats a live stream as an unbounded table and runs your query incrementally.

![Three parts: source, query, sink.](assets/figures/spark/section-6-map.svg) — Figure 6.1 — Source, query and sink.

## A tiny streaming aggregation (illustrative)

Uses the built-in `rate` source (a test generator) and an in-memory sink. I did not capture this run, so no output is shown; counts keep growing as batches arrive. For real jobs you would write to Kafka or Delta with a `checkpointLocation`.

```python
import time
q = (spark.readStream.format("rate").option("rowsPerSecond", 20).load()
     .withColumn("bucket", F.col("value") % 3)
     .groupBy("bucket").count()
     .writeStream.outputMode("complete").format("memory").queryName("rates")
     .trigger(processingTime="1 second").start())

time.sleep(5)
spark.sql("SELECT bucket, count FROM rates ORDER BY bucket").show()
q.stop()
```

## Always set a checkpoint in production

Without a checkpoint a restarted query starts over (or fails to resume correctly). Give each query its own durable checkpoint folder, and do not delete or share it casually.

**Quiz:** What does a checkpoint location store?

- [ ] Only the final output files
- [ ] Your passwords
- [ ] The Spark binaries
- [x] Query progress and state so a restart can resume

*Answer:* Query progress and state so a restart can resume. Offsets and state in the checkpoint let a restarted query continue where it stopped.
