A workflow scheduler is given one step per input file and there are 40,000 files - what breaks, and where does that fan-out belong?
answer
- count the launches, not the work
- a fixed cost per start, times 40,000
- bookkeeping per step, not per record
- boxes in the graph should not scale with data
basics
~20 sFixed per-step costs dominate: 40,000 program launches and 40,000 rows of scheduler bookkeeping for a few seconds of real work each, throttled by whatever concurrency limit the scheduler enforces. Fan-out over data belongs inside one submission, as parallelism.
solid answer
~50 sA workflow scheduler starts other programs in a required order and records whether each finished - its unit is a whole program, and every unit carries a fixed start-up and a durable record of its outcome. Multiply those fixed costs by 40,000 and they swamp the work: the graph takes longer to schedule than to compute, the scheduler's own bookkeeping becomes the bottleneck, the concurrency limit serialises what looked parallel, and the picture is unreadable to a human. A cluster execution engine's natural unit is the opposite: hand it the whole file list as one submission and the fan-out becomes parallelism inside that submission, where a piece of input costs no launch and no bookkeeping row. The rule: fan-out proportional to the *data* belongs inside one submission; fan-out proportional to distinct *programs* belongs to the scheduler.
go deeper
Recall the two units: a scheduler's unit is a whole program it starts, an engine's unit is a piece of input inside one submission. Then notice which one you are creating 40,000 of.
Explain the mechanics that multiply: fixed start-up per step, a durable bookkeeping record per step, a shared concurrency cap, and a graph no human can read. Then say the correct shape is one submission over the whole list.
Show that you would look for the general signature rather than this instance - a graph whose box count grows with the data - and price the collateral damage to other pipelines sharing the scheduler's capacity.
Make it a standard others can apply without you: the workflow graph's size is bounded by the architecture, never by the data, and data-sized fan-out is the engine's job by default.
## Two different units of work The error is a units error, and naming the units is most of the answer. A **workflow scheduler** - a system whose only job is to start other programs in a required order and record whether each one finished - has exactly one unit: a step, meaning a whole program it starts as its own process (sometimes a unit run inside a shared worker instead, but the accounting is the same). Every step carries: - a fixed start-up before any of your code runs; - a durable record of its existence and its outcome, written and updated by the scheduler; - a slot against whatever concurrency limit the scheduler enforces; - a box in a picture that a human is expected to read. A **cluster execution engine** - a system you hand a whole program to, which splits that program's work across many machines - has a different unit: one piece of input, handled by one worker thread inside one submission. A piece costs no program launch, no durable bookkeeping row and no diagram box. Its cost is scheduling a small unit of work onto a thread that is already running. ## What 40,000 steps actually does 1. **Fixed cost times 40,000.** Whatever a launch costs - and how much it costs depends entirely on how the machines are supplied - it is now the dominant term. A few seconds of real work per file against a fixed launch cost per file means most of the wall clock is spent starting things. 2. **The bookkeeping becomes the workload.** The scheduler writes and updates state for every step. Its own storage and its own loop were sized for graphs of tens or hundreds of steps, and 40,000 of them makes the scheduler the slow component in a pipeline that computes almost nothing. 3. **The concurrency limit serialises it.** A scheduler caps how many steps run at once, for the entire installation or per queue. That cap is now your parallelism, and it is shared with everybody else's pipelines - so this graph both runs slowly and starves its neighbours. 4. **The picture stops being a picture.** Nobody can look at 40,000 boxes. The one thing the workflow graph is genuinely good at - showing a human the shape of the pipeline - has been destroyed by using it as a loop. 5. **Per-file outcomes become 40,000 decisions.** Every step reports individually, so the pipeline's outcome is now an aggregate somebody has to compute and interpret. What the scheduler then does about a failed step is its own subject; the point here is that you created 40,000 of those events out of one logical computation. ## Where the fan-out belongs | the fan-out is over... | belongs to | why | |---|---|---| | input files, records, keys, time slices of one dataset | one submission to the engine | it is data parallelism; the engine's unit is a piece of input and it costs no launch | | distinct programs, different runtimes, different owners | the workflow graph | each really is a separate thing to start, and there are few of them | | a handful of independent datasets, each with its own pipeline | the workflow graph | tens of steps, not tens of thousands, and each box means something to a reader | The shape that works is one submission handed the *list* of 40,000 files - or the directory - so the engine reads them as one input and runs the same logic over every piece in parallel. How that input is then cut into pieces, and whether many tiny files should be coalesced first, is the engine's own splitting question and a different subject; what matters here is that the count stopped being 40,000 *programs*. ## The same error one level down Once you see it, the family is obvious: a step per customer, a step per hour of history, a step per row of a control table, a step that starts a program which processes exactly one record. Each of them hands per-record or per-item work to something that can only launch programs. The signature is always the same - **the number of boxes in the workflow graph is proportional to the size of the data.** A workflow graph should have the number of boxes a person could draw on a whiteboard; if it grows when the input grows, the fan-out is in the wrong system. There is a legitimate middle case worth conceding: a small, bounded fan-out over a handful of genuinely independent datasets - five sources, five steps - is fine and often clearer than one submission that reads all five. The line is not 'never fan out in the scheduler'. It is that the fan-out must be bounded by something structural, not by the data's cardinality.
- Would raising the scheduler's concurrency limit fix it?No, it moves the pain. The fixed launch cost per file is still paid 40,000 times, the scheduler's own bookkeeping grows the same way, and a higher cap means this graph now crowds out every other pipeline sharing the installation. The count of steps is the defect, not the throttle on it.
- When is one step per input actually the right shape?When the fan-out is small and structural rather than data-sized - five distinct sources with five genuinely different programs, say. The test is whether the number of boxes grows when the input grows. If it does, the loop is in the wrong system; if it is fixed by the architecture, the boxes are meaningful.
- What is the equivalent error with a managed query service?Starting one declarative statement per file against a service that owns its own storage, plan and capacity. Each statement carries its own fixed overhead and the service plans each one in isolation, so the same launch-cost multiplication applies - the fix is again one statement over the whole set rather than a loop that issues many.
saying these in an interview costs you the question
- Thinks 40,000 steps is fine because the scheduler runs steps in parallel.
- Treats per-step start-up and bookkeeping as free because each step is small.
- Proposes raising the concurrency limit instead of removing the fan-out.
- Assumes a workflow graph should mirror the data, one box per file.
- Believes an engine cannot be handed a whole list of files as one input.
- Says any fan-out in a scheduler is wrong, including five independent sources.