00
START HERE
Introduction to Distributed Data Processing
The shop and its two consumers, why one machine stops being enough, where execution actually happens, and how to read cost before you tune anything. Readable with no Spark or Flink experience.
PLANNED · R1the required resultwhen one machine stops being enoughwhere execution happenshow to think about costbounded input & evolving answersthe shop fixturehow to use the pathlab entry
01
Know where work runs, what moves, and how a program becomes execution.
BEGINNER FOCUSGATE · SA01 · EXECUTION TRACING
P01Partitions & ParallelismThe large cluster whose source still has one taskWork units and execution resources constrain concurrency separately, and spare cluster memory does not help a result collected at the driver.
P02Narrow & Wide DependenciesThe group-by estimated as if it were a filterPartition-local work moves nothing; a different distribution has to bring related contributions to the same owner.
P03Plans & ExecutionThe rewrite that changed the code but not the expensive planSource order does not decide the physical plan, where a function runs, or whether an action repeats work.
P04Bounded Input & Streaming ExecutionThe recurring job whose successful runs leave new events unprocessedBoundedness describes the input, not the execution mode — and run frequency alone defines neither correctness nor freshness.
↓ NEXT MOVING DATA
02
Reduce avoidable movement, choose a join strategy, and diagnose before tuning.
JUNIOR FOCUSGATE · SA02 · RESULT-PRESERVING OPTIMIZATION + A SMALL BATCH JOB
P05Combine Before You MoveThe distributed average built from unweighted partition averagesCombining early is valid only when partial state merges back into the required result: an average has to carry sum and count.
P06Joins at ScaleThe supposedly small join side that exceeded the receiving process's budgetReplicating a reduced side avoids redistribution but spends the receiver's memory, and the memory budget belongs to a specific process.
P07Data SkewThe stage whose final task keeps the rest of the job waitingKey skew, partition size, join multiplicity, and uneven per-record cost are different causes, and the legal way to split work depends on the operator.
P08Reading the JobThe runtime regression blamed on volume without checking distributionA runtime chart is a symptom: compare task distributions, bytes, records, and phases, and take one discriminating measurement before a risky change.
↓ NEXT TIME
03
Choose the right clock, define finite groups, read progress, and handle corrections.
JUNIOR CORE · SENIOR EXTENSIONSGATE · SA03 · MEMBERSHIP, PROGRESS, CORRECTION + A BASIC STREAMING JOB
P09Event, Ingestion & Processing TimeThe sales period chosen by the time a delayed phone finally syncedThe required business clock decides grouping, and replay can change a processing-time result even when the logical events are identical.
P10Windows & the Output RowThe report that sums overlapping window totals as if each order appeared onceAssignment defines membership — one record can belong to several windows, and an updated aggregate is not automatically a new independent fact.
P11Watermarks & Event-Time ProgressThe event-time result that never becomes eligible for outputA watermark is progress under a stated policy, not proof that older events are impossible; generation, propagation, and operator use are separate.
P12Late Data & Correction ContractsThe corrected window total that the consumer adds twiceA lateness policy has to say what is accepted, what happens to retained state, and how a correction is read downstream — a threshold is not a policy.
Two short shared primers land before P10: G01 on state and what an output row asserts, G02 on event-time progress and late arrival.
↓ NEXT STATE
04
Know what is retained, estimate its size, and define how it expires or recovers.
SENIOR FOCUSGATE · SA04 · STATE ESTIMATE, EXPIRY EXPERIMENT, VERSIONED ENRICHMENT
P13Stateless & Stateful ProcessingThe deduplicator that retains every event identity without a horizonState is the information an operator must remember; ownership, working storage, durable recovery copies, and external business data are separate concerns.
P14Bounding StateThe memory estimate that multiplies a key count by hours without a rateCount the retained entries for that operator, then multiply by representation cost — and say which clock expiry runs on, because it changes the answers.
P15Streaming Joins & Versioned EnrichmentThe historical order enriched with a price version that never applied to itSnapshot reads, current lookups, event-time temporal joins, and broadcast rules carry different validity contracts under delay and replay.
↓ NEXT CORRECTNESS
05
Specify delivery and business-effect guarantees, survive failures, and establish a usable recovery point.
SENIOR FOCUSGATE · SA05 · EFFECT/PROGRESS FAILURE TRACE + A RECOVERY EXERCISE
P16Delivery Guarantees & BoundariesThe engine guarantee mistaken for proof that a seller cannot be paid twiceSource replay, state recovery, and sink commit each bound the guarantee, and a transport duplicate is not the same thing as a duplicate business effect.
P17Idempotence & Retry-Safe EffectsThe retried refund that causes a second business effectWriting the result and recording the progress are two steps, and an increment is not an idempotent upsert however many times you retry it.
P18RecoveryThe valid checkpoint whose required input is no longer retainedRecovery needs available state, replayable input, compatible code — and enough spare capacity to catch up while new input keeps arriving.
P19Checkpoints & Recovery ProgressThe frequently requested checkpoints whose last usable recovery point is oldA configured interval measures requests, not completed progress; completion, timeout, in-flight work, and sink coordination decide what you can restore.
G03, a short backpressure primer, arrives before the barrier-alignment examples in P19.
↓ NEXT FLOW AND RESOURCES
06
FLOW AND RESOURCES
R6 · P20–P22
Explain backlog and latency, identify the limiting resource, and test a bounded change.
SENIOR FOCUSGATE · SA06 · BOTTLENECK INVESTIGATION + READING A MEASURED CHANGE
P20Backpressure & LagThe upstream scaling change that leaves a slow sink and a growing backlog untouchedBacklog means admitted work exceeded completed work over the interval; it does not say which component changed, so locate the constraint first.
P21Latency, Throughput & FreshnessThe shorter trigger interval that misses the same freshness targetVisible latency spans scheduling, processing, queueing, output eligibility, and sink visibility — a trigger interval is not a latency floor.
P22Memory, Spill & CachingThe cached intermediate that raises memory pressure without useful reuseBudgets belong to particular processes, some operators spill and others fail, and reuse buys saved recomputation with retained capacity.
↓ NEXT DESIGNING FOR SCALE AND FAILURE
07
DESIGNING FOR SCALE AND FAILURE
R7 · P23–P24
Choose an architecture from consumer contracts, then transfer the reasoning to an incident you have not seen.
SENIOR FOCUSGATE · SA07 · CONSUMER-DRIVEN DESIGN + AN UNFAMILIAR INCIDENT
P23Batch or Stream?The always-running pipeline built before anyone defined the consumer's freshness needCompare candidates on freshness, correction, durability, cost, and recovery — two consumer contracts do not prescribe two independent pipelines.
P24Incidents & the Transfer ProbeThe three incident tickets ranked by alarming charts instead of business impactTriage on impact, urgency, reversibility, and evidence, and remember a familiar fix only transfers when the new operator keeps the same semantics.
↓ NEXT END OF THE PATH
✓
WHILE THE PATH IS BEING BUILT
The distributed dimension of the diagnostic is live now. Take it to see which of these patterns you already reason through — and which ones to watch for here.
Take the free assessment