skip to content

Bounding an Endless Stream

An endless input supplies no boundary of its own, so every computation over one has to impose one: an interval for an aggregate, a time bound for a match, the record's own moment for a lookup.

on this pageshow

explore

questions

21

While a five-minute grouping computes a running total per key, what must the job hold, and what changes for a median?

level: juniorimportance: must knowfreq 62%

answer

  1. two ways to hold an open group
  2. carry a value or carry the records
  3. do two partial values merge
  4. total folds, exact median keeps everything

basics

~20 s

A running total needs one number per key while the group is open. A median needs every value in the group, because no single carried value can be merged into the answer. Same interval, very different memory.

solid answer

~50 s

There are only two ways to hold a group that is still open. You can **fold** each arriving record into a per-group value that already carries everything the answer needs and then discard the record — for a total, one running number — or you can **keep the records themselves** until the group is declared finished and compute the answer from all of them. A total folds: add the amount, drop the record, and two partial totals add together. A median does not: which value ends up in the middle depends on records that have not arrived yet, so every value stays. That is the whole cost difference — `one value per open group` against `every record in every open group`. Most aggregates people actually ask for fold; the ones that do not are the ones that need the group's whole distribution.

go deeper

for a junior

Recall the two options for an open group: carry one value folded from the records, or keep the records themselves. Know that a running total is the first and an exact median is the second.

for a middle

Explain the property that separates them — whether two partial values merge into one without growing — and show that a mean folds when carried as a total and a count.

for a senior

Show the production consequence: per-group size times open groups, and that the aggregate chosen, not the boundary chosen, is often what decides whether the job fits at all.

for a principal

Treat it as a requirements question. An exact median asked casually is worth challenging when an approximate one costs a constant per group and the decision it feeds tolerates the error.

## Two ways to hold a group that is open An endless input has no end of its own, so a computation over it finishes only when the author imposes a boundary — for example a boundary that cuts the input into back-to-back spans of equal length, so that every record falls in exactly one group. Between the arrival of a group's first record and the moment the group is declared finished, that group is **open**: more records for it are expected, and its answer cannot be handed downstream yet. While a group is open there are only two things the job can do with an arriving record. 1. **Fold it.** Combine the record into a single per-group value that already carries everything the final answer needs, then discard the record. A value like this is a **foldable accumulator**, and it has two properties that matter: any two of them, built over different records of the same group, merge into one; and its size does not grow with how many records have been folded into it. 2. **Keep it.** Add the record to the group's **retained record set**, and compute the answer from the whole set when the group is declared finished. Real jobs sometimes do both side by side — a folded total next to a retained set for a second output that genuinely needs the records — but there is no third mechanism. ## Why a total folds and a median does not A running total's carried value is one number. When a record arrives you add its amount, and the record is then worthless to the answer: nothing about it can change the total again. Count, smallest and largest behave the same way. A mean behaves the same way too, once you carry a total **and** a count as a pair rather than carrying the average itself. A median is different in kind. The middle value is a property of the whole set's shape, and which record turns out to be the middle one depends on every record that arrives afterwards. There is no fixed-size value you can carry such that the next record folds into it and the exact middle is still recoverable, so every value is kept until the group closes. Exact distinct counting fails for a related reason — to know whether an arriving identifier is new, you must have kept the identifiers — and an exact top-k by frequency needs a count for every distinct key, not only for the k that end up ranked. | Computation | Carried while the group is open | Cost per open group | |---|---|---| | running total, count | one number | constant | | smallest, largest | one value | constant | | mean | a total and a count | constant | | exact median or percentile | every value in the group | grows with records | | exact distinct count | every distinct identifier | grows with distinct values | | exact top-k by frequency | a count per distinct identifier | grows with distinct values | ## The arithmetic that follows The memory a grouping costs is what one open group holds, multiplied by how many groups are open at the same time. Folding collapses the first factor to a constant. That is why the same boundary over the same input can be a job that runs comfortably or a job that cannot be made to fit: the difference is the aggregate, not the interval. What a runtime does when the total will not fit in a worker's memory, and where those bytes physically live, is a separate subject from this one. ## Where engines of this class differ The fold-or-keep distinction is universal. Three things around it are not, and an answer that states one runtime's behaviour as the model will be wrong for a candidate whose next job uses a different one: - **Where the carried value lives between records.** Where each record advances individually through long-lived operators, the accumulator sits in the operator and is updated per arrival. Where arrivals are instead collected for a short span and one finite job runs over the collected set, the accumulator is handed from one such job to the next. In the oldest model in this family — a finite pass that materialises its intermediate result to disk between phases — no group persists at all, and a windowed result is produced by re-running the pass. - **When the group's value is read out.** Some runtimes hand a value downstream once, when the group is declared finished. Others emit an early value and corrections afterwards. Others restate the group's current value on every input. A group that may still re-emit has released nothing, so output appearing is not the same event as memory falling. - **Whether your computation is recognised as foldable.** A declarative aggregate expression is usually compiled into a fold. A hand-written per-record function that appends to a collection is usually executed literally as written and folds nothing, however foldable the underlying mathematics is. ## What an interviewer is listening for Not "a median is expensive". The mechanism: name what is carried, say whether two partial values merge into one without growing, and then state the memory as a product of per-group size and open-group count. A candidate who jumps straight to giving the workers more memory has skipped the only question that was asked.

  • Does a mean fold, given that an average of averages is usually wrong?
    Yes, provided the carried value is a total and a count rather than the average itself. Merging adds total to total and count to count, and the division happens once, when the group's value is produced. Carrying the average alone is exactly what breaks, because two averages cannot be combined without knowing how many records each covered.
  • A group emits its value and memory does not fall. What happened?
    Emission and release are different moments. Where the runtime allows corrections, a group that has been declared finished still holds its carried value through a stated extra span, so it can accept a record that turns up afterwards and re-emit. Memory falls when the carried value is dropped, not when output appears.

A shopkeeper can keep a running till total on one slip of paper and throw each receipt away, and still answer how much the shop took today. To answer what the middle-sized sale was, every receipt has to stay in the box until closing time.

saying these in an interview costs you the question

  • Assumes every aggregate can be carried as one running number.
  • Says a median can be folded by carrying a running average.
  • Thinks the interval's length alone decides what the grouping costs.
  • Believes records are always discarded once counted, whatever the aggregate.
  • Treats the memory question as answered by giving workers more memory.
open as a page

When each arriving order is enriched from a reference table, what supplies the bound for the computation and when is a result emitted?

level: juniorimportance: must knowfreq 60%

basics

~20 s

The record's own moment supplies the bound: it selects one version of the reference row, the one whose span of validity contains that moment. Nothing accumulates and no group closes, so each record yields its output once its match can be resolved.

open as a page

Two endless inputs are joined on a shared key. Why must the job hold records from both inputs, not just one?

level: juniorimportance: must knowfreq 62%

basics

~20 s

Neither input ends, so a record on either side may meet its partner later. Both sides are therefore held while the other is awaited, and a bound on how far apart their two moments may be is what lets a held record ever be dropped.

open as a page

A grouping uses ten-minute spans restarted every two minutes. How many groups does one record belong to, and why?

level: juniorimportance: must knowfreq 70%

basics

~10 s

Five. A span of fixed length restarted every step shorter than that length overlaps its neighbours, so a record's moment falls inside length-divided-by-step spans at once and is counted in every one of them.

open as a page

Over an endless input, what three ways can a grouping boundary be defined, and what sets each group's extent?

level: juniorimportance: must knowfreq 76%

basics

~20 s

Three shapes: back-to-back spans of equal length, where a record lands in exactly one group; fixed-length spans restarted every shorter step, where a record lands in several at once; and groups the data itself closes after a stated quiet period.

open as a page

Which three contracts can a job offer for emitting a grouped interval's value downstream, and what does each demand of the sink?

level: middleimportance: must knowfreq 62%

basics

~20 s

Emit once when the group is declared finished, emit early and correct later, or restate the running value on every input. Only the first suits an insert-only sink; the other two need a sink that replaces a value it already holds.

open as a page

Which aggregates fold into one mergeable value per group, which force keeping every record, and where does mean sit?

level: middleimportance: must knowfreq 58%

basics

~20 s

Sum, count, smallest, largest and mean fold into a fixed-size value per group. Exact median, exact distinct count and exact top-k do not. Mean folds only as a total-and-count pair, never as a running average.

open as a page

In an enrichment join, what separates matching the reference row that is current now from matching the version valid at the record's own moment?

level: middleimportance: must knowfreq 57%

basics

~20 s

Current-row matching asks what is true now; as-of matching asks what was true when the record happened, by picking the version whose validity span contains the record's moment. The two agree only while the stream is caught up and the value has not changed since.

open as a page

Two inputs each arrive at 5,000 records a second under a 30-minute match bound. What is retained, and what does doubling the bound cost?

level: middleimportance: must knowfreq 58%

basics

~20 s

Each side holds arrival rate times the bound: 5,000 a second over 1,800 seconds is 9 million records per side, 18 million across both. At a fixed rate and a symmetric bound, doubling the bound doubles the records held and the memory they occupy, while the extra pairs it recovers follow the tail of the gap distribution.

open as a page

In a job grouping an endless input into intervals, what rule, separate from the boundary, decides when a group's value goes downstream?

level: juniorimportance: should knowfreq 55%

basics

~20 s

The firing condition: the rule picking the moments a group's current value is handed downstream. It is a decision separate from the boundary, which only says which records belong together, so one group may be handed downstream several times.

open as a page

A job folds one small value per group, yet runs out of room — which two quantities multiply to set that memory, and what inflates each?

level: middleimportance: should knowfreq 46%

basics

~20 s

Memory is what one open group holds multiplied by how many groups are open at once. A tiny accumulator still costs plenty when the grouping key has millions of values, or when intervals restarted before they end put each record into several groups.

open as a page

A job holds a local copy of a reference table - which three approaches keep that copy current, and what does each cost?

level: middleimportance: should knowfreq 48%

basics

~20 s

Reload the whole reference periodically, apply a feed of its committed changes, or look each record up in the reference store directly. They trade staleness, load on that store and per-record latency against each other, and are often combined as an initial load plus a feed.

open as a page

A per-user grouping closes after five quiet minutes. What sets each group's length, and what can keep one open indefinitely?

level: middleimportance: should knowfreq 58%

basics

~20 s

The arrivals set it. The group stays open for one user while records keep coming and closes only once five minutes pass with none, so its length is data rather than definition — and a user who never pauses for five minutes keeps a group open forever.

open as a page

A stakeholder asks for an alert on errors in the last five minutes. What must be pinned down before it can be computed?

level: middleimportance: should knowfreq 60%

basics

~20 s

The phrase fixes only a length. Still undefined: whether the five minutes advances continuously or resets on a grid, how often it is recomputed, what groups the records together, which moment it is measured against, and what counts as an error.

open as a page

An hourly per-customer total emits early and corrects later, and the sink upserts on customer id: what breaks, and what identity is right?

level: seniorimportance: should knowfreq 47%

basics

~20 s

Each hour's value overwrites the previous hour's, so the sink ends up holding one row per customer showing only the latest hour. The identity must be the grouping key together with the interval's start, so revisions land on their own row.

open as a page

A job re-reads last month's records from its append-only source, but the enrichment attaches today's prices - what in the join caused that, and what makes the re-run reproducible?

level: seniorimportance: should knowfreq 40%

basics

~20 s

The join takes the reference value current at processing time, and the copy holds only current values, so old records are enriched with today's numbers. Reproducibility needs a reference that retains each version with its validity span and a predicate on the record's own moment.

open as a page

A payment arrives whose matching order never appears within the match bound. What becomes of that unmatched record, and who decides?

level: seniorimportance: should knowfreq 52%

basics

~20 s

Its fate is a declared property of the join, not a runtime accident: dropped silently, or emitted once with a null counterpart after the bound passes. The pipeline author decides, and a design that does not state which has chosen the silent drop by default.

open as a page

A report sums per-group results into daily totals. Which grouping shapes may it sum safely, and which will double count?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Only shapes whose groups partition the input may be summed: back-to-back equal spans globally, and gap-closed groups per key. Overlapping spans cover each record length-divided-by-step times, so summing them inflates the total by that factor.

open as a page

A dashboard needs numbers within seconds, but a group is only declared finished minutes later: how do you decide what to publish?

level: principalimportance: should knowfreq 40%

basics

~20 s

Decide what wrongness you will publish, not whether to publish any. Emitting before a group is declared finished is defensible when the size and duration of the error are stated, the number is visibly provisional, and each consumer gets the contract it can actually apply.

open as a page

Two percent of pairs miss a one-hour match bound and the business asks for a 24-hour bound. How do you decide?

level: principalimportance: should knowfreq 45%

basics

~20 s

Cost is linear in the bound and recall is not: 24 times the held records on both sides buys only the pairs whose moments lie one to 24 hours apart. Measure that gap distribution first, then weigh a wider bound against reconciling the remainder elsewhere.

open as a page

An exact distinct count per group holds every identifier seen — what does a bounded-size summary change, and what must you agree first?

level: seniorimportance: nice to knowfreq 32%

basics

~20 s

A fixed-size summary answers distinct-count, quantile or heavy-hitter questions within a stated error and merges the way an accumulator does, so a group costs a constant instead of one entry per distinct value. The price is an error budget somebody must accept.

open as a page