skip to content

Which decisions does an execution engine own once a nightly aggregation is handed to it as a description?

level: middleimportance: must knowfreq 54%

answer

  1. strategy moves, meaning stays
  2. ordering, chunking, placement
  3. one pass or two is not yours
  4. grouping is a relationship, not an algorithm
  5. visibility and data statistics decide quality

basics

~20 s

The engine owns strategy, not meaning: the order independent stages run in, how the input is chunked and partitioned, where each stage runs, whether described stages execute as one pass, which algorithm implements a grouping, and what gets materialised.

solid answer

~50 s

Everything the description does not state is the engine's to decide. Concretely that is the order of independent stages, how the input is chunked and partitioned across workers, where each stage runs, whether stages written separately are executed together, which algorithm implements a described relationship such as a grouping, and whether an intermediate is materialised or streamed. What stays yours is the meaning: which records are in scope and what the result is defined to be. The mirror image matters just as much - having handed those decisions over, you must stop assuming a visit order, a pass count, or that a stage's cost is proportional to how short it looks on the page. And the engine can only use the freedom well where it can see through the stages and knows something about the shape of the data.

code

pseudocode · 12 lines
pseudocode
description:
    orders.keep(recent).project(toRegionAmount).groupAndSum(by = region)

one legal execution:
    for each partition p of orders, in parallel:
        scan p once, testing and projecting each record,
        adding surviving amounts into a local map keyed by region
    merge the local maps

another legal execution:
    order the recent orders by region, then walk the ordered run once,
    summing each block of equal regions as it ends

go deeper

for a junior

Learn the line first: meaning belongs to the description, strategy belongs to the engine. If you can name two decisions on the strategy side - the order of the work and how it is split - you can hold this conversation.

for a middle

Enumerate the decisions concretely and pair each with the assumption it costs you. The strongest point to make is that a grouping is a relationship, and an engine may implement it by accumulating or by ordering, both correct.

for a senior

Bring the operational half: describe a run where the partition count or the grouping algorithm changed under you because the data's shape changed, and how you noticed it was strategy and not meaning.

for a principal

Talk about what the platform owes the engine. Visibility and data statistics are what turn delegated strategy into value, so standards about how jobs are described are worth more than any single job's tuning.

## Meaning stays, strategy moves A description handed to an engine draws a line through the job. On one side is **meaning**: which records are in scope, what a regional total is defined to be, which relationships hold between stages. Meaning is yours, permanently - an engine that altered it would be producing a wrong answer, not a faster one. On the other side is **strategy**: everything about how the meaning is realised. Strategy is what you did not write down, and that is precisely what the engine now owns. The practical skill is being able to enumerate what is on the strategy side, because that list is both the value of the approach and the source of every surprise it produces. ## The decisions that move - **Order of independent stages.** Stages that do not depend on one another may be executed in any order, and a narrowing stage may be pulled ahead of an expensive one - when the engine can see through both well enough to know that is safe. - **Chunking.** One record at a time, a block at a time, or the whole input in memory. This decision sets the memory profile more than anything else in the plan. - **Partitioning and parallel width.** Whether the input is split, on which key, and into how many parts. A key with few distinct values caps how wide the job can go. - **Placement.** Whether a stage runs next to the data or the data is moved to the stage. - **Fusion.** Whether two stages written as two are performed together rather than one after the other. - **Physical strategy for a relationship.** A grouping says *values with the same key belong together*. That is a relationship, not an algorithm: one implementation accumulates into a keyed map, another orders the records so that equal keys become adjacent runs. Both satisfy the description. - **Materialisation.** Whether an intermediate is written down - to memory, to disk - or streamed straight into the next stage. | Decision | What the description says | What the engine decides | |---|---|---| | Visit order | nothing | any order that preserves the result | | Chunk size | nothing | record, block or whole input | | Parallel width | nothing | partition count, split key | | Grouping | keys belong together | accumulate, or order and walk | | Passes | stages in sequence | how many passes actually occur | ## What you must stop assuming Handing those decisions over costs you a set of assumptions that step-by-step code lets you keep. Four of them cause most of the confusion: 1. **A visit order.** Nothing promises records are processed in input order, or that two records near each other in the input are processed near each other in time. 2. **A pass count.** Stages in the text are not passes over the data. Two described stages may be one pass; one described stage may involve more than one. 3. **Stable resource use between runs.** The same description may be executed with a different partition count or a different grouping algorithm tomorrow, because the engine's picture of the data changed. 4. **Cost proportional to text.** The most expensive thing in a description is frequently the shortest phrase in it. A side effect buried inside a stage collides with all four at once: its ordering, its count and its timing are all now the engine's business rather than yours, which is exactly why descriptions handed to a planner are expected to be built from stages that do not smuggle effects in. ## What the engine needs to use the freedom well The freedom is only worth something if the engine can act on it, and that needs two inputs. - **Visibility.** Stages the engine can see through - described relationships and operations it recognises - can be reordered, fused, or implemented by a better algorithm. A stage that is an opaque callback is a wall: the engine can neither estimate its cost nor establish that moving work across it is safe. - **Information about the data.** How many distinct keys, how the sizes are distributed, how much survives a narrowing stage. Without it the engine guesses, and a guess that is wrong about skew produces a plan that is correct and slow. That is the reason two descriptions with the same meaning are not equally optimisable. Expressing the intent in terms the engine recognises leaves it room to work; wrapping the same intent in code it cannot inspect leaves it executing your stages in the order you happened to write them - all the unpredictability of the intent-first form with none of the benefit. ## The answer an interviewer wants Name four or five decisions, say plainly that meaning is not among them, and then close the loop: the same delegation that lets the engine improve the job is what stops you predicting its cost, so the description alone will never tell you what the run will do.

  • Which assumptions does the author have to give up in exchange?
    The order records are visited in, the number of passes over the data, the timing of anything a stage does besides returning a value, and stable resource use from run to run. Also the intuition that a short phrase is cheap work: in a description, text length says nothing about cost.
  • Why are two descriptions with the same meaning not equally optimisable?
    Because optimisation needs visibility. Intent expressed in operations the engine recognises can be reordered, executed together, or implemented by a better algorithm. The same intent wrapped in a callback the engine cannot inspect is a wall it must plan around, so it ends up running your stages in the order you wrote them.

saying these in an interview costs you the question

  • Assumes records are processed in the order they appear in the input
  • Thinks each stage in the text is one pass over the data
  • Expects a stage's cost to match how short it looks
  • Believes an engine may drop records it judges irrelevant
  • Cannot separate a described relationship from an algorithm that computes it
  • Puts side effects inside a stage and expects their order to hold