Processing Engine Concepts
A cluster running a program over work one machine cannot finish - too much data, or too slow: how the work is split, what it costs when data must move, and what survives a crash.
on this pageshowhide
explore
- Choosing a Processing Engine24 questions
- The Single-Machine Ceiling4 questions
- Bounded and Unbounded Input4 questions
- Engine or Query Service4 questions
- Engine or Orchestrator4 questions
- One Runtime, Two Modes4 questions
- Cluster Overhead4 questions
- Job as a Dataflow25 questions
- Declared Shape or Loop4 questions
- Nothing Runs Until Asked4 questions
- Plan Rewrites5 questions
- The Opaque Step4 questions
- Reading a Job Plan4 questions
- Repeatability of a Result4 questions
- Splitting the Work28 questions
- The Unit of Parallelism4 questions
- Choosing How Many Pieces4 questions
- Narrow and Wide Steps4 questions
- Layout the Reader Inherits4 questions
- Arranging to Not Move4 questions
- The Output Shape4 questions
- Compute Near the Data4 questions
- Moving Data Between Workers27 questions
- The Stage Barrier4 questions
- Write and Fetch Sides4 questions
- Sending the Small Side4 questions
- Local Pre-Aggregation4 questions
- Costing a Movement3 questions
- Two Large Inputs4 questions
- Producing an Ordered Result4 questions
- Skew and Stragglers20 questions
- Reading the Spread4 questions
- The One Heavy Key4 questions
- Splitting a Heavy Key4 questions
- The Slow Worker4 questions
- Replanning at Runtime4 questions
- Memory, Spill and Caching20 questions
- The Worker's Budget4 questions
- Spilling to Disk4 questions
- Failing on One Piece4 questions
- Reusing a Computed Result4 questions
- Record Representation4 questions
- Time in a Stream20 questions
- Three Clocks4 questions
- Claiming Completeness5 questions
- Late Arrival4 questions
- Disorder and Waiting4 questions
- The Quiet Source3 questions
- Bounding an Endless Stream21 questions
- Shapes of a Window5 questions
- Emitting a Result4 questions
- Incremental or Buffered4 questions
- Joining Two Streams4 questions
- Joining Against a Table4 questions
- Long-Lived State21 questions
- The Retained Set4 questions
- Key-Bound State4 questions
- Memory or Local Disk4 questions
- State That Never Shrinks5 questions
- Changing a Running Job4 questions
- Surviving Failure22 questions
- Snapshot of a Job5 questions
- Rewinding the Input4 questions
- The Visible Effect5 questions
- Retry Granularity4 questions
- Stopping on Purpose4 questions
- Machines for the Job23 questions
- Standing or Per-Job Cluster4 questions
- Shaping a Worker4 questions
- Adding Capacity Mid-Flight4 questions
- Many Jobs, One Cluster4 questions
- Interruptible Workers3 questions
- The Coordinating Process4 questions
- Operating Jobs in Production27 questions
- Rerunning Against History5 questions
- Throughput, Lag and Failures5 questions
- Reading Backpressure4 questions
- Backlog and Catch-Up4 questions
- Debugging Remote Work5 questions
- Paying for a Job4 questions
- Testing and Changing Jobs22 questions
- Logic Outside the Cluster4 questions
- Time-Dependent Assertions4 questions
- Representative Samples4 questions
- Checking the Answer5 questions
- Altering a Live Pipeline5 questions
questions
300 · 13 sectionsWhat makes an input bounded — yesterday's finished files against a feed that never stops — and what does that ending buy?
basics
~20 sA bounded input has a last record the reader can reach; an unbounded one does not. Reaching that last record is what lets a job produce one total, one global sort or one rank, call it final, and stop.
A program is handed to a cluster and reads no input for forty seconds. What is the cluster doing?
basics
~20 sBefore the first read, a cluster must be granted capacity, launch a worker process on each machine, copy the program and its libraries to each, build a plan and list the input. None of that touches data.
A workflow scheduler and a cluster execution engine both draw a graph of steps with no cycles - what does an edge mean in each?
basics
~20 sIn a workflow scheduler's graph an edge means only 'start after the previous step reported success'; nothing travels along it. Inside one program handed to a cluster execution engine, an edge carries the records themselves from one step to the next.
When you hand a program to a cluster engine versus a statement to a query service, who owns the storage, the plan and the capacity?
basics
~20 sA cluster execution engine runs your program on machines that are part of your deployment, over storage you point it at, along a plan your code shapes. A managed query service owns all three and shows you none of them.
A pipeline over a never-ending input emits results only every thirty seconds: which execution shape is it on, and what fixes that floor?
basics
~20 sA fixed reporting cadence usually means the runtime executes repeated small finite runs: it collects whatever arrived in an interval, runs an ordinary finite job over that slice, and repeats. The floor is the interval plus that run.
What does a cluster engine know about a named operator that it does not know about a function you hand it to run per record?
basics
~20 sA named operator carries its meaning - the engine knows it keeps rows, or groups by a key - so it may reorder, narrow, fuse or skip it. A handed-over body means only 'call this per record', so it runs exactly where it stands.
Two runs of one job over the same input write the same rows in a different order — why?
basics
~20 sNothing in the job asked for an order, so none was kept. The work is cut into pieces processed by many workers at once, and rows appear in whatever sequence those pieces finished and were collected. Only a declared ordering step makes row order stable.
In an engine that runs no declared step until something demands an answer, what distinguishes a describing call from a demanding one?
basics
~20 sA describing call only adds a step to the graph the engine will run and computes nothing. A demanding call asks for an answer or for output to be written, and that demand is what makes the recorded steps execute.
A program reads a hundred-column input, keeps two columns and filters on a date - which two edits does the engine make before running it?
basics
~20 sTwo: read-time column narrowing, which tells the read to produce only the two columns later steps use, and filter at the read, which moves the date condition down so rejected rows are never produced. Both shrink what enters the job.
Which facts in the text an engine prints about a job exist only because the work is spread over machines?
basics
~20 sFour: the points where records must be sent between workers, where the listing breaks into runs of steps between those points, how many pieces each step processes at once, and what the read was told to skip.
In a cluster where storage and compute share machines, why place a piece of work on the machine already holding its bytes?
basics
~20 sThe program is kilobytes and the input is gigabytes, so shipping the work to the bytes is far cheaper than pulling the bytes across a shared network to wherever a thread happens to be free.
A job's only input is one 40 GB file compressed as a single stream — how many worker threads can read it?
basics
~20 sOne. A stored form that cannot be decoded from an arbitrary offset yields exactly one piece of the input — one slice a single worker thread reads end to end — so every other thread in the cluster stays idle during that read.
A worker thread holds one piece of the input. What decides whether it can finish a step alone or needs records other workers hold?
basics
~20 sDependency shape decides: if each output depends only on records already in the thread's piece, it finishes alone — a narrow step. If one output needs records spread across every piece, the step is wide.
A finite job's stored input is divided into eight pieces and the cluster offers 200 worker threads; how many can work at once?
basics
~10 sEight. The piece count is the ceiling on how many worker threads can be busy; the machine count only decides how many threads exist, so the other 192 lanes stay idle.
A finite job's final step holds 2,000 pieces and writes to a directory. How many files does it leave, and what set that number?
basics
~20 sRoughly one file per piece, so about 2,000. Each worker thread writes the records it already holds, where it is; nothing gathers them first. The piece count at the final step, not the machine count, set the file count.
Two 500 GB inputs are matched on a shared account key, and neither fits on one worker - what happens to both inputs first?
basics
~20 sBoth inputs are cut on the matching key by the same routing rule into the same number of destinations and sent across the cluster, so records sharing a key value land on one worker, which then matches them locally.
Why is the cost of redistributing data across a cluster counted in bytes moved rather than in records?
basics
~20 sA redistribution costs bytes: what crosses the network, what the producing side writes to local disk where output is materialised, and processor time spent encoding it. Record counts predict none of those unless every record is the same width.
A job sorts each of its 200 pieces locally and writes them out; why is the combined result not ordered end to end?
basics
~20 sSorting inside a piece orders only that piece. Unless the routing that filled the pieces gave each one a contiguous span of the key range, piece 3 can hold keys that belong before piece 2's, so reading the pieces in order interleaves ranges.
Why does folding records together on the machine that produced them cut what a grouped aggregation sends across the network?
basics
~20 sRecords sharing a grouping key are combined where they were read, so each producing machine sends one partial result per key instead of every record. The collecting side combines partials into the same final answer, over far less traffic.
When a large input is matched against a tiny one, why does copying the tiny input to every worker avoid a redistribution?
basics
~20 sMatching records on a key normally forces every worker to send each record to whichever worker owns that key. Copying the whole small input to every machine instead lets the large input be matched where it already sits.
A step's 400 pieces average 9 seconds each, yet the step runs 52 minutes — what do the median and maximum piece durations add?
basics
~20 sThe median shows the typical piece finished in about a second while the maximum shows one ran nearly the full 52 minutes: a long tail of one or two outsized pieces, not a step that is uniformly slow.
In a step that gathers every record with one key onto a single worker, why does doubling the cluster not shorten the commonest key's work?
basics
~20 sA key-to-destination rule sends every record carrying the same key to one destination, and that mapping is fixed. Extra machines add destinations, not a way to split one key, so the commonest key still runs on a single worker.
How does appending an artificial suffix to a heavy grouping key spread its records, and what second pass does that force?
basics
~20 sAppending a small integer to the grouping key turns one key value into several, so its records reach several destinations instead of one. A second pass then re-groups those partial results under the original key to produce the real answer.
A step finished and reported 200 MB where the plan assumed 200 GB - what can a runtime that re-plans mid-run change now?
basics
~20 sRuntime replanning - re-deciding unrun steps from the sizes a finished step actually produced - can change only what has not started: the shape of the next step's pieces, or a join switched to copying the small side. Finished work stands.
Two later steps each read the same computed intermediate — what does pinning that result change about the work done?
basics
~20 sPinning tells the engine to keep the computed intermediate after the first reader finishes, so the second reader reads the kept copy instead of re-running every step that produced it. Without a pin, that branch is computed twice.
One unit of work fails for lack of memory four times and kills the job while the cluster sits idle — what does that rule out?
basics
~20 sMemory is owned per worker process, so a failure inside one unit of work says nothing about total cluster capacity. Idle machines cannot lend memory to the process that is short; what reached that one unit has to change instead.
Why can one record held as an ordinary language object occupy several times the memory of the same record packed as bytes?
basics
~20 sEach field can become a separate object carrying the runtime's own per-object bookkeeping - a header, alignment padding, and a pointer to it from the record. A record of many small fields therefore costs a multiple of its raw bytes, not a small addition.
A sort in one unit of work wrote 12 GB to the worker's local disk, yet the job finished correctly — why is that by design?
basics
~20 sSpilling is planned degradation: an operator that cannot hold its working set writes part of it to disk attached to the worker and reads it back to finish. The answer is identical; only the runtime grows.
A worker process is given one fixed memory budget. What competing demands share it while the job runs?
basics
~20 sOperators borrowing room while they run, any results the job asked the engine to keep, and runtime overhead the engine never counts - all shared by every unit of work inside that one process, with no lending from idle workers.
A streaming job advances a timestamp asserting nothing older than it is still expected. What does that assertion license?
basics
~20 sA completeness claim - a watermark - is a timestamp the job carries with the records, asserting no record older than it is still expected. It makes any group ending at or before it eligible to close.
A job groups records by the moment stamped in the payload — where does one that arrives after its group closed go, and who would notice?
basics
~20 sNowhere visible: a record whose group has already closed is usually discarded without an error, so the totals are simply short. An explicit count of records dropped for lateness is what makes the loss visible at all.
A continuous job keeps receiving records stamped earlier than ones it already handled — why is that normal, and what causes it?
basics
~20 sOut-of-order arrival is a property of the transport, not a fault: buffered devices, retries and parallel inputs merged together all deliver an earlier-stamped record after a later one. It costs nothing by itself until a time-based grouping has to be treated as finished.
A daily order count runs over a continuous stream. Which three moments could each order be counted under, and which should the count use?
basics
~20 sThree moments compete: when the order was placed, stamped in the payload; when the receiving system wrote the record down; and the wall clock of the worker that handled it. A daily count belongs to the payload moment.
A job carries a timestamp asserting nothing older than it is still expected. Why is that a guess rather than an observed fact?
basics
~20 sNothing in an endless input reports what has not yet been sent. The job infers the claim from the moments it has seen plus a standing assumption about how far behind a record may run, so it can be too early or too late.
While a five-minute grouping computes a running total per key, what must the job hold, and what changes for a median?
basics
~20 sA running total needs one number per key while the group is open. A median needs every value in the group, because no single carried value can be merged into the answer. Same interval, very different memory.
When each arriving order is enriched from a reference table, what supplies the bound for the computation and when is a result emitted?
basics
~20 sThe record's own moment supplies the bound: it selects one version of the reference row, the one whose span of validity contains that moment. Nothing accumulates and no group closes, so each record yields its output once its match can be resolved.
Two endless inputs are joined on a shared key. Why must the job hold records from both inputs, not just one?
basics
~20 sNeither input ends, so a record on either side may meet its partner later. Both sides are therefore held while the other is awaited, and a bound on how far apart their two moments may be is what lets a held record ever be dropped.
A grouping uses ten-minute spans restarted every two minutes. How many groups does one record belong to, and why?
basics
~10 sFive. A span of fixed length restarted every step shorter than that length overlaps its neighbours, so a record's moment falls inside length-divided-by-step spans at once and is counted in every one of them.
Over an endless input, what three ways can a grouping boundary be defined, and what sets each group's extent?
basics
~20 sThree shapes: back-to-back spans of equal length, where a record lands in exactly one group; fixed-length spans restarted every shorter step, where a record lands in several at once; and groups the data itself closes after a stated quiet period.
A long-running job holds per-key running totals. Why does deploying new code against it need a plan a stateless service does not?
basics
~20 sNew code has to read back the retained set the previous build wrote — everything the job holds between records. A stateless redeploy discards nothing; here, entries the new build cannot match or decode mean a refused start, or history silently lost.
In a job that keeps a running total per customer, why can a stateful step read only the current record's customer total?
basics
~20 sKey-bound state is addressable only under the grouping key of the record in hand. Each key's entry lives with the worker that owns that key, and the step is handed no way to reach another key's entry.
A job deduplicates events by holding every id it has seen. Why does that retained set grow forever, and what bounds it?
basics
~20 sA deduplication check only ever adds: each new id creates an entry and nothing deletes one, so the retained set grows with distinct ids for the life of the job. Only a removal rule bounds it.
A continuous job runs a filter, a per-user running total and a repeat-identifier dropper — which of the three hold something between records?
basics
~20 sThe running total and the repeat-identifier dropper hold entries between records: one accumulator per user, one entry per identifier already seen. The filter holds nothing, because its verdict depends only on the record in hand.
A job keeps everything it remembers between records as live objects inside the worker process - what caps that, and what happens at the cap?
basics
~20 sThe cap is the memory of that one worker process, shared with everything else the process is doing. There is normally no overflow path, so passing it degrades and then kills the worker rather than quietly moving entries to disk.
A long-running job crashes 40 seconds after its last recovery point. What does resuming from that saved picture restore, and what work is done again?
basics
~10 sResuming restores each worker's accumulated values and the recorded read position exactly as they stood at one coherent moment. Everything the job did in the following 40 seconds is read again and recomputed.
A job crashed while writing its results and restarted from its last recovery point — why can the destination hold some rows twice?
basics
~10 sRecovery is reprocessing: a restart re-reads input the job had already handled and writes it a second time. The destination remembers nothing of the first attempt, so a plain append adds those rows again.
A worker dies mid-piece in a 400-piece finite job: what does the engine re-run, and what must be true for that to be correct?
basics
~20 sThe engine re-runs only the failed unit of work — the smallest piece it hands one worker — on whichever machine has capacity. That is correct only when the piece can be recomputed from re-readable input and its partial output is discarded.
A job crashes mid-run and restarts. Why does its recovery depend on the input being re-readable from a recorded position?
basics
~20 sRecovery is reprocessing: a restart re-reads input it already consumed. So the source must retain those records and let a reader resume from a durably recorded position — a source that forgets each record after handing it over makes the lost work unrecoverable.
Parallel workers never pause at the same instant, so how does a running job save parts that add up to one coherent picture?
basics
~20 sA marker injected at the sources travels with the records; each worker saves its own part as the marker reaches it and then passes it on, so every part covers the same prefix of the input.
A job needs machines before it can run. What are the three ways a platform supplies them, and who sizes them under each?
basics
~20 sA standing pool that is already up and shared by many jobs; a cluster raised for one job and torn down with it; or a managed compute service that shows no machine at all. Sizing moves from submitter to provider across the three.
A job over input that ends gains extra worker processes halfway through its run. Where can that new capacity first do useful work?
basics
~20 sOn pieces of work that have not started yet, normally from the next round of work onward. A piece already running on another worker process is not cut in half and shared, and finished work is not redone to spread it more evenly.
The provider takes back a discounted machine in the middle of a run. Besides the piece it was computing, what does the job lose?
basics
~20 sEverything that machine held goes at once: every worker process on it, the pieces they were running, anything cached there, and intermediate output produced there that other workers had not yet fetched. Finished work can therefore have to be produced again.
A job's worker processes each run several concurrent work slots — what do the slots inside one process share?
basics
~20 sWork slots in one worker process share that process's single memory pool, its local scratch space, its fixed start-up overhead and its fate: they compete for the same memory, and all of them die together when the process does.
A step in a running job spends most of its time unable to hand its output on — what does that indicate?
basics
~20 sBackpressure: a step further along the job graph cannot accept work as fast as this one produces it, so this one blocks. The figure locates a bottleneck ahead of the resisting step, not a fault inside it.
Why can the diagnostic lines a worker process printed be unrecoverable once a job run has ended and its machines are released?
basics
~10 sA worker's diagnostic lines are written on the machine that ran it, and the cluster releases that machine when the work ends. Only what the run carried off its machines stays readable.
A continuous job reports it is four million records behind its input — why does that number alone not say whether it will recover?
basics
~20 sFour million records behind is a level, and recovery depends on the trend. Compare two readings taken minutes apart: a distance that is growing means the job is losing ground, a shrinking one means it is already draining.
A teammate says a running job 'feels slow' and gives no other detail. Which numbers do you read first, and why those?
basics
~20 sFive numbers cover almost every 'slow' report: processed throughput in records and in bytes per second, duration per step of the job graph, distance behind the input, and failed or retried units of work. Read them before changing anything.
A job written for a live feed is replayed over six months of stored input - which of its parts depend on the wall clock?
basics
~20 sAnything the code reads from the present: the current time written into output, timeouts and inactivity gaps measured in real seconds, rate thresholds per minute, expiry of held state, and relative date filters. Compressed into one hour, none behaves as it did live.
The first 100,000 rows of an input file are used as a development cut. What does that cut hide about the whole input?
basics
~20 sThe first rows of a file are whatever was written first - one day, one source, one writer's share - so they carry neither the whole input's key distribution nor the rare record shapes that actually break the job.
Which part of a distributed job can a plain unit test call without a cluster, and what stops it?
basics
~20 sThe per-record rule can: a plain function taking ordinary values and returning ordinary values. It stops being callable once its signature mentions a runtime type, or it reads configuration, storage or the clock itself instead of taking them as arguments.
Why is sleeping in a test a poor way to prove a job groups records by when they happened?
basics
~20 sSleeping proves only that the machine's wall clock advanced. Instead hand the job records stamped with chosen moments and advance from the test the job's own claim that nothing older will arrive, so groups close on command.
Why does comparing a distributed job's output line by line against a stored expected result fail even when every value is right?
basics
~10 sA distributed run promises the right records, not an order, a file layout or bit-identical arithmetic. A line-by-line comparison asserts all three at once, so it fails on arrangement while every value is correct.
What does running two versions of a live job's logic over the same input, with only one publishing, prove that a test cannot?
basics
~20 sRunning both versions over the same production input, with only one publishing, measures the change against real volume, real key distribution and real dirty records — evidence a test cannot give, because a test holds only the cases its author imagined.