1a
Run the cell of Step 1. How many events are waiting?
Check your output
events waiting : 12000The number is lower when your flight data holds fewer than 20000 usable rows.
Block 2 · Exercise 6
Open the notebook 18-structured-streaming-starter in your workspace and run the setup cell at the top. It starts the Spark session, which takes a few minutes ; read Step 1 while it starts. Steps 2, 5 and 8 have a cell with None where your code goes, and the other steps run as they stand. Everything reads files, so no broker is needed.
Do one part at a time, 1a, 1b and so on, and compare with the expected output before you go on. The solution is in 18-structured-streaming-solution, under every part on the exercise site, and at the back of the exercises PDF.
The flights are a table, so there is nothing to stream yet. The cell of Step 1 turns a slice of them into JSON files with an event time.
Run the cell of Step 1. How many events are waiting?
events waiting : 12000The number is lower when your flight data holds fewer than 20000 usable rows.
Read the cell. Find the filter on crs_dep_time and the folder the events are written to. Why does the cell filter out 2400, and why is the folder events/part1 and not events? Write both answers in a comment.
crs_dep_time is an HHMM integer that the cell turns into a timestamp. What is the time 2400 to a timestamp parser?
Step 4 adds a second batch of events. Where will it go?
# 2400 is not a valid time of day, so the timestamp parser gives null and the row has no event time.
# part1 leaves room for part2 in Step 4 : the second batch of files lands next to the first
# in the same events folder, and the query has to pick up only the new ones.Build a streaming DataFrame events over the EVENTS folder, with SCHEMA as its schema. Add .option("recursiveFileLookup", "true"), because the events sit in a subfolder. Print events.isStreaming and events.printSchema().
The slide Sources shows the read. A batch read starts with spark.read ; a stream starts with the same word in the other form.
The slide Schemas says where the schema goes.
is streaming : True
root
|-- carrier: string (nullable = true)
|-- origin: string (nullable = true)
|-- dest: string (nullable = true)
|-- dep_delay: double (nullable = true)
|-- event_time: timestamp (nullable = true)events = spark.readStream.schema(SCHEMA).option("recursiveFileLookup", "true").json(EVENTS)
print("is streaming :", events.isStreaming)
events.printSchema()Read the folder as a stream again into a variable no_schema, without .schema(SCHEMA). Keep the recursive option. Read the error.
Use another variable name, so events stays as it is for the steps below.
An error that says the schema has to be specified. The exact wording depends on the Spark version.no_schema = spark.readStream.option("recursiveFileLookup", "true").json(EVENTS)Read the folder as a stream a third time into no_recursion, with SCHEMA and without the recursive option. Write it to a memory sink named no_recursion with .trigger(availableNow=True), wait for the query to end and count the rows of spark.table("no_recursion"). Does the query report an error?
The slide Sinks shows the memory sink. The sink is read back with spark.table and the query name.
0The query starts and ends without an error.
no_recursion = spark.readStream.schema(SCHEMA).json(EVENTS)
q = (no_recursion.writeStream.format("memory").queryName("no_recursion")
.trigger(availableNow=True).start())
q.awaitTermination()
print(spark.table("no_recursion").count())Compare the read of 2a with the batch read of Step 1. Which call is the only difference between a streaming and a batch DataFrame? Write the answer in a comment.
# spark.read against spark.readStream. The DataFrame API after that is the same.The cell of Step 3 keeps the flights delayed by more than fifteen minutes and writes them to parquet with .trigger(availableNow=True).
Run the cell. How many rows did the query read, and how many did it write?
rows read : 12000
rows written : ...Rows written is lower than rows read, and depends on the data.
Print q.isActive in a new cell. Did the query end by itself, or did you have to stop it?
The slide Triggers says what availableNow does once the waiting data is read.
Falseprint(q.isActive)
# False : availableNow reads what is waiting, writes it, and the query ends by itself.The cell of Step 4 writes the rest of the events as a second subfolder, part2, and runs the query of Step 3 again.
Run the cell. How many rows did the second run read, and how many rows are in the sink now?
rows before : ...
read this run: ...
rows after : ...The second run reads at most 8000 rows, far fewer than the 20000 events in both folders. Rows after is larger than rows before.
What told the query where to start this time? Write the answer in a comment.
The slides Checkpointing and Restart and Recovery show what is stored between two runs.
# The checkpoint location : it records which files the query has processed, so the second
# run read only the new ones. On a Kafka source it stores offsets in place of file names.Fill in the cell with windowed = None. Count the departures per hour of event time, per origin, with a two-hour watermark. Add the average dep_delay, rounded to one decimal, as avg_delay. Name the counted column departures.
The slide Windowing shows F.window("event_time", "1 hour") as a grouping key. Put it next to origin.
The watermark goes on events before the groupBy. The slide Watermarks shows the call.
windowed = (events
.withWatermark("event_time", "2 hours")
.groupBy(F.window("event_time", "1 hour"), "origin")
.agg(F.count("*").alias("departures"),
F.round(F.avg("dep_delay"), 1).alias("avg_delay")))Write windowed as parquet to win_out in append mode, with win_ckpt as the checkpoint location and .trigger(availableNow=True). Wait for the query to end. Read the result back into result and print its number of rows.
Define win_out and win_ckpt as folders inside WORK, as Step 3 does. The slide Output Modes names the mode.
window rows : ...The number depends on the data.
win_out, win_ckpt = f"{WORK}/windows", f"{WORK}/ckpt-windows"
q = (windowed.writeStream.format("parquet")
.option("path", win_out).option("checkpointLocation", win_ckpt)
.outputMode("append").trigger(availableNow=True).start())
q.awaitTermination()
result = spark.read.parquet(win_out)
print("window rows :", result.count())Show the eight busiest airport and hour combinations, with the window start and end, the origin, departures and avg_delay, busiest first. Which airport and hour is busiest?
window is a struct column. window.start and window.end pick its two parts.
+-------------------+-------------------+------+----------+---------+
|start |end |origin|departures|avg_delay|
+-------------------+-------------------+------+----------+---------+
|... |... |... |... |... |
+-------------------+-------------------+------+----------+---------+Eight rows, with departures falling from the top. The values depend on the data.
(result.select("window.start", "window.end", "origin", "departures", "avg_delay")
.orderBy(F.desc("departures")).show(8, truncate=False))Run the cell of Step 6. It runs the same aggregation on the batch DataFrame sample. Are the two row counts the same?
streaming rows : ...
batch rows : ...
they agree : FalseThe two counts are close, and the batch count is the larger.
Where did the missing streaming rows go? Add up the departures column of result and compare it with sample.count(). Write your explanation in a comment.
The slide Watermarks says when a window is appended.
Look at the latest event time in sample. Which windows can still receive events at that point?
rows accounted for by the streaming windows : ...
rows in the sample : 20000
rows still sitting in open windows : ...The last number is small next to 20000.
accounted = result.agg(F.sum("departures")).first()[0]
print("rows accounted for by the streaming windows :", accounted)
print("rows in the sample :", sample.count())
print("rows still sitting in open windows :", sample.count() - accounted)
# A window is appended only once the watermark has passed its end. The newest windows are
# still open when the query stops, so their rows are in the batch result and not yet in the
# streaming one. Append mode gives a correct result that lags behind the input.Run the cell of Step 7. It takes the watermark out of the aggregation and starts the query in append mode. What does the error say?
AnalysisException
Append output mode not supported when there are streaming aggregations on streaming
DataFrames/DataSets without watermark ...The rest of the message names the plan. The wording depends on the Spark version.
Why can Spark not append without a watermark? Write the answer in a comment.
An event with an old timestamp can always arrive. What does that do to a window that has closed?
# Without a watermark no window is ever final, because an event could always arrive for it.
# So there is no row that Spark can safely declare finished and append.Fill in the cell with enriched = None. Join events with lookup, the airports table with the city as city, so that every event has the city its flight left from.
The slide Joins covers a stream joined with a static table.
The join condition is F.col("origin") == F.col("iata_code").
enriched = events.join(lookup, F.col("origin") == F.col("iata_code"))Write event_time, carrier, origin, city and dep_delay of enriched as parquet to join_out, with join_ckpt as the checkpoint location and .trigger(availableNow=True). Wait for the query to end. Count the rows in join_out and compare that with the number of events that went in. Show five rows.
Count the events that went in with a batch read of EVENTS, using the schema and the recursive option.
events : 20000
joined : ...
lost to the inner join : ...
+-------------------+-------+------+----------+---------+
|event_time |carrier|origin|city |dep_delay|
+-------------------+-------+------+----------+---------+
...Five rows follow the header. joined is smaller than events and depends on the data.
join_out, join_ckpt = f"{WORK}/joined", f"{WORK}/ckpt-joined"
q = (enriched.select("event_time", "carrier", "origin", "city", "dep_delay")
.writeStream.format("parquet")
.option("path", join_out).option("checkpointLocation", join_ckpt)
.trigger(availableNow=True).start())
q.awaitTermination()
joined = spark.read.parquet(join_out)
events_in = spark.read.schema(SCHEMA).option("recursiveFileLookup", "true").json(EVENTS).count()
print("events :", events_in)
print("joined :", joined.count())
print("lost to the inner join :", events_in - joined.count())
joined.show(5, truncate=False)If rows were lost, which ones, and why? Write the answer in a comment.
The same thing happened in the joins of module 15. Which side of the join can miss a key?
# An inner join drops an event whose origin has no row in the airports table. The same
# reason as in module 15 : a code that the OurAirports extract does not list.Run enriched.explain() and look for the join in the plan.
== Physical Plan ==
...BroadcastHashJoin...The plan holds a BroadcastHashJoin : the small airports side is broadcast, as it would be in batch. The details depend on the Spark version.
Run the cell of Step 9. It lists spark.streams.active and stops whatever is left.
active queries : 0
active queries : 0Any query still running is listed between the two lines and stopped.
Click the session indicator and stop the session.
foreachBatch in place of a parquet sink, and print a summary for each batch. The batch is an ordinary DataFrame inside foreachBatch, so everything from modules 15 and 16 works on it. The solution notebook has cells for this and for the fifteen-minute window.