skip to content

In a stream pipeline, how do a mapping step, a predicate filter, and a running-total accumulator differ?

level: juniorimportance: must knowfreq 70%

answer

  1. how many outputs per input
  2. one in, one out
  3. predicate decides keep or drop
  4. accumulator folds previous result forward
  5. stateful step remembers earlier elements

basics

~20 s

A mapping step returns exactly one output element per input; a predicate filter returns zero or one, dropping the rest; a running accumulator emits a value derived from every element seen so far, so it carries state between elements.

solid answer

~40 s

I separate them by **how many elements leave for each one that arrives** and **whether the step remembers anything**. A mapping step is one-in, one-out and stateless: a timesheet row becomes a pay line and the element count is unchanged. A predicate filter is one-in, zero-or-one-out and also stateless: it decides keep or drop, it never rewrites the element. A running accumulator is one-in, one-out as well, but stateful: each output folds the previous result together with the current element, so a year-to-date gross grows as pay lines go by. The related whole-sequence fold emits a single value at completion instead of one per element. All three are element-wise, so they preserve the upstream order and introduce no concurrency of their own.

code

pseudocode · 10 lines
pseudocode
pipeline = timesheetRows
    .transformEach(function(row) {
        return payLineFor(row)              // 1 in -> 1 out, no memory
    })
    .keepIf(function(line) {
        return line.payableHours > 0        // 1 in -> 0 or 1 out, no memory
    })
    .accumulate(0, function(yearToDate, line) {
        return yearToDate + line.grossPay   // 1 in -> 1 out, carries the total
    })

go deeper

for a junior

Be able to say in one line each how many elements leave for every element that arrives: exactly one for a mapping step, zero or one for a predicate filter, a running result for an accumulating step.

for a middle

Explain the stateless versus stateful split and name the state: the accumulating step holds the result so far and folds the arriving element into it, which is what gives it a lifetime the other two do not have.

for a senior

Show the judgement. Reach for the stateless shapes first, and treat every stateful step as state you must size, scope to one run and be able to explain when a pipeline misbehaves in production.

for a principal

The trade-off worth owning is how much accumulated state lives inside pipelines at all rather than in a store that survives a failed run: in-stream state is invisible during an incident and cannot be inspected after the fact.

## What an element-wise step is A stream pipeline is a chain of steps sitting between a **source** that pushes elements over time and a **subscriber** that consumes them. An **element-wise step** looks at one element at a time and decides what, if anything, leaves it. Three such shapes do most of a pipeline's everyday work, and they are told apart by two properties: - **cardinality** — how many elements leave for every element that arrives; - **state** — whether the step remembers anything between elements. Use a payroll run as the setting. A source pushes raw timesheet rows; the pipeline must turn each row into a pay line, drop the rows that should not be paid at all, and report a year-to-date gross that grows as the lines go by. ## Mapping: one in, one out A **mapping step** applies a function to each element and emits the result: `row -> payLineFor(row)`. Exactly one element leaves for each one that arrives, so the length of the sequence is unchanged; only the *type and content* of the elements change. It is **stateless** — the output for a row depends on that row alone — which is why it can be tested, reasoned about and re-run one element at a time. A mapping step makes no judgement about whether an element belongs and no judgement about when the sequence should end. ## Predicate filtering: keep or drop A **predicate filter** asks a yes/no question of each element and forwards it untouched when the answer is yes: `line.hours > 0`. Zero or one element leaves per input, so the sequence can only shrink. It is stateless too, and it is deliberately weaker than a mapping step: it may not rewrite the element it inspects. A filter that rejects everything still produces a perfectly valid sequence — one that simply completes without ever emitting a value. Note what it cannot do: it judges each element in isolation, so it never discovers that no further element can match, and therefore never ends the sequence early. ## Accumulation: the step that remembers An **accumulating step** combines the result so far with the arriving element: start from a seed, then `combine(runningTotal, payLine)`. This is the only one of the three that is **stateful** — it must hold the previous result to produce the next one. Two forms exist and confusing them is a common slip: - a **running accumulation** emits the intermediate result as each element arrives, so roughly one value leaves per element; - a **whole-sequence fold** emits a single value only when the source completes — on an endless source, that value never arrives. Implementations differ over whether the seed itself is emitted first, so treat that as a detail to check rather than a rule to recite. The state belongs to a single run of the chain: each subscription accumulates its own year-to-date total, starting from the seed. ## Side by side | Step | Outputs per input | Holds state | Can end the sequence early | Payroll use | |---|---|---|---|---| | Mapping | exactly 1 | no | no | timesheet row becomes a pay line | | Predicate filter | 0 or 1 | no | no | drop rows with no payable hours | | Running accumulation | about 1 | yes | no | year-to-date gross after each line | | Whole-sequence fold | 1, at completion | yes | no | total cost of the whole run | ## What these steps give you for free 1. **Order is preserved.** Each element is handled to completion before the next is offered downstream, and none of the three brings in a second source or starts work in parallel, so pay lines leave in the order the timesheet rows arrived. 2. **No concurrency is introduced.** Whatever thread or worker carries the elements carries the transform too; these steps neither add workers nor move work between them. 3. **They are cheap to combine.** Because a stateless element-wise step needs nothing but the current element, adjacent steps can often be fused into a single pass over each element rather than handing values through a queue at each stage. The stateful one is the odd member: it must be given a place for its running result. What you do *not* get for free is purity. Nothing stops a mapping function from writing a row somewhere, and nothing warns you that it did; keeping effects out of transforming steps is a rule you keep, not one the pipeline enforces. ## What interviewers listen for A weak answer lists three names. A solid answer names the cardinality of each shape, points out that the first two are stateless and the third is not, and adds the consequence: the stateless ones can be moved, re-run and reasoned about element by element, while the accumulating one owns state and therefore owns a lifetime. The strongest answers reach for the stateless shape first and treat every stateful step as a deliberate decision.

  • Does an accumulating step emit one value per element, or only a final value?
    Both forms exist. A running accumulation emits the intermediate result as each element arrives, so roughly one value leaves per element. A whole-sequence fold emits one value only when the source completes, which on an endless source means never. In a long-lived pipeline the running form is usually the one you want.
  • Which ordering guarantee do these three element-wise steps give you for free?
    They preserve the upstream order. Each element is handled to completion before the next is offered downstream, and none of them introduces a second source or a second worker, so pay lines leave in the order the timesheet rows arrived. Ordering is only at risk once a step starts several pieces of work at once.
  • When does reaching for an accumulating step instead of a stateless one cost you?
    As soon as the step is stateful you own the state: it has to be created for each run of the chain, it cannot be re-run for one element in isolation, and its correctness now depends on every element that came before. If an element's output depends only on that element, keep the step stateless.

saying these in an interview costs you the question

  • Says a predicate filter rewrites the element instead of keeping or dropping it
  • Claims a one-to-one mapping step may emit two elements for one input
  • Assumes every accumulating step emits only once, at completion
  • Thinks these element-wise steps may reorder elements or add concurrency
  • Reaches for a stateful accumulator where a stateless mapping would do