TT Lab
Get started
Learn Learning paths Courses

Apache Spark — The answer to a slow job is in the plan and the event log

Streaming is a series of small batches, and the checkpoint does the remembering

Continue in TT Lab

In one line

Structured Streaming is an engine that runs the same query repeatedly in small batches on an endlessly growing table; the checkpoint remembers how far it has processed, and the sink's record remembers what it wrote out. The watermark decides how long to wait for late data.

Why you should not rerun the batch every day

Suppose partner files drop into a landing folder dozens of times a day. If you solve it with a batch job, there are only two ways. Either you re-read the whole folder every time (it gets slower as data accumulates), or you manage the list of files already read yourself (if the job dies midway, the list and the result disagree). To do the second way properly, you end up hand-building a mechanism that "first writes down what it decided to process, and after writing it all, writes down that it is finished".

Structured Streaming is that mechanism built into the engine. The getting started documentation sees the incoming data as an input table to which rows keep being appended, and the query as a result table on top of it. At each trigger, it updates the result using only the new rows. You write the query exactly as for a batch, and the engine takes care of running it incrementally. According to the overview documentation, the default execution method is micro-batch, which handles this work as a succession of small batch jobs.

How it works: the file source

The file source in the API documentation reads files that newly appear in a directory. There are four rules to know.

The third rule creates duplicates in practice. If the partner resends yesterday's bundle with only the name changed, the engine accepts it as a new file. "Exactly once" at the file level is guaranteed, but not exactly once at the content level. That is solved later with deduplication.

Checkpoints and the sink's record

Diagram of the order in which one micro-batch runs. After picking new files from the landing folder, it first writes into the checkpoint's offsets the range this batch will process, writes the processed result as files to the output folder, writes the list of those files into _spark_metadata, and finally writes into commits that it is finished. If it dies midway, it reruns the batch that has no commits entry over the same range

The fault tolerance section of the getting started documentation summarizes the design in one sentence. Each source has an offset that represents the position read, and the engine records the range of offsets to process at each trigger with checkpoints and a write-ahead log. The sink is designed so that the result is the same even if it receives the same batch again (idempotently). A source that can be re-read and an idempotent sink meet to give end-to-end exactly-once.

If you open the checkpoint directory, you can see this design as files. offsets/0 is the range that batch 0 decided to process, and it is written before processing. commits/0 is the mark that the batch is finished, and it is written after processing. A restarted query finds the batches that are in offsets but not in commits and reruns them over the same range. The file sink has a matching record too. _spark_metadata/0 in the output folder lists the files that batch 0 wrote out, and when you read that folder with Spark, you see only the files on this list. This is why the half-files left by a dead batch do not get mixed into the result. The file sink being marked exactly-once in the table of the API documentation is also thanks to this record.

So deleting the checkpoint is deleting the memory. If you start the same query with a new checkpoint, it reprocesses every file in the folder from the beginning. The same documentation lists changing the output path of the file sink or the deduplication columns between restarts as not-allowed changes too.

Triggers: when to run a batch

If you do not give a trigger, it runs the next batch as soon as the previous one finishes. If you give an interval, it runs at each interval. For a use like this lab, which processes only what is there and then stops, availableNow fits. According to the API documentation, it processes all the data present at the time of execution and then stops by itself; depending on source options (maxFilesPerTrigger for a file source), it processes it split into several batches, and it processes first the batches that could not be committed in a previous run. The old once trigger is scheduled for deprecation.

q = (spark.readStream.schema(schema).json("/root/landing")
       .withWatermark("event_time", "10 minutes")
       .dropDuplicates(["event_id", "event_time"])
       .writeStream.format("parquet")
       .option("path", "/root/out")
       .option("checkpointLocation", "/root/chk")
       .trigger(availableNow=True)
       .start())
q.awaitTermination()

Deduplication and the watermark

In streaming, dropDuplicates means the same as in batch but has a different cost. According to the dropDuplicates documentation, in streaming it must hold every key it has already seen as state across triggers. Without a watermark, that state grows endlessly.

A watermark is the line saying "data arriving later than this is no longer waited for". The withWatermark documentation defines that line as the maximum event time seen so far minus the threshold. This line does two jobs. It tells which windows of a windowed aggregation are finalized, and it clears the state of the finalized windows. That is why a windowed aggregation in Append mode does not come out right away when the window ends. It comes out once, only after the watermark passes the end of the window.

Be sure to remember that the guarantee goes in only one direction. The API documentation guarantees that a 10-minute watermark never drops data that is late by less than 10 minutes, but it does not say that data later than that is necessarily dropped. It is usually dropped, but it may be aggregated. The watermark is not a filter for late data but a criterion for clearing state. And to clear state in an aggregation, you must put the watermark before the aggregation, on the same time column used in the aggregation, and the output mode must be Append or Update.

What it looks like in the field

First, putting the checkpoint in a temporary folder. With one reboot the memory disappears and the query reprocesses all the files. A checkpoint is data as precious as the output.

Second, writing in place in the landing folder. If the trigger runs while an upload tool is writing a file, it reads a half file. Have it write completely under a temporary name and then move it.

Third, a resend creates duplicates. The file source judges new files by path, so a bundle resent under a changed name comes in as is. Remove duplicates by a unique ID, and apply a watermark together so that the state does not grow endlessly.

What really matters in practice

What you will do in the next lab

You put bundles of device events into the landing folder, run the first batch with availableNow, and confirm that record 0 appears in the checkpoint's commits and in the output folder's _spark_metadata. You see that when you put in a new bundle only that one is processed, and that when you put in two bundles at once they are processed as one batch. You count the event_id values that came in as duplicates from a bundle resent with only its name changed, and then run a dropDuplicates stream with a new checkpoint to leave only the unique count. Finally, in a second landing folder copied with modification times preserved, you run a 10-minute watermark and a 10-minute window aggregation with one file per trigger, and see in the progress records how the watermark moves and how an event that arrived about 40 minutes late is dropped.