A stage in an export pipeline fetches ledger rows a page at a time but emits them one by one — how should it convert downstream demand into upstream requests?
answer
- demand is not conserved across a stage
- two counters, upstream and downstream
- ask upstream in pages
- replenish at a low-water mark
- a dropped value needs a replacement request
basics
~20 sNot one for one. The stage keeps two counters — what it owes downstream and what it has asked upstream — requests a whole page ahead, and replenishes upstream only when the page has drained to a low-water mark, never asking for more than it has room to hold.
solid answer
~50 sDemand is **not conserved** across a stage. A rewriting stage keeps two separate counters: the outstanding count it owes downstream, and the outstanding count it has asked of upstream. Downstream asks for one row; upstream it asks for a page and buffers the result, because a request signal per value is pure chatter on a paged source. It then replenishes at a low-water mark — after roughly three quarters of a page has been emitted, ask for that many more — so the total it holds or is owed never exceeds its own storage. Two invariants keep backpressure intact: it never emits more than its downstream count allows, and it never asks upstream for more than it can absorb. A filtering stage adds a third: a dropped value must be replaced with a fresh upstream request, or the pipeline stalls.
code
pseudocode · 18 linespageSize = 64
replenishAt = 48 // ask again once 48 of the 64 have left
on subscribeToSource():
requestUpstream(pageSize)
on requestSignalFromBelow(n):
outstandingDown = outstandingDown + n
drain()
function drain():
while outstandingDown > 0 and buffer not empty:
emit(buffer.removeFirst())
outstandingDown = outstandingDown - 1
emittedSinceAsk = emittedSinceAsk + 1
if emittedSinceAsk == replenishAt:
emittedSinceAsk = 0
requestUpstream(replenishAt) // 16 still held + 48 owed = 64go deeper
Know that a value passing through a pipeline crosses several stages, and that the number a stage asks of its own source need not equal the number it was asked for.
Explain the two counters and why a paged source is asked for a page while the consumer is asked one at a time. Be able to state what bounds the upstream request: the stage's own storage.
Diagnose the real failures: the pipeline that stalls because a filtering stage never replaced dropped values, and the stage that quietly ended backpressure by asking upstream for more than it could hold.
Decide the defaults every stage in the codebase inherits. Buffer size and replenishment threshold set the latency-versus-chatter point for the whole system, and a reviewable rule beats each author picking a number.
## Demand is not conserved across a stage A naive mental model has a request signal travelling upstream unchanged, like a token passed hand to hand. That is true only of the simplest pass-through. Any stage whose input and output counts differ — a stage that batches, filters, expands, or fetches in pages — must **translate** the number, and how it translates is a design decision with real failure modes on both sides. The setting: a stage in a ledger export sits above a paged source and below a report writer. The writer commits one row at a time and therefore asks for one row at a time. The source answers in pages of sixty-four. A request signal per row would mean sixty-four round trips to fetch one page's worth of data, and on a remote source each of those is latency the export pays for nothing. ## The two counters a rewriting stage keeps 1. **Downstream outstanding** — how many values it is currently permitted to emit. Incremented by requests arriving from below, decremented by each value it emits. 2. **Upstream outstanding** — how many values it has asked its own source for and not yet received. Incremented by each request it sends up, decremented by each value that arrives. These are independent numbers, and the stage's job is to run a policy between them. Conflating them into one counter is the single most common way a rewriting stage is written wrong: the stage then either asks upstream for exactly what it was asked for (chatter) or forwards a large upstream request every time something small is asked below (it will be handed more than it can hold). ## Replenishing at a low-water mark The workable policy is: on subscribe, ask upstream for a full buffer's worth. Emit downstream whenever there is both demand below and a value in hand. Count the emissions, and once some fraction of the buffer has drained, ask upstream for exactly that many more. With a buffer of sixty-four and a replenishment threshold of forty-eight, at the moment the stage asks for forty-eight it still holds sixteen — so the most it can ever hold or be owed at once is sixty-four, exactly its capacity. That arithmetic is the whole safety argument, and it should be checkable by hand in any stage you write. The threshold trades two costs against each other: - **Threshold near the full buffer** — one request per large batch, minimal signalling, but the buffer runs closer to empty before the replacement arrives, so a slow source can leave the stage idle. - **Threshold near one** — the buffer stays full, but you are back to a request signal per value, which is the chatter the batching existed to avoid. ## Different stage shapes translate differently | Stage shape | Downstream asks for 10 | What it should ask upstream | |---|---|---| | one in, one out | 10 | 10, or a batch with a low-water mark | | fetches in pages | 10 | a page, buffered, replenished at a mark | | filters | 10 | 10, plus one more for every value it drops | | expands one into many | 10 | 1, and the next only once that expansion is drained | | combines N inputs into one output | 10 | 10 from each input it must pair | The filtering row is the one that breaks pipelines in production. A stage that receives a value and drops it has spent a unit of its upstream demand and produced nothing downstream. If it does not request a replacement, upstream demand ratchets down with every dropped value until it hits zero, and the stream goes quiet with the consumer still waiting. The symptom is a pipeline that works in testing, where few values are filtered, and stalls in production, where most are. The expanding row is the mirror hazard. A stage that turns one input into many outputs must not pull the next input while the current expansion is undelivered, or it accumulates output it has no permission to emit. ## The invariants worth stating out loud - Never emit more values downstream than the downstream count permits. - Never request upstream more than the stage can hold if all of it arrives at once. - Never let a consumed-but-not-emitted value silently reduce the upstream count. - Replenish from what has actually left the stage, not from what has arrived. Hold those four and backpressure survives the stage: a slow report writer at the bottom still eventually slows the paged source at the top, just at page granularity instead of row granularity. Break the second one and the stage becomes the place where flow control quietly ends, while every stage above and below it still looks correct in isolation.
- A filtering stage drops a value it received. Why does the pipeline eventually stall if it does nothing else?The drop consumed a unit of the stage's upstream demand but produced nothing to satisfy the demand below. Repeat that and the upstream count walks down to zero while the downstream count is still unfilled: the source is no longer permitted to send, and the consumer is still waiting. The stage must request a replacement for each value it drops.
- What actually bounds how much a rewriting stage may request upstream?Its own storage. The request is a promise to accept that many values whenever they arrive, so the safe bound is what the stage can hold with everything outstanding delivered at once. Requesting more than that is where flow control ends, even though every neighbouring stage still looks correct.
- How does the replenishment threshold change the behaviour?It trades signalling against idleness. Replenishing late means few request signals but a buffer that can run dry before the replacement arrives; replenishing early keeps the buffer full but approaches one signal per value, which is the cost the batching was meant to avoid.
saying these in an interview costs you the question
- Assumes a request signal passes upstream unchanged through every stage
- Keeps one counter for both the upstream and the downstream side
- Forwards a full page request upstream for every single value asked below
- Forgets that a dropped value must be replaced with a new upstream request
- Requests upstream more than the stage has room to hold
- Pulls the next input while an expansion of the current one is still undelivered