A Minimal Streaming Query

A Minimal Streaming Query Against a File Source

The listing plays a producer that lands BookNest order events as JSON Lines files in a directory, and a streaming query that keeps a running count per event type. Each run_once() processes whatever has arrived and stops, the way a scheduled job (Orchestration and Pipelines) would run it.

Counting BookNest order events from a directory of arriving filesJavaScript
import os, shutil
src, chk = "out/stream_in", "out/stream_chk"
for d in (src, chk):
    shutil.rmtree(d, ignore_errors=True)
os.makedirs(src)
with open("data/raw/order_events.jsonl") as f:
    events = [next(f) for _ in range(3000)]                 # BookNest's first 3,000 events
def drop_file(i):                                           # a producer lands one file
    with open(f"{src}/events-{i}.jsonl", "w") as out:
        out.writelines(events[i * 1000:(i + 1) * 1000])
schema = "event_id BIGINT, ts TIMESTAMP, type STRING, order_id BIGINT, customer_id INT"
counts = (spark.readStream.schema(schema).json(src)         # an unbounded input table
          .groupBy("type").count().orderBy("type"))
def run_once():                                             # process what arrived, then stop
    q = (counts.writeStream.outputMode("complete").format("console")
         .option("checkpointLocation", chk).trigger(availableNow=True).start())
    q.awaitTermination()
    print(f"batch {q.lastProgress['batchId']}: {q.lastProgress['numInputRows']} new rows")
drop_file(0); drop_file(1)
run_once()
drop_file(2)
run_once()                                                  # only the new file is read
Output
-------------------------------------------
Batch: 0
-------------------------------------------
+---------------+-----+
|           type|count|
+---------------+-----+
|order_cancelled|   31|
|     order_paid|  959|
|   order_placed| 1009|
|  order_shipped|    1|
+---------------+-----+
batch 0: 2000 new rows
-------------------------------------------
Batch: 1
...
|order_cancelled|   55|
|     order_paid| 1414|
|   order_placed| 1501|
|  order_shipped|   30|
...
batch 1: 1000 new rows

The second run read only the new file, 1,000 rows, yet printed totals over all 3,000: the checkpoint remembered the processed files and held the aggregation state. The file source needs a schema and treats files as immutable, so producers should write elsewhere and move each finished file in.