In a job's printed plan, how do you locate the points where records cross the network, and what does a long unbroken chain mean?
answer
- find the cuts in the listing
- which operations force records to meet
- unbroken means one piece in, one out
- compare crossing count to key changes
- a crossing is not always materialised
basics
~20 sLook for where the listing breaks: a crossing is a point where records must be sent between workers, and most engines cut the printed steps into runs around it. An unbroken chain means every step builds each output piece from one input piece.
solid answer
~50 sA **network crossing** is a point where records must be sent between workers because the next step needs records another worker is holding — a shuffle. In the printed text you find them two ways: by the explicit break most engines print between runs of steps, and by the operations that force one, since a grouping by key, a join on a key, a global ordering or a deduplication all require records that must meet to arrive in one place. A run of steps between two crossings — many engines call it a stage — is a group the engine can execute without moving anything. So a long chain with no break tells you every step in it builds each output piece from exactly one input piece, a narrow dependency, and the whole chain can be walked over each piece in a single pass. Compare the crossing count against how many times your program genuinely changes grouping key.
go deeper
Know that a crossing is a point where records must be sent between workers, and that the printed listing is normally cut into runs around those points. Recognise grouping and joining as the usual causes.
Find the crossings both from the printed structure and from the operations that force them, and explain that an unbroken run means every step builds each output piece from one input piece.
Reconcile the printed crossing count with the key changes the program actually requires, and name where the extra ones came from. State plainly which of your claims about a crossing hold only for a materialising engine model.
Decide whether crossing counts are a metric worth standardising across a platform when the engines in use build crossings differently, and what a team should measure instead of reading shape.
## What you are looking for A **network crossing** is a point in the plan where records have to be sent from one worker to another, because the step above needs records that a different worker is holding. The industry term for that movement is a *shuffle*. It is the single most consequential thing the printed text tells you about a distributed job, because it is the one part of the work whose cost has nothing to do with how clever the computation is. A **phase** is a run of steps between two crossings that the engine can execute without moving any records at all. Many engines call such a run a *stage*, and most print the listing already cut into them. Locating the crossings and locating the phases is the same act seen from two sides: the crossings are the cuts, the phases are what lies between them. ## Two ways to find them 1. **Structurally, from the printed text.** Engines mark the cut. Depending on the engine you may see the listing grouped into numbered sections, a named node standing between two runs of steps, or an indentation break with the redistribution described on its own line. Whatever the spelling, the pattern is the same: the steps above the mark cannot begin on a piece of data until records have arrived from below. 2. **Semantically, from the operations.** You can predict the crossings before you look, which is why this question is asked. Records that must meet have to be moved to one place, so a crossing is forced by: - grouping by a key, where every record sharing a key has to be brought together; - joining two inputs on a key, unless both are already arranged by it or one side is small enough to be sent whole to every worker; - producing one globally ordered result, which needs a range of values per worker; - deduplicating or counting distinct values across the whole input; - any explicit re-arrangement of the data into a different number or arrangement of pieces. The interesting reading is the **difference between the two counts**. If your program changes grouping key twice and the plan shows four crossings, two of them came from somewhere you did not intend, and that gap is the finding. ## What an unbroken chain tells you A long run of steps with no cut in it is a positive fact, not an absence. It says that for every step in the run, each output piece is built from exactly one input piece — a *narrow* dependency, as opposed to a *wide* one where a step's output piece draws on many input pieces. Consequences worth stating out loud: - The whole run can be executed as one pass over each piece; engines that do **step fusion** collapse several adjacent per-record steps into a single pass, so no intermediate collection exists between them. - The number of steps printed in such a run says very little about cost. Ten fused per-record steps in one phase are usually cheaper than one crossing. - Its width is inherited. A chain with no crossing in it cannot change how many pieces are being worked on, so everything in it runs at the width the input arrived with. ## Where the models genuinely disagree What a crossing *is made of* is not the same everywhere, and this is the part candidates get wrong by generalising from the engine they know. | Engine model | What a crossing looks like at run time | |---|---| | One grouping step at a time, every intermediate written to disk | The producing side writes its output out, the consuming side fetches it afterwards; the crossing is a materialised handover | | One fixed graph kept running, each record passed through as it arrives | Records are pushed across the connection as they are produced, with no materialisation; the crossing is a live connection, not a handover point | | Continuous work as a fast succession of small finite jobs | Each small job crosses the network the same way a finite job does, once per crossing per small job | Two cautions follow. First, do not assert that everything downstream waits for everything upstream: whether a crossing paces the job that way depends on the model, and the rule about what a phase boundary does to pacing is a subject of its own. Second, a crossing is not automatically present just because a grouping is: an input that already arrives arranged by the grouping key can be grouped with no movement at all, and a good plan reading notices when the engine has spotted that. ## Using the finding Once you have the crossings located, the reading is comparative rather than absolute: - Count them, and account for each one against something in your program. - Note which are between the read and the first aggregation, since work done before a crossing travels less far than work done after it. - Note where two crossings sit adjacent with little between them, which usually means the data was re-arranged twice by different keys. What the plan will not give you is how much each crossing actually moves. That is a measurement, and costing a movement is a separate subject with its own owner.
- A plan shows two crossings in a row with only a projection between them. What does that suggest?That the data was re-arranged twice by different keys — typically a grouping on one key followed by a join or ordering on another. Each rearrangement is a full redistribution, so the usual question is whether one of the two keys can be dropped, or whether an earlier arrangement could serve both.
- The plan shows a join with no crossing on either side. What happened?Either both inputs were already arranged by the join key so matching records are already co-located, or the engine decided to send one side whole to every worker so no redistribution of the large side is needed. Both are visible from the fact that the large input's width does not change across the join.
- Does a long chain with no crossing mean the job is cheap?No. It means no records move between workers in that chain. The chain can still read far too much, decode expensive records, or call a costly body once per record. Absence of movement bounds one cost, not all of them.
saying these in an interview costs you the question
- Says every grouping forces records to move between workers
- Counts printed steps as a proxy for cost
- Assumes a crossing always writes to disk and is fetched later
- Thinks an unbroken chain can change how many pieces run
- Reads a crossing as the engine failing to optimise
- Confuses this graph with the ordering between separate scheduled jobs