skip to content

Two steps a workflow scheduler starts in order pass a 400 GB intermediate through object storage - what does that boundary cost?

level: middleimportance: should knowfreq 58%

answer

  1. the arrow carries no bytes
  2. write it, read it, start again
  3. count the name and its lifetime too
  4. what crosses must be materialised

basics

~20 s

A full write and a full read of 400 GB, plus a second program's start-up and a path someone must name, encode and eventually delete. An edge between two steps a scheduler starts carries nothing, so whatever crosses it has to be materialised.

solid answer

~50 s

A workflow scheduler only starts other programs in a required order and records whether each finished; its edges move no records. So the 400 GB has to land in storage every machine can read, in an agreed encoding, under a name the downstream step is told - that is a full write and a full read of every byte, the start-up of a second program, and a path with an owner and a lifetime. Fold the two steps into one submission - one whole program handed to the cluster and run end to end - and the engine owns the handoff instead: depending on the design it may keep it in memory, run the two steps as one pass, or still write it out. You are buying the engine's *choice*, not a promise that nothing is written. The boundary is still right when that intermediate is a dataset other consumers read anyway.

go deeper

for a junior

Recall that an edge between two scheduler-started steps moves nothing, so anything one produces for the next must be written somewhere and read back. That alone explains where the 400 GB is going.

for a middle

Itemise the cost rather than calling it slow: bytes out and in, encoding twice, a second program's start-up, a named artefact with an owner and a lifetime, and no rewriting across the boundary.

for a senior

Demonstrate the judgment call - ask who else reads that intermediate. If the answer is nobody, the boundary is an artefact of where the picture was drawn, and the repair is one submission, not faster storage.

for a principal

Set the rule of thumb others apply: a boundary that materialises 400 GB must be justified by a consumer, a different program or an external wait, and someone must own that artefact's schema, cost and deletion.

## What the edge is not doing A **workflow scheduler** is a system whose only job is to start other programs in a required order and record whether each one finished; it moves no records itself. So an edge between two steps it starts is an ordering claim and nothing more. The 400 GB is not travelling along the arrow. It is travelling because somebody wrote it to **object storage** - storage every machine in the cluster can read - and told the second program where to look. That is a perfectly legitimate thing to do. It is also never free, and the first job in an interview answer is to itemise the bill rather than say 'it is slow'. ## The bill, itemised 1. **Bytes out and bytes in.** 400 GB is written once and read once, over the network in both directions, and neither half was in the logical computation you wanted. 2. **Encoding both ways.** The intermediate has to become a format on the way out and be parsed on the way in. Whatever the format's cost per record is, you pay it twice, and a format chosen for convenience rather than for bulk can cost more than the transfer. 3. **A second start-up.** The downstream step is a fresh program: asking for machines, getting them, loading code. How large that fixed delay is depends entirely on how the machines are supplied, which is a different subject - but it is never zero and it is paid per boundary. 4. **A named, owned artefact.** Somebody chose the path, somebody owns the schema at that path, somebody must decide what happens when a rerun writes it again, and somebody pays for it until it is deleted. This cost is the one that is always forgotten and the one that is still there in a year. 5. **No rewriting across the boundary.** On engines that rewrite a declarative program into a physical plan, two steps in one submission can be reorganised together - a filter pushed back before the expensive part, two passes collapsed into one. Across a scheduler edge nothing can see both sides, so none of that is available. Note the hedge: engines that execute per-record functions almost literally, and the oldest disk-to-disk designs, would not have rewritten much anyway. ## What folding the boundary away actually buys Put both steps in one submission and the engine owns the handoff. What that means concretely differs by design, and saying so is the honest answer: - some engines keep the intermediate in memory and run the two steps as a single pass over the records; - some write it to the producing worker's own local disk, which is still a write but a local one, without a name, an owner or a cleanup story; - the oldest disk-to-disk designs write each round's output to durable storage anyway - so on those, folding the boundary away saves the second start-up and the naming, but not the bytes. **You are buying the engine's choice, not a guarantee.** The reliable wins are the second start-up, the external name and its lifetime, and the possibility of rewriting across the boundary. The bytes may or may not survive the fold. ## When the boundary is right anyway | the intermediate is... | verdict | |---|---| | a published dataset other teams or tools already read | keep the boundary - you were writing it regardless, so the write is not a new cost | | consumed only by the very next step and by nothing else | fold it in - it is an engine edge wearing a scheduler's clothes | | produced by one program and consumed by a different one (different runtime, different language) | keep the boundary - no single submission can hold both | | the point where a human or an external arrival gates the pipeline | keep the boundary - what is waited on is not a record | ## The smell The tell that a boundary is accidental rather than chosen is that nobody can say who else reads the intermediate, and the answer turns out to be 'nobody'. That usually happens because the whole pipeline was drawn in the scheduler - it is where the team looks, where the picture lives, so every stage of the computation became a step in it, and every arrow between stages silently became 400 GB through storage. The repair is not to make the write faster. It is to notice that those two stages have a data dependency and no reason to be separate programs, and to let one submission carry the records between them. The converse error exists too and is worth naming: a genuinely shared, consumed-by-many dataset that somebody hides inside a submission to save a write, so nothing downstream can read it without rerunning the whole computation. The question is never 'materialise or not'. It is 'does anything but the next step need this'.

  • If you fold the two steps into one submission, is the 400 GB write guaranteed to disappear?
    No. You reliably remove the second start-up, the externally named artefact and its lifetime, and you allow rewriting across the boundary on engines that rewrite. The bytes depend on the design: some engines hold the handoff in memory, some spill it to a worker's local disk, and the oldest disk-to-disk models write each round out regardless.
  • What makes a 400 GB intermediate worth keeping on a scheduler edge?
    That something other than the next step reads it. If it is a dataset other teams query, another pipeline consumes, or an audit needs, you were going to write it anyway - the boundary costs you only the second start-up. If the sole consumer is the next step, the write exists because of where the picture was drawn.
  • Does splitting the pipeline into more scheduler steps ever make it faster?
    Not by itself - each extra boundary adds a materialisation and a start-up. It can help indirectly when the two sides want genuinely different capacity, so a short expensive stage stops holding a large allocation open for a long cheap one. That is a resource argument, not a throughput one.

saying these in an interview costs you the question

  • Says the scheduler passes the intermediate from one step to the next.
  • Assumes a large intermediate is fine between scheduler steps because it always has been.
  • Counts only the bytes and forgets the second start-up and the named path.
  • Claims folding the steps together always removes the write entirely.
  • Thinks a faster storage tier fixes a boundary that should not exist.
  • Treats every materialised intermediate as waste, including ones other consumers read.