skip to content

Your program contains three wide steps. How does that number read differently on a runtime that schedules in units versus one where all operators run at once?

level: seniorimportance: should knowfreq 42%

answer

  1. one classification, two predictions
  2. boundaries against hops
  3. four scheduled units from three regroupings
  4. no boundary means no re-planning
  5. key choice decides the worker either way

basics

~20 s

On a runtime that schedules the work between regroupings, three wide steps mean four separately scheduled units of steps. On a runtime where every operator runs at once, they mean three routing hops each record may make in flight.

solid answer

~60 s

The classification is the same on both; what it predicts is not. Where the runtime schedules work in **step groups** — the run of steps between two wide steps, scheduled and retried as one unit — three wide steps cut the program into four groups. That number predicts how many scheduling units exist, how many handovers are materialised, where progress becomes visible, and how much is re-run when a worker is lost. Where every operator is live at once and records are pushed downstream as they are produced, a wide step is instead a routing rule on an edge: each record goes to the worker that owns its key. Three of them means up to three network hops per record, no boundary to re-plan at, and a slow consumer slowing its producer rather than a group finishing — the effect called backpressure. What is identical is that a wide step is where your key choice decides which worker gets each record, so an uneven key distribution hurts on both.

go deeper

for a junior

The takeaway is that counting wide steps is useful everywhere, but what the count means depends on the runtime: in one style they cut the program into separately scheduled units, in the other they are points where records are routed by key.

for a middle

Explain both readings and their consequences — scheduling units, retry granularity and observable progress on one side; network hops per record and upstream slowdown on the other. Avoid stating either as the universal model.

for a senior

Show judgment on the system in front of you: name which style it is, say what that makes the unit of retry and whether re-planning is even possible, and handle the case of a never-ending workload executed as a succession of small finite runs.

for a principal

The portable part of a design is the key choices and the number of regroupings; the rest is runtime-specific. Standardising on that distinction is what lets an organisation move logic between execution styles without re-deriving its cost model.

## Same classification, two different predictions A **wide step** is one a worker cannot finish from the records it already holds, because producing one output needs records currently on other workers. That definition is common to this whole class of engine. What a runtime *does* at such a step is not common at all, and the count of wide steps in your program therefore answers a different question depending on where it runs. ## Where work is scheduled between regroupings Some runtimes cut the program at every wide step into **step groups**: the run of steps between two wide steps, scheduled and retried as one unit. Three wide steps give four groups. On such a runtime the count tells you: - **How many scheduling units exist**, and therefore how coarse the unit of retry is — losing a worker generally costs you work at the granularity of the group, not of the individual operator. - **How many handovers are produced**, since the boundary between two groups is where one group's output becomes the next group's input rather than staying in a single pass. - **Where progress is observable.** Completed groups are the natural checkpoints of a run's progress; a long program with one wide step gives you far less to look at than one with five. - **Where a runtime may re-plan.** Some runtimes measure what a finished group actually produced and adjust the next one; that opportunity exists only where such a boundary exists. The oldest model in this family is the extreme case: two halves with a single regrouping between them, everything materialised to disk in the middle. A program with three logical wide steps does not fit in one such job at all — it becomes a chain of them, which is precisely why the model was superseded for multi-step work. The mechanics of what happens at a boundary — who must wait for whom, what is written on one side and fetched on the other, and what that costs in bytes — is a subject of its own and not this one. Here the boundary matters only as the thing that makes a wide step a *structural* fact about the program. ## Where every operator runs at once On a record-at-a-time runtime the whole graph is live: every operator has running instances and records are pushed to the next operator as they are produced, with no materialised handover in between. A wide step is then not a split in time at all. It is a **routing rule on an edge** — each record is sent to the instance that owns its key — and the count tells you something else entirely: - **How many network hops a record may make** between entering and leaving the job, which is a latency statement rather than a scheduling one. - **How congestion propagates.** With no boundary to finish, a consumer that cannot keep up slows its producer instead, and that slowdown travels upstream to the source. The term of art is backpressure: a downstream operator unable to keep up makes the upstream one slow down rather than buffering without bound or dropping records. - **Why re-planning is largely unavailable.** There is no completed unit whose real output could be measured, so adjusting the shape of the run typically means stopping and restarting it. - **Why the width is a declared number.** Nothing about an endless input can be measured in advance, so the author states how many instances of each operator run, and that number stands until the job is restarted — which means changing it is an operation, not a tuning knob. One more source of confusion is worth naming: part of this market runs continuous workloads as a rapid succession of small finite jobs. Such a job is *continuous* in what it consumes and *scheduled in units* in how it executes, so it sits in the first column of the table below even though the workload never ends. ## Side by side | | scheduled between regroupings | every operator live at once | |---|---|---| | what a wide step is | the boundary that ends one step group | a routing rule on an edge | | three wide steps predict | four scheduled units | three hops per record | | unit of retry | typically the group | typically the whole job or its recovery point | | re-planning at a boundary | possible on some runtimes | largely unavailable | | congestion shows as | a group taking longer to finish | upstream slowdown, i.e. backpressure | | changing the width | can often be decided per run | usually fixed until restart | ## What does not change On both, a wide step is the moment your key choice decides which worker receives each record. If one key value carries far more records than the rest, that worker gets disproportionate work — the condition called skew, which is a subject of its own — and it hurts in both columns, as one long-running piece in the first and as one permanently hot instance in the second. Also on both: the number of wide steps is the first-order description of the program's cost, and reducing it is the rewrite with the largest effect. ## What an interviewer is listening for Not a recitation of one runtime's execution model. They want the candidate to notice that "how many wide steps" is a question with two useful answers, to say which one applies to the system in front of them, and to avoid the confident universal — "the runtime writes the intermediate result and the next group reads it" is true of one lineage and false of the other, and a candidate who states it flatly has revealed the only engine they have used.

  • Why is runtime re-planning largely unavailable where every operator runs at once?
    Re-planning needs a completed unit whose actual output can be measured and fed into the next decision. Where the whole graph is live and records are pushed as produced, nothing ever completes in that sense, so the shape is what the author declared; changing it generally means stopping the job and starting it again with a different width.
  • A workload never ends but executes as a rapid succession of small finite jobs. Which column is it in?
    The scheduled-in-units one. It is continuous in what it consumes but each small run is scheduled, executed and completed like a finite job, so wide steps still act as boundaries within a run. The consequence is a latency floor set by the size of those runs, which a record-at-a-time runtime does not have.
  • Does the number of wide steps predict recovery cost the same way on both?
    No. Where work is scheduled in units, losing a worker usually costs re-running at the granularity of a group, so more boundaries can mean less re-done work. Where everything runs at once there is no such granularity: recovery works from a saved recovery point or from replaying the input, which is a different subject entirely.
  • If the same logic must run in both styles, what should you hold constant?
    The key choices and the number of distinct regroupings, because those are the parts of the design that mean something in both. What should not be carried across is any assumption about materialisation, about a boundary existing to re-plan at, or about being able to change the width mid-flight.

saying these in an interview costs you the question

  • Says every wide step writes intermediate results to disk and the next step reads them.
  • Assumes a boundary exists to re-plan at on every runtime.
  • Describes the width as a knob that can be changed mid-flight everywhere.
  • Thinks a never-ending workload cannot be executed as scheduled units.
  • Claims wide steps cost nothing where operators all run at once.