skip to content

Mapping each queued item to its own asynchronous call yields a source of sources: what does a flattening step then do?

level: middleimportance: must knowfreq 64%

answer

  1. one level too many
  2. a source per element
  3. nothing runs until subscribed
  4. subscribe, then re-emit as one
  5. strategy decides order and concurrency

basics

~20 s

A flattening step subscribes to each inner source the mapping produced and re-emits their elements onto one outer sequence. Its strategy also fixes the ordering of results, how many inner calls run at once, and whether superseded ones are cancelled.

solid answer

~50 s

Mapping is one-to-one: element in, value out, and the shape of the sequence is unchanged. When the value is itself an asynchronous call, the mapping produces a sequence whose elements are *sources*, so a subscriber downstream would receive sources rather than results. The flattening step removes that extra level. It subscribes to each inner source - which is what actually issues the call, since a cold source does nothing until subscribed - re-emits every element the inner source produces onto a single outer sequence, and completes only once the outer source and every inner source it started have completed. The interesting part is not the flattening but the strategy: whether inner sources are subscribed one at a time, several at a time up to a bound, or one at a time with the previous cancelled. That choice decides result ordering, the load placed on the dependency, and whether work already in flight is abandoned.

code

pseudocode · 9 lines
pseudocode
// mapping alone: each element becomes a source, not a verdict
sources = queue.map(function(item) { return classify(item) })
// elements of `sources` are SOURCES; nothing has been sent yet

// flattening: subscribe to each inner source, re-emit what it produces
verdicts = flatten(sources, strategy = SEQUENTIAL)

verdicts.subscribe(function(verdict) { record(verdict) })
// without the flatten, record() would be handed an unsubscribed source

go deeper

for a junior

Recall that mapping an element to another asynchronous call produces a source of sources, and that a separate flattening step is what turns those back into one sequence of results.

for a middle

Explain the three jobs the flattening step performs - subscribing to each inner source, re-emitting its elements into one sequence, and deciding when the whole sequence is complete - and name the strategies it can use.

for a senior

Show that you treat the strategy as a production decision: inner calls in flight are load on a dependency, and the ordering guarantee is a correctness property somebody downstream may already be relying on.

for a principal

The angle to own is the default. When nobody states a strategy the pipeline still has one, and an unbounded or order-losing default quietly becomes an unwritten contract that other teams build against.

## One level too many A mapping step is **one-to-one**: it takes an element, returns a value, and the sequence keeps its shape. That holds only while the value is a value. The moment the mapping returns **another asynchronous call**, the shape changes, because the call is not a result - it is a *source* that will push its own result later. The mapping now produces a sequence whose elements are sources, and anything subscribing downstream receives sources, not answers. Take a moderation queue where every arriving item must be sent to an external classifier. Mapping each item to its classification call gives you a source of calls. Nothing has been classified. Nothing has even been sent. Removing that extra level, and issuing the calls in the process, is the whole job of a **flattening step**. ## What the flattening step actually does Three jobs, and only the first is obvious. 1. **It subscribes.** A **cold** inner source defines work per subscriber and does nothing until something subscribes to it, so the flattening step is what issues the call. This is why an element that the strategy decides to discard costs nothing downstream - its call is never made at all. 2. **It re-emits.** Whatever an inner source produces is republished onto a single outer sequence, element by element as it arrives. It is not collected into a list, not paired with the element that produced it, not deduplicated. 3. **It decides completion.** The flattened sequence completes when the outer source has completed **and** every inner source still alive has completed. Completing as soon as the last element was mapped would truncate results that were still in flight. The second and third points are what separate a flattening step from a plain transformation: it owns subscriptions, so it owns lifetime. ## The strategy is the real content Every flattening step answers two questions, and the answers are what interviewers are after: **how many inner sources may be live at once**, and **what happens to work already in flight when a new element arrives**. | strategy | inner sources live at once | order of results | what it gives up | |---|---|---|---| | sequential concatenation | one | element order | throughput: one call per round trip | | bounded concurrent merge | up to the bound | completion order | element ordering | | switch to newest | at most one, previous cancelled | only the latest element's result | results of superseded elements | | ignore while busy | one | element order of the kept elements | elements that arrived during a call | None of these is a default worth accepting without thought. In the moderation console, the queue itself wants throughput and probably tolerates interleaving, so bounded merging fits. The detail pane beside it, which loads whichever item is highlighted, wants the newest and nothing else, so switching fits. A button that triggers one expensive re-classification wants ignore-while-busy, so a double press costs one call. The same pipeline can legitimately use three different strategies in three places. ## Symptoms when the flattening is wrong or missing - **Results arrive as sources.** In a typed pipeline this is a compile error; in an untyped one it renders as something meaningless. The mapping was never flattened. - **The sequence completes before any result arrives.** Something waited on the outer source only and ignored the inner ones. - **A burst of arrivals becomes an identical burst downstream.** The merge has no useful bound, so every queued element has a call in flight simultaneously. - **Results are recorded in an order nobody expected.** A concurrent merge was chosen where element order was part of the requirement, and the defect appears only when inner latencies differ - which is rarely true in a test with stubbed calls. ## Why this is asked The map-versus-flatten distinction is a screening question because it separates someone who has assembled a pipeline from someone who has only read about one. The real interview, though, is the follow-up: *which strategy, and why that one here*. An answer that stops at "it flattens the nested sources" is incomplete, because the flattening is mechanical and the strategy is the design decision. A strong answer names the two axes - concurrency and what happens to in-flight work - and then picks a strategy for the specific case in front of it, saying out loud what that choice costs. One last property worth internalising: because the flattening step owns the inner subscriptions, cancelling the flattened sequence propagates down into them. Whether that actually stops the outbound calls depends on whether each inner source honours cancellation, which is a property of how that source was built rather than of the flattening strategy.

  • When does the flattened outer sequence complete?
    Only when the outer source has signalled completion and every inner source still alive has completed too. A strategy that cancels superseded inner sources still waits for the survivor. Completing as soon as the last element was mapped would cut off results that had not yet arrived.
  • What happens to an inner source the strategy never subscribes to?
    Nothing runs. A cold source does no work until a subscription starts it, so a strategy that discards an element while a call is in flight costs the dependency nothing - the call is never issued. An inner source that was already running independently is the exception, since its work proceeds regardless.

saying these in an interview costs you the question

  • Thinks a plain mapping step waits for the inner call to finish
  • Calls flattening a formatting convenience rather than a subscription policy
  • Assumes every flattening strategy preserves the order elements arrived in
  • Believes a cold inner source starts work before anything subscribes
  • Forgets the outer sequence waits for inner sources to complete