TT Lab
Get started
Learn Learning paths Courses

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

Process each arriving file exactly once and drop late data

Continue in TT Lab

Goal

Each time a file arrives in the landing folder, you process it with Structured Streaming and write it to a Parquet sink, and confirm how the checkpoint and the sink log remember "how far it has processed". You see that resent files with only a different name create duplicates, prevent them with dropDuplicates, and then see in a window aggregation how the watermark drops late data.

Why it matters

The difficulty of streaming is not computation but memory. When a job dies and comes back, it has to know what it has already done and what it has not. Structured Streaming writes this in the checkpoint folder. Before starting a batch, it writes the range to read into offsets, and after writing the results fully to the sink, it writes into commits. The file sink leaves in its own folder's _spark_metadata the list of files written for each batch, and the reader sees as the result only the files on that list. That memory has limits. The file source remembers whether something was processed by file name, so if the partner resends the same content under another name, it processes it as new data. To remove duplicates by content, you have to hold the keys you have seen as state. Aggregating windows by event time creates another problem: when to close a window. The watermark is "the latest time seen so far − the allowed delay"; data older than that is dropped as late, and only windows whose end has passed the watermark are emitted as results (append mode). The watermark moves only between batches, so how many batches the data arrives in changes the result.

Steps

  1. Create the landing folder /root/spk/stream/in and copy /data/stream/batch-01.jsonl into it.
  2. With /root/spk/stream/stream.py (application spk-stream-run), run once a stream that reads the landing folder and writes Parquet to /root/spk/stream/out/events, with the checkpoint /root/spk/stream/ckpt/events and trigger(availableNow=True).
  3. Add batch-02.jsonl to the landing folder and run the same stream again. Only the new file should become a new batch.
  4. Add batch-03.jsonl and batch-04.jsonl together and run again. The two files should be processed as one batch.
  5. Add the resent file batch-06.jsonl (its content is the same as batch-03) and run again, then write the number of event_id values that went into the sink twice to /root/spk/stream/out/dups.txt as an integer.
  6. With /root/spk/stream/dedup.py (application spk-stream-dedup), run a new stream (checkpoint /root/spk/stream/ckpt/dedup) that applies dropDuplicates(["event_id"]) to the landing folder and writes to /root/spk/stream/out/dedup.
  7. Copy batch-01 to batch-05 with cp -p into the second landing folder /root/spk/stream/late_in, and with /root/spk/stream/window.py (application spk-stream-window), write a count aggregation over 10-minute windows with maxFilesPerTrigger=1 and a 10-minute watermark to /root/spk/stream/out/windows in append mode (checkpoint /root/spk/stream/ckpt/windows). Write the progress records to /root/spk/stream/out/progress.json and the late-data analysis to /root/spk/stream/out/late.json.
  8. In /root/spk/stream/report.md, write three sections: ## 한 번씩만 처리하기, ## 재전송과 중복 and ## 늦은 자료 (use exactly these Korean headings in this order; they mean "Processing exactly once", "Resends and duplicates" and "Late data"). Put the number of duplicates from step 5 in the second section and the number of late events from step 7 in the third.

Notes

Put the first bundle in the landing folder

Create the landing folder /root/spk/stream/in and copy /data/stream/batch-01.jsonl into it (/root/spk/stream/in/batch-01.jsonl).

The streaming file source watches the folder and picks up newly appearing files for the next batch. A file must appear complete, all at once: if it picks up a file that is being written, a half is processed (that is why people usually write elsewhere and then move).

The first batch: checkpoint and sink log

Create /root/spk/stream/stream.py with the application name spk-stream-run, read /root/spk/stream/in with the schema (event_id STRING, device STRING, event_time TIMESTAMP, value INT), start a stream that writes Parquet to /root/spk/stream/out/events with checkpointLocation=/root/spk/stream/ckpt/events and trigger(availableNow=True), and wait until it finishes.

availableNow means "process what is there now and stop". When it finishes, the checkpoint has offsets/0 and commits/0, and the sink folder's _spark_metadata/0 has the list of files this batch wrote. The grader reads the files on that list and checks whether they match the event_id values of batch-01.

Only the new file becomes the next batch

Add /data/stream/batch-02.jsonl to the landing folder and run the step 2 stream again as it is. Only the events of batch-02 should go into the new batch.

The checkpoint remembers that batch-01 was already processed, so it does not read it again. If you delete the checkpoint, the memory disappears too and everything is processed from the start, and one more copy of the same data piles up in the sink. The grader looks in the sink log for a batch containing only the ids of batch-02.

Two files in one batch

Add batch-03.jsonl and batch-04.jsonl together to the landing folder and run the stream again. The events of the two files should be processed as one batch.

How many files one batch picks up is decided by maxFilesPerTrigger, and if you do not give it, it picks up all the new files present at that time. Batch boundaries affect delay and the watermark, so you meet them again in step 7.

A resend with only a different name creates duplicates

Add the resent file /data/stream/batch-06.jsonl (its content is the same as batch-03) to the landing folder and run the stream again. Then read the committed files of the sink (the _spark_metadata list of each batch) and write the number of event_id values that appear two or more times to /root/spk/stream/out/dups.txt as an integer.

The file source's memory is the file name. The name is new, so it processes it as new data. You can also just read the sink folder (Spark reads it by looking at _spark_metadata), but remember that a tool that directly scans the part files of the directory may also pick up the debris of a failed batch.

dropDuplicates: remember the ids you have seen

Create /root/spk/stream/dedup.py with the application name spk-stream-dedup, read the same landing folder, and run a new stream that writes the result of dropDuplicates(["event_id"]) to /root/spk/stream/out/dedup, with checkpointLocation=/root/spk/stream/ckpt/dedup and availableNow. Every event_id in the result should appear exactly once.

Deduplication is a stateful operation. It holds the seen ids in a state store and drops an id when it comes again. Without a watermark the state grows endlessly, so in production you use it together with withWatermark or use dropDuplicatesWithinWatermark. A checkpoint is separate for each stream.

Watermark: drop late data and emit only closed windows

Create /root/spk/stream/late_in, copy batch-01 to batch-05 into it with cp -p, create /root/spk/stream/window.py with the application name spk-stream-window, read with maxFilesPerTrigger=1, apply withWatermark("event_time", "10 minutes"), count over 10-minute windows, and write them with the columns start,end,count to /root/spk/stream/out/windows in append mode (checkpointLocation=/root/spk/stream/ckpt/windows, availableNow). After it finishes, from query.recentProgress, write as a list to /root/spk/stream/out/progress.json the batch, rows, watermark and dropped (the sum of the state operator's numRowsDroppedByWatermark) for each batch, and write the number of batch-05 events earlier than the watermark at the time batch-05 came in to /root/spk/stream/out/late.json as {"watermark": "yyyy-MM-dd HH:mm:ss", "late_events": 정수, "dropped_rows_metric": 정수} (the last two values are integers).

The watermark rises at the end of a batch to "the latest time seen − 10 minutes" and is used from the next batch. That is why you have to feed it one file per batch so that when batch-05 comes in, the watermark has already reached around 09:29. Also notice that there are a few dozen late events but the dropped metric is 1: this is because partial aggregation before the state operator has already collapsed them into one row per window.

Record memory, duplicates and lateness in numbers

In /root/spk/stream/report.md, write three sections: ## 한 번씩만 처리하기, ## 재전송과 중복 and ## 늦은 자료 (use exactly these Korean headings in this order; they mean "Processing exactly once", "Resends and duplicates" and "Late data"). Put the number of duplicates from step 5 in the second section and the late_events and dropped_rows_metric from step 7 in the third section.

In the first section, write what the checkpoint's offsets and commits and the sink log each remember; in the second, why a file with only its name changed becomes a duplicate; and in the third, why the number of late events and the metric differ.