skip to content

Why must a stream step that carries a running year-to-date total create its state per subscription?

level: middleimportance: should knowfreq 46%

answer

  1. state always has an owner
  2. a chain is only a description
  3. each subscription is its own run
  4. captured once means shared by every run
  5. seed per subscription, discarded with it

basics

~20 s

Each subscription is an independent run of the chain. State created once when the chain was built is shared by every subscriber, so one run's totals leak into another; state created per subscription starts from the seed and stays private to that run.

solid answer

~40 s

A chain of steps is a **description** that can be run many times; each subscription is one run of it. A running year-to-date total is state, and where that state is created decides who sees it. If the total lives in a variable captured when the chain was assembled, every subscriber shares one counter: a second payroll run starts at the first run's total, and two concurrent runs interleave their updates. If the accumulating step creates the state when the subscription starts, each run gets its own total, starting from the seed, and it disappears with the run. That is also why in-stream accumulation never resumes a failed run where it stopped — the new subscription gets new state, and continuity has to come from somewhere durable instead.

code

pseudocode · 11 lines
pseudocode
// shared: one total, captured when the chain was built
yearToDate = 0
pipeline = payLines.transformEach(function(line) {
    yearToDate = yearToDate + line.grossPay   // every run updates the same value
    return withYearToDate(line, yearToDate)
})

// per subscription: the accumulating step makes a fresh total for each run
pipeline = payLines.accumulate(0, function(yearToDate, line) {
    return yearToDate + line.grossPay
})

go deeper

for a junior

Remember that a chain of steps can be run many times, and that a step holding a running total must start that total afresh for each run rather than keeping one value on the side.

for a middle

Explain the two possible homes for the state — captured when the chain is built, or created when a subscription starts — and describe exactly what a second subscriber sees in each case.

for a senior

Demonstrate the production angle: name what a shared counter does to concurrent runs, and say where continuity comes from when a run fails part-way, since in-stream state does not survive it.

for a principal

The call to own is how much accumulated state belongs inside pipelines at all. State that matters to the business usually wants a durable home with a seed read at the start of a run, not a value living inside a subscription.

## Stateless steps, and the one that is not Most element-wise steps in a stream pipeline are **stateless**: a mapping step's output depends only on the element in front of it, and a predicate filter's decision does the same. An **accumulating step** is different. To emit a year-to-date gross after each pay line it must hold the result so far and fold the arriving line into it. That held value is state, and state always raises the same two questions: **who can see it** and **how long does it live**. ## A chain is a description; a subscription is a run Building a pipeline does not move any element. The chain of steps is a description of the work, and nothing happens until a subscriber asks for values. Each such subscription is an independent run of that description, and for a source that produces its values per subscriber, two subscriptions mean two separate passes over the data. That gives the accumulating step's state exactly two possible homes: 1. **Captured when the chain is assembled.** One variable exists, created once, shared by every subscription that ever runs the chain. 2. **Created when a subscription starts.** A fresh value exists per run, reachable only by that run's elements, discarded when the run ends. Only the second is correct for a running total, and the first is one of the most common bugs in reactive code. ## What the shared counter actually does With one variable shared by every run of a payroll pipeline: - **The second run starts dirty.** A payroll run for the next period begins with the previous period's total already in the counter, so every year-to-date figure is wrong by a constant. - **Concurrent runs corrupt each other.** Two subscriptions active at once interleave their updates into one counter; neither total is meaningful, and the values are not even wrong in a stable way. - **The bug is invisible with one subscriber.** In development there is usually exactly one run and the code looks correct, which is why this reaches production. - **It survives a failure.** After a failed run leaves the counter part-way, the next attempt inherits a value that belongs to no complete run. - **The state outlives its purpose.** A captured value is reachable for as long as the chain is, so anything it holds — a collection being accumulated, for instance — is never released. ## Side by side | | Captured when the chain is built | Created per subscription | |---|---|---| | One subscriber | looks correct | correct | | Two subscribers | one shared, interleaved total | two independent totals | | Second run of the same chain | starts from the previous run's value | starts from the seed | | After a failed run | keeps the part-way value | discarded with the run | | Lifetime | as long as the chain exists | as long as the subscription | | Testing | order of tests matters | each test is isolated | ## How to keep the state where it belongs 1. **Let the accumulating step own it.** A step whose contract is "fold the result so far into each element" is given its seed and creates its own working value per subscription. That is the whole reason the shape exists rather than you keeping a variable on the side. 2. **Do not mutate anything declared outside the chain from inside a step.** A counter, a list being appended to, a map being filled — if it was created before the subscription, it is shared by every subscription. 3. **If a run must genuinely resume, make that explicit.** Seed the accumulation from a durable value read at the start of the run. Then continuity is a deliberate input to the run, not an accident of a variable that happened to survive. 4. **Prove it with a second subscriber.** The cheapest test of this property is to run the same chain twice and assert that the second run produces the same totals as the first. ## What interviewers listen for The phrase that earns credit is "per subscription". A candidate who says the state must be created when the run starts, and can then explain what a shared counter does to a second subscriber, has understood that a pipeline is a reusable description rather than a running object. The follow-up that separates levels is the failure case: in-stream accumulation does not resume where a failed run stopped, and the candidate who knows that will reach for a durable seed instead of hoping the value survived.

  • A run fails after 800 of 1,000 pay lines and the subscriber subscribes again. What happens to the running total?
    It starts again from the seed. The new subscription is a new run of the chain, so the accumulating step creates new state and recomputes from the first line it sees. Nothing in the pipeline stores the part-way value, so if the resumed run must continue from line 800 the total has to be seeded from a durable record.
  • Two subscribers subscribe to the same chain at the same time. How many running totals exist?
    Two, one per subscription, each starting from the seed and invisible to the other — provided the state is created when the subscription starts. If it was captured when the chain was built there is exactly one, and the two runs interleave their updates into it, which is the bug this rule exists to prevent.

A running total is like a till roll: every checkout lane needs its own. One roll shared between two lanes prints a total that belongs to neither of them.

saying these in an interview costs you the question

  • Keeps the running total in a variable captured when the chain was built
  • Assumes a second subscriber continues the first subscriber's total on purpose
  • Thinks in-stream accumulator state survives a failed run and resumes
  • Mutates a collection declared outside the chain from inside a step
  • Argues the shared counter is safe because there is only one subscriber today