Readineer
BLOG  /  DISTRIBUTED DATA PROCESSING
DISTRIBUTED DATA PROCESSING · 7 STAGES · 24 PATTERNS PLANNED

Batch and streaming, one pattern at a time

Partitions, shuffles, clocks, state, delivery guarantees, and the incidents they cause. Each pattern starts from a failure a real pipeline produces, shows the mechanism behind it, and ends with a fix and its cost.

0 PUBLISHED  ·  24 PLANNED  ·  RELEASES R1 → R7
This is the path as it is being written, in the order it will be published: each stage ships with the assessment that gates it. The lessons are not live yet, so nothing below is a link. The diagnostic that covers this dimension is live.
THE LESSONS ARE BEING WRITTEN. THE DIAGNOSTIC IS LIVE.
Find out where your distributed reasoning actually stands — shuffles, skew, state, late data, and delivery semantics, scored on the reasoning you show.
No card required  ·  Save your progress
Take the free assessment
ONE SCENARIO, ALL 24 PATTERNS
The same shop as the SQL path: the Wireless Mouse sale, the mobile order that arrives forty minutes late, prices that change under you, and two consumers who disagree — a dashboard that wants a one-minute answer it can correct later, and a settlement ledger that must never pay twice. Every pattern is a question that shop asks.
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 · R1
the required resultwhen one machine stops being enoughwhere execution happenshow to think about costbounded input & evolving answersthe shop fixturehow to use the pathlab entry
01
FOUNDATIONS
R1  ·  P01–P04
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
MOVING DATA
R2  ·  P05–P08
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
TIME
R3  ·  P09–P12
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
STATE
R4  ·  P13–P15
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
CORRECTNESS
R5  ·  P16–P19
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