Why choose a send-driven generator pipeline over a pull-based one for an ETL export?
answer
- Who owns the driving loop?
- Callback sources cannot be looped over
- Each stage holds the next one
- A plain call, so no queue in between
- One raise closes the whole chain
basics
~20 sChoose push when the rows arrive from something you do not drive, or when one row must reach several sinks. You pay for it: no backpressure, manual wiring, and one failing stage kills the whole chain.
solid answer
~50 sA pull pipeline needs your code to own the loop, so it only works when you can ask for the next row. A send-driven pipeline inverts that: each stage is a primed consumer generator holding a reference to the next stage and calling `send()` on it, so whoever *has* the rows drives - a cursor callback, a change feed, a parser. It also buys stage-local state for free, which is how a warehouse writer batches rows and flushes every N without a class, and it fans out when one row must reach several sinks. The costs are real: `send()` is a plain synchronous call, so there is no backpressure; an exception in any stage propagates out of the sender's `send()` and leaves both generators closed; and shutdown ordering is yours. For anything new, an async consumer draining a queue is usually the better shape.
code
python · 29 linesfrom functools import wraps
def primed(func):
@wraps(func)
def start(*args, **kwargs):
gen = func(*args, **kwargs)
gen.send(None)
return gen
return start
@primed
def writer(batch_size):
batch = []
while True:
row = yield
batch.append(row)
if len(batch) == batch_size:
print("flush", batch)
batch = []
@primed
def redact(target):
while True:
row = yield
target.send(row | {"email": "***"})
stage = redact(writer(2))
for i in range(4):
stage.send({"id": i, "email": f"u{i}@example.com"})go deeper
Focus on the shape rather than the tradeoffs: each stage is a consumer generator that holds the next stage and calls send() on it, so data is pushed forward instead of being pulled by a loop at the end.
Explain the mechanics you can demonstrate: stages are primed at construction, the successor arrives as a constructor argument, batching state lives in the suspended frame, and send() is a plain synchronous call with no queue in between.
Bring the operational judgement. Name the failure semantics - one raise unwinds and closes every stage, buffered rows are lost, restart means rebuild - and say when you would refuse the pattern in favour of an explicit writer object or an async consumer.
Weigh it as a maintenance decision. Across an 11-person team a clever chain of suspended frames is read far more often than it is written, so decide whether the push shape is forced by the source at all, and set the boring default for everything it is not.
### The shape In a pull pipeline each stage is a generator that iterates the stage before it, and the final `for` loop at the end drags rows through the whole chain. That requires your code to own the iteration. In a push pipeline the arrow reverses. Each stage is a consumer generator - `while True: row = yield` - primed at construction and holding a reference to the *next* stage. It does its work and calls `target.send(row)`. There is no driving loop inside the pipeline at all; the pipeline is a chain of parked frames, and whoever holds the rows kicks the first one. ### When the inversion is genuinely worth it **The source drives, not you.** If rows for a warehouse export arrive through a callback, a change feed, a SAX-style parser or any API that calls *you*, there is nothing to write a `for` loop over. You can only be pushed to, and a push pipeline matches that without an intermediate thread or queue. **Stage-local state across rows.** A writer that accumulates a batch and flushes every N rows keeps `batch` as an ordinary local in a suspended frame. No class, no attribute, no reset logic - the linear narrative of 'take a row, append it, flush when full' stays readable as straight-line code. **Fan-out.** One stage can hold several targets and send each row to all of them - the warehouse writer, an audit log, a metrics counter - which is awkward to express with a pull pipeline, where a single consumer owns the iteration. ### What it costs **No backpressure.** `send()` is an ordinary synchronous call. Every stage runs to completion inside the sender's call, so the chain is only as fast as its slowest stage and there is no queue, no buffering, and no way for a downstream sink to say 'slow down' other than blocking. A pipeline described as 'streaming' can quietly mean 'the source blocks on the warehouse round trip'. **Failure closes the whole chain.** If the last stage raises, the exception propagates out of the `target.send(row)` call inside the stage upstream of it, escaping *that* generator's frame too, and so on up to the code that pushed the row. When the dust settles every generator in the chain has been left in a closed state - you can confirm this with `inspect.getgeneratorstate` - so there is nothing to retry into. Recovery means rebuilding the pipeline, and any batch buffered in a stage is gone with it. **Shutdown is yours.** Rows sitting in a partly filled batch are not flushed by anything automatic. Ending the export cleanly means an explicit finish protocol that reaches every stage, in order, and it is easy to get right during a demo and wrong during an exception. **Debuggability.** A traceback through four suspended generator frames is harder to read than four method calls, and the pipeline's real structure exists only at runtime, in whatever object was passed to whatever constructor. ### Wiring, and one trap worth naming Each stage needs its successor. Give it as a constructor argument - `redact(writer(batch_size))` - and never let a stage reach for its successor by importing the module that defines it. Stages that import each other at module top level to look up 'the next one' turn a linear data flow into a cycle in the import graph, and you find out with a circular import at startup, usually after the wiring already looks correct. Passing the target as an argument keeps composition a runtime concern and keeps every stage independently testable: feed it a recording stub and assert what it forwarded. ### The honest recommendation For an export nobody has written yet, reach for the boring version first: a source generator and an explicit writer object with `write()` and `flush()`, or, if the source is genuinely push-shaped and concurrency is involved, an `async def` consumer draining an `asyncio.Queue`, where backpressure is a bounded queue rather than a hope. Use the send-driven chain when the push shape is forced on you and the stage-local state really does read better as a suspended frame - and say so explicitly, because the pattern's cleverness is exactly what makes it expensive for the next person to modify.
- How does an exception in the last stage of a push pipeline surface to the code feeding the first?It travels back along the same call chain the rows travelled forward on. The raise escapes the last stage's frame, comes out of the `target.send(row)` call in the stage above it, escapes that frame too, and so on until it reaches the pusher. Because the exception left every frame, all of those generators end up closed, so the pipeline cannot be resumed - you must rebuild it, and anything buffered mid-stage is lost.
- Does a send-driven pipeline give you backpressure?No. `send()` is a synchronous call, so a row is processed all the way to the sink inside the source's own call - there is no queue between stages and nothing to fill up or block on. What looks like streaming is really the source running at the speed of the slowest stage. Real backpressure needs a bounded buffer: a bounded queue between threads, or an `asyncio.Queue` with a maxsize in an async design.
- How would you unit-test a single stage of such a pipeline?Give it a recording stub as its target - anything with a `send()` method that appends to a list - then push a few rows in and assert on what was forwarded. That works only if the stage receives its successor as an argument; a stage that imports the next stage itself cannot be isolated, which is one more reason to wire the chain at construction time.
saying these in an interview costs you the question
- Claims a send-driven pipeline provides backpressure
- Thinks the chain can be reused after a stage raises
- Wires stages by importing each other at module top level
- Assumes buffered batches are flushed automatically at the end
- Calls it concurrent when every stage runs inside the sender's call
- Chooses the pattern for speed rather than for a push-shaped source