A finite job's 800 worker threads write rows straight into one live operational database at the final step, and the database stalls. What has to change?
answer
- the slow side is not the cluster
- the destination sets the ceiling
- rows per round trip
- cap writers, not pieces
- a measured rate, not a headline figure
basics
~20 sThe limit is now the destination's, not the cluster's. Three things change: rows go in batches rather than one request each, the number of simultaneous writers is capped independently of the piece count, and the run holds a rate ceiling the store was measured to survive.
solid answer
~40 sA file destination absorbs parallel writers because each output is independent. A live store does not: connections, locks, per-row index maintenance and replication are one shared resource, and 800 writers arriving at once queue against each other and against the store's other users. So you stop sizing the write from the cluster. Batch rows so per-request overhead is paid once per few hundred rows; cap simultaneous writers at a number the destination was measured to take, which has nothing to do with the piece count that was sized for compute; and hold a rate ceiling for the run as a whole so a finite job's final step is a plateau rather than a burst. Where the store has a bulk path, writing files and loading them in one step beats all three.
go deeper
Recall that a live store is not a directory: many writers hitting one database at once make it slower, not faster, and the job must be held back rather than widened.
Explain the three levers and why each works: per-request overhead amortised by batching, a writer cap decoupled from the piece count, and a rate ceiling that flattens a finite job's burst.
Demonstrate the diagnosis — queueing at the store, latency for its other users, retries adding load — and know that a bulk load path beats tuning writers at all.
The angle is the shared destination: who owns that store's headroom, how the safe rate is discovered, and what it costs the organisation when one run can degrade everyone else's reads.
## The ceiling moved off the cluster When the destination is files, parallel writers are close to free: each worker thread writes its own output, the outputs do not contend, and more writers mostly means more throughput. A live operational store inverts every part of that. It is one shared system with a connection budget, locks, index maintenance on every row, a durability log and replication to followers — and all of it is shared with whatever else uses that store. The job's parallelism is now aimed at a resource that does not grow when the cluster does. The failure has a recognisable shape: - connections are queued or refused once the writers outnumber the store's budget; - per-request latency climbs, which makes each writer hold its connection longer, which deepens the queue; - lock and log contention makes throughput fall as concurrency rises, past a point; - the store's **other** users see the latency too, which is usually the actual incident; - the job's own retries add load exactly when the store has least to give. ## The three levers 1. **Batch.** Send a few hundred rows per request instead of one. Per-request overhead — the round trip, parsing, transaction bookkeeping — is then paid once per batch rather than once per row, and that overhead usually dominates at small sizes. It is not monotonic: a very large batch is one long transaction holding locks, and a single failure re-sends the whole thing, so there is a size beyond which it gets worse. 2. **Bounded concurrency.** Cap the number of writers that may be in flight, independently of the piece count. This is the lever candidates miss, because the two numbers feel like one. The piece count was sized for compute — bytes per piece against available worker threads — and nothing about that reasoning knows the destination's connection budget. Eight hundred pieces can be written by sixteen writers, in turn. 3. **A rate ceiling.** Rows or requests per second for the run as a whole, so that a finite job — whose pieces all reach the final step at roughly the same time — presents a plateau instead of a burst. The number comes from measuring the store, not from a headline figure. ## File destination against live store | | files | live operational store | |---|---|---| | what scales with writers | throughput, nearly linearly | nothing, past the store's own budget | | the binding limit | the piece count you chose | connections, locks, log and replication | | the symptom of too much | a huge number of small files | rising latency for every user of the store | | the lever | the piece count at the final step | batch size, writer cap, rate ceiling | ## The job's shape changes the shape of the pressure A finite job concentrates the write: every piece reaches the final step in roughly the same window, so the store sees a spike whose height is set by the piece count. A **continuous job** — a run over an input with no end, whose author declares an operator width that stands until restart — presents sustained pressure instead. There the store being slow does not usually topple anything: on runtimes where a consumer that cannot accept records makes its upstream wait, that wait propagates back towards the source, which is **backpressure**, and the visible symptom is a growing backlog and rising lag rather than a broken store. Some runtimes expose that signal directly; others simply appear to stall, which is why the same destination problem is diagnosed differently in the two modes. ## The answer above the three levers Where the destination offers a bulk load path, the strongest answer is to stop writing rows from the cluster at all: have the job write files, then load them in one step. Thousands of small transactions become one load the store can plan, and the cluster's parallelism goes back to being free. The costs are real and worth naming — a staging location, a second step, and the fact that ordering that step relative to the job belongs to `workflow scheduling` rather than to the job itself. Two boundaries are worth knowing and not stepping over in the answer. What the destination ends up seeing when a failed attempt is retried — duplicated or not — is owned by `The Visible Effect`. Running out of room inside a worker while writers hold buffers is owned by `Memory, Spill and Caching`. This subject is narrower: when the destination is a live store rather than a directory, the write is limited by the destination, and the job must be shaped to that limit instead of to the cluster's.
- When is writing files and loading them separately the better answer?Whenever the store offers a bulk path. It converts thousands of small transactions into one load the store can plan, and the cluster's parallelism becomes free again. The costs are a staging location and a second step, and ordering that step relative to the job belongs to `workflow scheduling`.
- Does a continuous job hit the same wall?It hits a different one. The pressure is sustained rather than bursty, and on runtimes where a consumer that cannot accept records makes its upstream wait, that wait travels back towards the source — backpressure. The symptom is a growing backlog and rising lag, not a toppled store.
- Why does batch size have an upper limit as well as a lower one?A large batch is one long transaction: it holds locks longer, grows the store's durability log in one push, and on failure the whole batch is re-sent. Per-request overhead falls sharply for the first few hundred rows and then flattens, so past some size you pay the risk without buying throughput.
saying these in an interview costs you the question
- Adds worker threads so the write finishes sooner.
- Says throughput scales linearly with the number of simultaneous writers.
- Sends one request per row and blames the network for the latency.
- Lets writer concurrency follow the piece count, which was sized for compute.
- Treats a shared operational store as having no ceiling worth measuring.
- Retries the failed writers immediately, while the store is still saturated.