Before choosing how long a job waits for out-of-order records, what do you measure on the live source, and in what shape?
answer
- measure the source, not the pipeline
- distance behind the newest stamped moment
- reference receipt, never the worker's clock
- percentiles and a maximum, not an average
- split per input, release and hour
basics
~20 sMeasure, per record, how far behind the newest stamped moment already seen on that input it arrived. Report the result as a distribution with percentiles and a maximum, broken down per input and per hour, rather than as a single number.
solid answer
~50 sThe quantity is per record: take the newest **occurrence time** — the moment stamped in the payload — seen so far on that input, subtract this record's occurrence time, and you have how far behind the front it arrived. Compute it against **receipt time**, the moment the system that first accepted the record wrote it down, never against the wall clock of the worker that handled it: the second measurement absorbs your own backlog, so a job catching up after an outage reports hours of "disorder" it created itself. The answer is a *distribution*, not a constant — typically a tight body of a few seconds, a shoulder at the 99th percentile, and a maximum measured in days. Break it down per parallel input, per producer release and per hour of day, and re-measure: the curve belongs to the client population, not to your code.
go deeper
Know that the wait a stream job uses should come from a measurement of how the source actually behaves, not from a number someone liked the look of.
Be able to define the per-record quantity precisely — distance behind the newest stamped moment seen — and explain why referencing the handling worker's clock contaminates it with your own backlog.
Show that you read the curve rather than a statistic: shoulder against maximum, breakdowns per input and per producer release, and a re-measurement cadence tied to client releases.
Treat the distribution as an asset the organisation maintains: who owns it, how often it is refreshed, and which downstream promises are re-derived when it moves.
## The quantity you actually want You are not measuring how slow the pipeline is. You are measuring how far out of order the **source** delivers, which is a property of producers and transport and is already fixed before your job reads anything. The per-record quantity is straightforward. Track the newest **occurrence time** — the moment the thing happened, stamped into the payload by whatever produced the record — seen so far on that input. For each arriving record, subtract its occurrence time from that running maximum. A record that arrives ahead of everything contributes zero. A record arriving when the front has already moved ten minutes on contributes ten minutes. Call that the record's out-of-orderness. That is the only quantity a wait can honestly be chosen from, because a wait is precisely a duration held behind the newest stamped moment the job has seen. Measuring anything else and then choosing a wait from it is choosing a round number and calling it analysis. ## Measure against receipt, not against your own worker's clock There are two tempting reference points and one of them is wrong. | reference point | what it actually measures | trustworthy while catching up? | |---|---|---| | occurrence time against **receipt time** (when the accepting system wrote it down) | the source and its transport | yes — both values are fixed before this job reads | | occurrence time against the **worker's wall clock** at handling | the source, plus this job's own backlog | no — a job eight hours behind reports eight hours of "disorder" | The test is simple: a job that is fully caught up and the same job an hour behind must produce the same disorder measurement, or the measurement means nothing. The second reference point fails that test, and it fails it exactly when you most want the number — during an incident. ## It is a distribution, not a constant Real sources do not have *a* disorder. They have a shape, and the shape is what you read. | statistic | what it is for | |---|---| | median | describes the healthy path and is nearly useless for choosing a wait | | 95th / 99th percentile | the shoulder — the region where a wait first looks affordable | | 99.9th percentile | where each further nine starts costing an order of magnitude of waiting | | maximum | the disconnected-producer population; tells you a wait will never reach it | Break the distribution down, because an aggregate curve hides the thing you need: - **per parallel input**, since one input may be served by a slower region or a different producer fleet; - **per producer version**, since one mobile release that changes buffering behaviour moves the whole curve; - **per hour of day and per day of week**, since overnight and weekend traffic are different populations; - **per geography or network type**, where connectivity differs sharply. A source whose aggregate 99th percentile is forty seconds may be two clean populations: connected producers at under a second and a small reconnecting fleet at hours. Those two need different treatment, and the aggregate percentile hides that they exist. ## Where and when to take the measurement 1. **Offline, over a stored span of raw records** read from the durable, re-readable input. You can compute the running front in one pass without touching the running job, and you can span a weekend and a release. 2. **Live, as a histogram the job itself emits**, so drift is visible rather than discovered during an incident. 3. **Repeatedly.** The curve is a property of the client population. A new producer release, a new market, a new device class or a change in a mobile platform's background-execution rules all move it, and none of them touches your code. ## What the measurement does not tell you Being clear about the boundary is part of a good answer: - It does not tell you what the wait **should** be. That is a cost decision — how much end-to-end latency the consumer of the output can afford — and it is decided against the distribution, not read off it. - It does not tell you what to do with whatever still misses. That is a separate mechanism and a separate promise to downstream consumers. - It does not tell you whether the job is **behind**. Backlog is unprocessed records piling up; a job can be entirely caught up and still deliberately waiting, and a job with a perfectly ordered source can still carry an hour of backlog. ## What varies between engines What you can do with the measurement depends on the execution mode, so say which you mean. On a runtime that executes continuous work as a rapid succession of small finite jobs, the granularity at which any wait can act is at least one of those intervals, so a measured spread far below that interval cannot be acted on at all. On a record-at-a-time runtime, the measurement translates almost directly into a hold behind the newest stamped moment. And for a single pass over a finished bounded input, the measurement is informational: the whole input is present when the read begins, so disorder costs a sort rather than a wait.
- Why is an average out-of-orderness a poor basis for choosing a wait?Because the distributions are heavy-tailed. A handful of records hours behind pull the mean far above the body while representing almost nothing, and a wait set at the mean is simultaneously far too long for the common case and far too short for the tail. Percentiles and the maximum describe a shape the mean destroys.
- The measured curve suddenly widens overnight with no deployment on your side. What is the likely cause?The producer population changed: a client release altering buffering or background upload behaviour, a new device class or market, a network or regional degradation, or one parallel input being served by a slower path. Break the distribution down per input and per producer version and the shift usually localises to one of them.
- Can you measure disorder on a source you cannot re-read?Only live, and only going forward. The measurement needs a running maximum over stamped moments plus a fixed reference such as receipt time; without the ability to re-read a stored span you cannot compute the curve retrospectively, cannot cover a release or a weekend, and cannot re-derive it after a change.
saying these in an interview costs you the question
- Picks a round number such as five minutes with no measurement at all
- Measures against the handling worker's wall clock and calls it source disorder
- Reports a single average and treats disorder as a constant
- Measures once at launch and never re-measures as producers change
- Reads a large maximum as a sampling error rather than a second population
- Confuses the measurement with the job's backlog during an incident