TT Lab
Get started
Learn Learning paths Courses

Data Pipelines

Yesterday's Number Changed Today: Watermarks and Late Data

Continue in TT Lab

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

  1. Create and run /root/wmark/gen_stream.py to create /root/wmark/work/stream.jsonl.
  2. In /root/wmark/wm.py, create skew so that it produces the lag distribution and the number of out-of-order events.
  3. Add watermark to move the watermark in arrival order.
  4. Add windows to divide and aggregate by event time.
  5. Add late to separate the data that came after a window closed.
  6. Add agg to produce the results of the drop policy and the apply policy.
  7. Add close to produce the value at the moment of closing and the corrections afterward separately.
  8. Write /root/wmark/work/watermark_report.json and /root/wmark/work/watermark_report.md.

Reference

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.