Skip to content
Tech Interview Prep home
Technical interview guide

Distributed Data Processing

How frameworks like Spark parallelize work across a cluster — partitioning, shuffling, and their costs.

Read
30 min
Practice MCQs
25
Interview QA
25
Edition
v4
Editorial status
Reviewed

Scope: Apache Spark 4.2 documentation, Apache Beam and Apache Hadoop current 2026-08-31.

Overview

Curated: · Written: · Reviewed:

Distributed processing is a movement, failure, and correctness problem

A distributed engine divides a logical computation into jobs, stages and tasks over partitions. Parallelism helps only when work can be divided usefully and the coordination, serialization, network, storage and recovery costs remain smaller than the saved compute time. Start with the data size, key distribution, required result, latency, cost and failure semantics. A cluster is not an automatic cure for an inefficient algorithm or an unclear output contract.

Transformations usually build a lazy logical plan; an action requests a result and causes execution. Narrow dependencies let a child partition read a small number of parent partitions and can often pipeline in one stage. Wide dependencies—grouping, repartitioning, many joins and global ordering—require a shuffle boundary. Shuffle writes intermediate blocks, transfers them over the network, sorts or hashes them, and may spill to disk. Count bytes and partition distributions across that boundary, not merely the number of transformations.

Partitioning determines concurrency and locality. Too few partitions underuse the cluster and create large failure/retry units; too many create scheduler overhead, tiny files and excessive metadata. Choose partition count from bytes, record cost, available cores and target task duration, then validate with production-shaped measurements. Storage partitioning can avoid data movement only when the operator's distribution requirements and the stored layout are compatible. Repartition deliberately when it reduces more downstream work than it costs.

Join strategy follows measured sizes and semantics. Broadcasting a genuinely small relation can avoid shuffling the large side, but every executor must hold the broadcast and stale statistics can make the plan unsafe. Sort-merge and shuffled-hash joins distribute both sides and have different memory and ordering costs. Inspect the physical plan and runtime statistics. Adaptive execution can coalesce small post-shuffle partitions, change strategies and split skewed partitions, but it does not excuse missing statistics, bad keys or an unbounded broadcast.

Skew is a tail problem: one hot key or oversized partition leaves most executors idle while a straggler controls completion time. Compare median, upper-percentile and maximum task input, shuffle, spill and duration. Fix the cause by filtering invalid hot keys, pre-aggregating, splitting heavy keys, salting with a correctness-preserving second aggregation, isolating exceptional tenants, or using tested adaptive skew handling. Adding executors does little when only one task owns the remaining work.

Cache only reused, expensive-to-recompute data whose retained footprint and lifetime justify eviction, serialization and garbage-collection cost. Unpersist it when done. Checkpointing truncates lineage and can support recovery or iterative algorithms, but it adds durable I/O and is not a substitute for a transactional output. Memory pressure can produce spill, long garbage collection, executor loss and recomputation; measure serialized size and per-task working set before changing heap sizes blindly.

Failure recovery re-executes tasks and may recompute lineage. Therefore transformations should be deterministic for the same versioned inputs, and external effects must be idempotent or protected by an attempt-aware commit protocol. Speculative execution may run duplicate attempts; it can mitigate heterogeneous stragglers but can duplicate unsafe side effects and waste capacity on deterministic skew. Driver loss, executor loss, shuffle loss, partial output and retry exhaustion need explicit tests and runbooks.

Resource design spans driver memory and availability, executor cores and memory, task concurrency, local and remote storage, network, cluster quotas and concurrent workloads. Huge executors can worsen garbage-collection pauses and loss impact; tiny executors add overhead. Dynamic allocation responds to backlog and idleness but needs compatible shuffle preservation and bounded minimums/maximums. Fair pools or workload isolation prevent one backfill from starving production, while admission control protects shared dependencies.

Correctness and observability are inseparable. Bind output to input snapshot or offsets, code/configuration, schema and run identity. Reconcile counts and control totals by meaningful segments. Track task/stage tail latency, input/output and shuffle bytes, records, spill, locality, garbage collection, retries, executor loss, scheduler delay, skew, cache hit/eviction, output files and cost. Retain event logs long enough for incident analysis, but protect plans, paths, samples, credentials and logs as sensitive metadata.

Spark and MapReduce mechanics that interviews actually probe

Spark's driver builds a DAG of stages at shuffle boundaries; tasks run on executors. RDD lineage recomputes lost partitions from parents; checkpointing truncates that lineage to reliable storage. Shuffle write produces map-output files that reducers fetch; spark.speculation may launch a second attempt of a slow task, which is why output committers and external writes must be attempt-aware. Spark SQL adaptive execution can coalesce post-shuffle partitions, switch broadcast versus sort-merge, and split skewed partitions using runtime stats—stale ANALYZE statistics can still broadcast a 20 GB dimension. Configuration knobs for spark.sql.shuffle.partitions, memory fractions, and serializer choice change bytes on the wire; they do not change whether a non-associative aggregate is correct.

Kubernetes mode runs the driver and executors as pods with service accounts, secrets, and volume mounts; killing the driver pod fails the application even if executors still hold shuffle files. Dynamic allocation needs shuffle tracking so executors holding map outputs are not removed. Hadoop MapReduce's output commit protocol distinguishes task attempt paths from job commit: a failed attempt's files must not become visible. Beam's runner translates a portable pipeline; windowing and state semantics must be tested on the runner you will operate, not only on the DirectRunner.

Figures, locality, and failure modes

A stage whose p50 task is 8 seconds and p99 is 22 minutes is skew or straggler, not “needs a bigger cluster.” Track shuffle read/write bytes, spill memory and disk, GC time, scheduler delay, locality (PROCESS_LOCAL versus ANY), failed and killed task counts, and speculative duplicates. Data locality collapses when input is in a different AZ than executors; the job still completes while network cost dominates. Nondeterministic map with System.currentTimeMillis or an unordered hash-set aggregation will not replay the same output after a fetch failure. Security: event logs, SQL, and explain plans in the Spark UI can leak predicates and paths; treat the monitoring surface as sensitive metadata.

Test locally for logic and on a cluster for distribution. Use production-shaped sizes, key cardinality and skew; compare expected output with an independent oracle. Inject executor and driver disruption, shuffle loss, slow workers, unavailable storage, duplicate attempts and concurrent backfills. Capacity tests should identify the sustainable workload with recovery headroom, not only the fastest happy run. The final design argument must connect correctness, tail latency, cost, recovery, security and evidence.

Worked example: eight more executors do not split Acme

Join orders to customers on tenant_id. Acme is 41% of keys. Stage 4 is the shuffle.

changetask p50task p99wall clock
spark.sql.shuffle.partitions = 2008 s22 min22 min
+8 executors8 s22 min22 min
salt Acme into 16 subkeys, second aggregation9 s40 s45 s

The remaining work is one hot partition. Capacity is not divisibility.