Yesterday's Number Changed Today: Watermarks and Late Data
Goal
You build a tool wm.py that aggregates by event time. You measure the gap between processing time and event time as a distribution, move the watermark in arrival order, separate the data that came after a window closed, compare the results of a drop policy and an apply policy, and leave the value at the moment of closing and the corrections afterward separately.
Why it matters
We dealt with the watermark once in the earlier lab of this course. It was writing the maximum timestamp of the loaded data in a table so that the next run reads only what comes after, and the question that watermark answers was just one: how far have I read. Here we deal with what comes next. A mobile app loses signal and sends events a few minutes later, and a device that was in airplane mode sends events a few hours later. The time at which those events occurred is a time that has already passed, and we have already published the aggregate for that time range. The watermark is a promise that you consider that no more data will come, not a fact. So there are three things to decide. What to set the allowed lateness to, whether to drop the data that came in breaking the promise or apply it as a late update, and, if you apply it, how to explain the difference between the number you published yesterday and the number you publish today. This is not a lab about measuring time. It is a lab that computes with two fields written in the event, event_time and ingest_time, and how many seconds the program actually takes has no bearing. The grader does not trust the text you write down. It sets up an event stream that the grader created in a temporary file, actually runs your tool, and compares the answers while changing the window size and the allowed lateness. The counts and the lag distribution change on every run.
Steps
- Create and run /root/wmark/gen_stream.py to create /root/wmark/work/stream.jsonl.
- In /root/wmark/wm.py, create
skewso that it produces the lag distribution and the number of out-of-order events. - Add
watermarkto move the watermark in arrival order. - Add
windowsto divide and aggregate by event time. - Add
lateto separate the data that came after a window closed. - Add
aggto produce the results of the drop policy and the apply policy. - Add
closeto produce the value at the moment of closing and the corrections afterward separately. - Write /root/wmark/work/watermark_report.json and /root/wmark/work/watermark_report.md.
Reference
- Do all the work under
/root/wmark. The data is/root/wmark/work/stream.jsonl. - One data line is one JSON object and holds
id,key,event_time,ingest_time, andamount. The two times are integer epoch seconds. There may be other fields. - The order written in the file is the arrival order. If you sort again by event time, every judgment in this lab collapses.
- Execution contract:
python3 /root/wmark/wm.py <명령> <파일> [--size=초] [--lateness=초] [--policy=drop|update](the placeholders stand for the command, the file, and seconds). Give the answer as a single JSON object on standard output. On success the exit code is 0, if the file is missing it is 3, and if the usage or the policy name is wrong it is 2. - Windows are fixed-size and do not overlap. An event's window start is
event_time - (event_time % size), and the window's range runs from the start up to just before the start plus the size. The JSON key is the window start written as a string. - Lag is
ingest_time - event_time. Percentiles use nearest rank — put the values in ascending order and pick theceil(건수 * p / 100)-th one (from 1), where the placeholder stands for the count. Do not interpolate. skew <파일>response:{"events": 정수, "out_of_order": 정수, "min_lag": 정수, "max_lag": 정수, "p50_lag": 정수, "p95_lag": 정수}(the placeholders stand for the file and integers). out_of_order is the number of events whose event time is earlier than the maximum event time of those that arrived before it.watermark <파일> --lateness=<초>response:{"lateness": 정수, "max_event_time": 정수, "advances": 정수, "final_watermark": 정수}(integers). advances is the number of arrivals that newly updated the maximum event time.windows <파일> --size=<초>response:{"size": 정수, "count": 정수, "windows": {"창시작": {"events": 정수, "amount": 정수}}}(the placeholder in the key stands for the window start). It counts everything without regard to late or early.late <파일> --size=<초> --lateness=<초>response:{"size": 정수, "lateness": 정수, "on_time": 정수, "late": 정수, "late_by_window": {"창시작": 정수}}. Which event was late is judged with the watermark computed only from those that arrived before it. If the watermark is at or beyond the end of that window, it is late. The first arrival has no earlier data to compare with, so it is not late.agg <파일> --size --lateness --policy=drop|updateresponse:{"policy": 문자열, "size": 정수, "lateness": 정수, "windows": {...}, "dropped": 정수, "restated": [창시작 문자열 오름차순]}(a string, integers, and the window starts as strings in ascending order). With drop, it leaves out the late data and counts it as dropped, and restated is an empty list. With update, it includes the late data too and dropped is 0, and the windows that late data went into are put in restated.close <파일> --size --latenessresponse:{"size": 정수, "lateness": 정수, "sealed": {...}, "corrections": [{"window": 창시작, "delta_events": 정수, "delta_amount": 정수}], "final": {...}}. sealed holds only what arrived before closing (a window with only late data is left at 0), corrections is in ascending window-start order, and final is the sum of the two.- The report JSON holds size, lateness, events, windows, on_time, late, dropped, restated, max_lag, p95_lag, and policy. Choose lateness from 1 up to less than max_lag, and at that value at least 1 late-arriving event must appear. If none appears, reduce the allowed lateness.
- The section headings of the report MD are
## 무엇을 재었나## 허용 지연을 얼마로 잡았나## 늦게 온 자료를 어떻게 했나## 어제 숫자가 바뀐 이유(in order: what was measured, what allowed lateness was chosen, what was done with the late data, and why yesterday's numbers changed). - Official documentation: Flink Generating Watermarks · Flink Windows · python json
- Common mistakes: sorting the data again by event time (late data comes out as 0 rows), letting the watermark go backward, not counting the dropped rows, and sorting window starts as strings.
Get the late-arriving events in hand
Create and run /root/wmark/gen_stream.py to create /root/wmark/work/stream.jsonl. There must be at least 60 lines, the file order must be ascending ingest_time, every line must have ingest_time greater than or equal to event_time, there must be at least 5 lines with a lag of 120 seconds or more, the span of event time must be at least 1200 seconds, and there must be at least 3 kinds of key.
If you make the lag from only one kind of distribution, there is nothing to look at later. Scatter it into three branches: most come in within a few seconds, some a few minutes later, and a very few flood in tens of minutes later. After you make it all, sort by ingest_time and write it to the file, and that becomes the arrival order. You must fix the seed so that the data does not wobble while you compare changing the allowed lateness.
First measure how much the gap is
In /root/wmark/wm.py, create skew <파일> (the placeholder stands for the file) so that it produces, as JSON, the count, the number of out-of-order events, and the minimum, maximum, median, and 95th percentile of the lag.
Lag is ingest_time - event_time. Use the nearest rank for percentiles and do not interpolate — they must come out as integers so that neither grading nor the meeting wobbles. For the out-of-order count, scan in arrival order and count those earlier than the maximum event time so far.
Move the watermark in arrival order
Add watermark <파일> --lateness=<초> (the placeholders stand for the file and seconds) to produce lateness, max_event_time, advances, and final_watermark. advances is the number of arrivals that newly updated the maximum event time.
The watermark does not go backward. If you let a late-arriving event lower the maximum event time, a window that was already closed repeatedly reopens and closes and no value is finalized. The first arrival has no earlier data to compare with, so it is itself one advance.
Divide windows by event time
Add windows <파일> --size=<초> (the placeholders stand for the file and seconds) to produce the count and amount per window. The window start is event_time - (event_time % size), and in this step you count everything without regard to late or early.
JSON keys must be strings, so write the window start as a string. When you sort later, compare as integers, not strings — if the digit counts differ, string sorting gives a wrong order.
Pick out the data that came after a window closed
Add late <파일> --size=<초> --lateness=<초> (the placeholders stand for the file and seconds) to produce on_time, late, and late_by_window. Which event was late is judged with the watermark computed only from those that arrived before it, and if the watermark is at or beyond the end of that window, it is late.
Lateness is not a property of the data but a property of the arrival order. If you sort once by event time to make it convenient to handle, the arrival order disappears and late data comes out as 0 rows — it is not that there is no problem; you have removed the eyes to see it. Scan in the order written in the file, but update the maximum event time after you finish the judgment. The first arrival has no earlier data to compare with, so it is not late.
Drop it or apply it
Add agg <파일> --size --lateness --policy=drop|update (the placeholder stands for the file). With drop, it leaves out the late data and counts it as dropped, and restated is an empty list; with update, it includes the late data too and dropped is 0, and the windows that late data went into are put in restated.
Even if you choose to drop, be sure to count the dropped rows. If you drop without counting, later there is no basis to explain why the totals do not match. Put restated in ascending order of window start, but compare as integers. If an unknown policy name comes, end with exit code 2.
Leave the value at closing and the corrections separately
Add close <파일> --size --lateness (the placeholder stands for the file) to produce sealed (only what arrived before closing), corrections (a list of corrections in ascending window-start order), and final (the sum of the two). The sealed of a window with only late data is 0.
If you produce only the two merged into one, there is no way to explain why yesterday's number changed today. If you leave them separately, the amount that changed is itself the answer. Put in corrections only the windows that actually have a correction, and put all windows in final.
A single page that explains the changed numbers
In /root/wmark/work/watermark_report.json, write size, lateness, events, windows, on_time, late, dropped, restated, max_lag, p95_lag, and policy, and in /root/wmark/work/watermark_report.md, write four sections: ## 무엇을 재었나 ## 허용 지연을 얼마로 잡았나 ## 늦게 온 자료를 어떻게 했나 ## 어제 숫자가 바뀐 이유 (in order: what was measured, what allowed lateness was chosen, what was done with the late data, and why yesterday's numbers changed).
Choose the allowed lateness from 1 up to less than the maximum lag, and at that value at least 1 late-arriving event must appear. If none appears, reduce the value. In the report, write the count of late-arriving events as a number — that one number is the answer to the first question that comes up in the next meeting. You can simply call the functions you built in the earlier steps.