In a payroll stream pipeline, what goes wrong when the mapping step also writes an audit row for each element?
answer
- a transform should only compute
- effects hide inside innocent-looking steps
- runs per element, per subscription
- failure leaves earlier writes applied
- declare the effect as its own step
basics
~20 sThe write becomes invisible in the chain and repeats with the run — once per element per subscription, again on resubscription, and a mid-stream failure leaves earlier writes applied. Keep the transform computing only and declare the effect separately.
solid answer
~50 sA transforming step's contract is that it computes an output from its input; anyone reading the chain assumes it does nothing else. Burying a write in it breaks three things. It is **invisible**: the chain says "turn a row into a pay line", so nobody expects rows in an audit store. It **multiplies with the run**: the write happens once per element per subscription, so a second subscriber or a re-subscription after a failure writes everything again. And it makes the transform **non-re-runnable**: a step you could previously replay for one element now has a consequence outside the pipeline, and a failure half-way leaves the first 800 writes applied and the last 200 not. Put the transform in a step that only computes, and declare the write in an observe-only step whose job is visibly an effect.
code
pseudocode · 9 lines// hidden: the effect lives inside a step that claims to compute
.transformEach(function(row) {
auditStore.write(row.id) // per element, per subscription
return payLineFor(row)
})
// declared: the transform computes, the effect is its own step
.transformEach(function(row) { return payLineFor(row) })
.observeEach(function(line) { auditStore.write(line.rowId) })go deeper
Learn the rule first: a step that turns one value into another should only compute. Writes, logging and anything touching the outside world belong in a step whose visible purpose is that effect.
Explain the mechanics: an element-wise step runs once per element per subscription, so a hidden write multiplies with every run of the chain, and a failure part-way leaves the earlier writes in place.
Diagnose it in a real system: duplicated audit rows traced back to a second subscriber, a transform whose latency is really an external store's, and recovery that must reconcile effects a failed run already applied.
The standard worth setting is where effects may appear in a pipeline at all, and the expectation that any effect which can be attempted twice is repeat-safe by design rather than by luck of how the run ended.
## What a transforming step promises An element-wise transforming step in a stream pipeline promises one thing: given an element, produce the corresponding output element. Everyone who reads the chain — and every tool that reasons about it — takes that promise at face value. `row -> payLineFor(row)` looks like arithmetic. When the same step also writes an audit row, the promise is broken silently: the code still compiles, the tests still pass, and the pipeline now has a consequence outside itself that its own shape does not advertise. ## What the hidden write actually costs 1. **Invisibility.** The chain reads as a computation. A reader looking for what touches the audit store will not look inside a mapping step, so the effect is discovered during an incident rather than during review. 2. **Multiplicity with the run.** An element-wise step runs once per element **per subscription**. A source that produces its values per subscriber will run the whole transform again for a second subscriber, and again for a re-subscription after a failure. The audit store gets one row per element per run, which nobody intended. 3. **Partial application on failure.** A failure signal travelling down the pipeline after 800 of 1,000 elements does not undo the 800 writes that already happened. The pipeline's output is abandoned; the effect is not. Any recovery has to reckon with the effects already applied. 4. **Coupling of failure modes and latency.** The transform now fails when the audit store fails, and takes as long as the audit store takes. A step that was a pure computation has acquired the availability and latency of an external system, and it did so inside the part of the chain least likely to be suspected. There is one more, subtler cost: the step is no longer safely **re-runnable**. Re-running a pure transform for one element to debug it is free; re-running this one writes a row. ## Three kinds of step, side by side | Step | Returns | Changes the sequence | Runs how often | Where a thrown error goes | |---|---|---|---|---| | Transforming step | the new element | yes, it replaces values | once per element, per subscription | becomes a failure signal downstream | | Observe-only hook | nothing | no, values pass through | once per element, per subscription | becomes a failure signal downstream | | Explicit effect step | the result of the effect | depends on how it is composed | once per element, per subscription | becomes a failure signal downstream | The table makes the honest point: moving the write does **not** change how often it runs. An observe-only hook sees exactly the elements the transform emitted, so the count is identical. What changes is that the effect is now **declared** — visible in the chain, movable, and removable without touching the transform — and the transform goes back to being a computation you can replay, test and reason about element by element. ## What an observe-only hook may and may not do An observe-only hook receives each element and returns nothing, so it cannot replace an element with another value and cannot drop one. It is not, however, harmless: an exception raised inside it still turns into a failure signal that ends the sequence, and slow work inside it still holds up the element that is passing. "Observe-only" is a statement about the *values*, not a promise that nothing can go wrong. ## Doing it in the payroll run - Keep `row -> payLineFor(row)` free of effects, so it can be run against a row anywhere, including in a test with no store attached. - Declare the audit write in its own step, placed where you actually want it — before the transform if you are auditing the input, after it if you are auditing the result. That placement is now a visible decision rather than an accident of where the code happened to sit. - If the effect must not repeat when the run is attempted again, make it repeat-safe on its own terms, keyed by something stable such as a run identifier plus the row identifier. Moving it out of the transform does not give you that for free. - If the write is itself asynchronous, an element-wise step is the wrong home for it: subscribing to inner work per element is a different shape of step and a different subject. ## What interviewers listen for A weak answer says "it is not pure" and stops. A strong answer is concrete: it names the per-subscription multiplication, the partial effects left behind by a mid-stream failure, and the borrowed latency and failure modes — and then does not overclaim, admitting that moving the write to a declared step improves clarity and separation but does not by itself reduce the number of writes or make them repeat-safe.
- Does moving the write into an observe-only step reduce how many audit rows are written?No. The hook sees exactly the elements the transform emitted, so the count per run is unchanged, and a second subscription still writes everything again. What the move buys is visibility and separation: the effect is a step you can place, remove or replace, and the transform becomes a computation you can safely re-run. Repeat-safety has to be built deliberately.
- The audit write must happen once per payroll run, not once per pay line. Where does it belong?Not in an element-wise step at all. Per-element steps run per element by definition, so a once-per-run effect belongs at a boundary of the run: triggered when the subscription starts, or on the completion signal at the end, or outside the pipeline entirely by the code that starts it.
- What happens to the audit rows already written when the pipeline fails half-way?They stay. A failure signal ends the sequence and abandons the remaining work, but it does not reverse effects that already happened outside the pipeline. Recovery therefore starts from a partially written store, which is why an effect that may be re-run needs to be repeat-safe or reconciled explicitly.
saying these in an interview costs you the question
- Treats a transforming step as a convenient place for logging or writes
- Assumes an effect inside a mapping step runs once overall, not once per subscription
- Believes an observe-only step can replace or drop the element it sees
- Expects a mid-stream failure to undo effects already applied
- Claims moving the write to its own step makes it repeat-safe by itself
- Ignores that the transform now inherits the store's latency and failure modes