skip to content

Which stages in a log-shipping pipeline retain items internally even though nobody declared a buffer?

level: middleimportance: should knowfreq 44%

answer

  1. retention hides in ordinary stages
  2. can it emit after one item
  3. windows, flattening, pairing
  4. bounded by a constant or by rate
  5. count times item size per stage

basics

~20 s

Any stage that cannot emit until it has seen more than one item retains them. Time-based grouping holds a whole window, concurrent flattening holds every live inner source and its results, and handing values to another worker needs a queue.

solid answer

~50 s

Audit a chain with one question per stage: can it emit its output after seeing a single input item? If not, it is holding items, whether or not the word buffer appears. Three families dominate. A **time-based grouping** stage retains every line of the window before it emits, so its depth is the input rate times the window length. A **flattening** stage that subscribes to several inner sources at once retains each live inner source plus its in-flight results, and with no concurrency limit that set is unbounded. A **combining** stage retains the latest item from each input — small and constant — but its pair-by-index variant retains an entire input's backlog when one side runs faster. Add anything that de-duplicates, sorts, or hands values across a worker boundary. Then classify each: retention bounded by a constant, or bounded by the input rate?

code

pseudocode · 7 lines
pseudocode
readLogLines()                          // about 12000 lines per second
  .groupWithinEachInterval(5 seconds)   // retains 60000 lines before it emits
  .flattenConcurrently(limit = NONE)    // one live inner write per group, unbounded
  .writeEach(batch -> archive.append(batch))

// declared holds: none
// actual retention: 60000 lines in flight, plus every batch whose write has not finished

go deeper

for a junior

Learn the one-item test: if a stage cannot produce output until it has seen several inputs, it is keeping those inputs in memory right now.

for a middle

Name the families and their bounds: a time window holds rate times duration, unlimited flattening holds one live inner source per arrival, pairing by position holds the faster input's backlog.

for a senior

Audit a chain the way you would audit a budget: item count and item size per retaining stage, every rate-proportional bound converted to a constant, and depth exported so the audit survives the next edit.

for a principal

What you standardise is the review question, not the operator list. Requiring every stage's retention bound to be written down makes memory a design-time number that a reviewer can check instead of an incident finding.

## The test that finds a retaining stage Engineers look for retention by searching for the stage whose name contains "buffer". That finds the declared holds and misses the ones that matter, because several ordinary operators retain items as an unavoidable consequence of what they compute. The reliable test is a question you ask of every stage in the chain: > **Can this stage emit its output after seeing exactly one input item?** If the answer is no, the stage must hold input somewhere until it can. That is the whole audit. It needs no knowledge of any particular implementation, and it works on a chain you are reading for the first time. ## The three families you will actually meet - **Time-based grouping.** A stage that collects lines for a fixed interval and emits them as one batch must retain every line in that interval. Its depth is `input rate x window length`: at 12,000 lines per second and a 5-second window, 60,000 lines are resident at the moment of emission — and if the downstream writer is slow, the previous batch is still alive while the next one fills. - **Concurrent flattening.** A stage that maps each incoming item to an inner source and merges the results retains one live inner source per in-flight item, plus whatever each has produced but not yet delivered. With a concurrency limit of `k` the retention is bounded by `k`; with no limit it is bounded by the arrival rate, which is the same unbounded problem wearing a different name. - **Combining inputs.** A stage that emits a combination whenever any input produces, using the latest value from each, retains exactly one item per input — constant and harmless. The variant that pairs the nth item of one input with the nth of another is the dangerous one: if one input runs faster, everything it produces is retained waiting for its partner, and the retention grows with the rate difference. Beyond the three, two more are easy to miss: a stage that **de-duplicates by remembering keys it has seen** retains a set that grows with distinct keys forever, and a stage that **moves values to another worker** needs a handoff queue by construction, so every switch of execution context is a retention point. ## Bounded by a constant versus bounded by the input rate The classification that matters is not "does it retain" but "what bounds the retention". | Stage shape | What it retains | Bound | |---|---|---| | Group by fixed time window | Every item in the window | Input rate times window length | | Group by fixed count | Items up to the count | A constant | | Flatten with a concurrency limit | Live inner sources up to the limit | A constant | | Flatten with no limit | One inner source per arrival | Input rate | | Combine using latest per input | One item per input | A constant | | Pair items by index | The faster input's backlog | Rate difference | | Distinct by remembered key | Every distinct key seen | Cardinality of the key space | | Handoff to another worker | Items in the handoff queue | The declared queue capacity | A constant bound is a cost you can compute once. A bound proportional to the input rate is the unbounded-hold failure with extra steps: it grows for as long as the mismatch lasts. ## Auditing a real chain Walk the log-shipping chain stage by stage and write down two numbers for each: **how many items it can hold** and **how large each is**. Sum them. If any line of that table says "depends on the input rate", that stage is where memory will go, and the total is not a number at all — it is a growth rate. The audit also tells you where to put a limit. Bounding the declared hold at the end while leaving an unlimited flattening stage in the middle moves the collapse rather than preventing it, because the retained items simply accumulate one stage earlier. ## What to do with what you find 1. Give every rate-bounded stage a constant bound: a concurrency limit on flattening, a size cap as well as a time cap on grouping, an eviction rule on a remembered-key set. 2. Convert window lengths into item counts before accepting them, using the peak input rate rather than the average. 3. Export depth per retaining stage, so that the audit stays true after someone edits the chain. The point an interviewer is testing is whether you read a pipeline as a set of retention points with bounds, rather than as a list of transformations. Once you do, a chain's worst-case memory is an arithmetic exercise instead of a surprise.

  • How do you turn a time-based window into a bound you can defend?
    Multiply the peak input rate by the window length to get an item count, then multiply by the retained bytes per item. Declare that count as an explicit size cap alongside the time cap, so the stage emits early under a burst instead of growing with it.
  • A stage keeps the latest value from each of two inputs. Is that a retention risk?
    No. Its retention is one item per input, constant regardless of rate, so it costs a fixed amount forever. The risky relative is the stage that pairs items by position, which retains the faster input's entire backlog while it waits for the slower one's partner.
  • Why is bounding only the final hold not enough?
    Because retention accumulates wherever the chain cannot emit. Capping the last stage makes it refuse items, and the items then pile up in the unbounded stage before it. The collapse moves upstream rather than disappearing; every rate-bounded stage needs its own constant bound.

saying these in an interview costs you the question

  • Only stages named buffer retain items
  • A time window is cheap because it holds one batch at a time
  • Flattening without a concurrency limit is fine if each inner source is small
  • Pairing two inputs by position costs the same as keeping the latest of each
  • Bounding the last stage bounds the whole chain