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