skip to content

Parallelism and Work Distribution

Splitting work to actually go faster: data versus task parallelism, fork-join with work-stealing schedulers, pipeline stages, and Amdahl's law setting the ceiling. Interviewers use Amdahl to check that you know why doubling the cores rarely halves the time.

part ofComputer science fundamentalsoverview, primer and where to startread it →
on this pageshow

questions

30

Explain the difference between data parallelism and task parallelism, give a concrete example of each, and say what determines how far each one can scale.

level: juniorimportance: must knowfreq 56%

answer

  1. same code / different data vs different code / same time
  2. data parallel width grows with input; task parallel width is fixed
  3. task parallelism floor = critical path
  4. data parallelism = throughput; task parallelism = latency hiding
  5. they nest: parallel stages, parallel partitions inside a stage

basics

~20 s

Data parallelism runs the same operation over different slices of one data set — its scaling limit is how finely you can split the data. Task parallelism runs different operations concurrently — its scaling limit is the number of independent tasks and their dependency graph. Resizing a million images is data parallel; fetching a user profile and their orders at once is task parallel.

solid answer

~60 s

**Data parallelism**: one operation, many data partitions. Split a collection into chunks, apply the same function to every chunk, combine. Example: computing a checksum for each of 10 million records, or resizing every image in a batch. The degree of parallelism is bounded by the number of partitions you can create — so it grows with the data, which is why it scales out across cores and machines. **Task parallelism**: different operations, run concurrently because they don't depend on each other. Example: a page render that fetches profile, orders and recommendations in parallel, then merges. The degree of parallelism is bounded by the number of independent tasks in the dependency graph — a fixed, usually small number that does *not* grow with input size. The critical path (longest dependency chain) sets the floor on completion time. They compose: a task-parallel stage can itself be internally data parallel. The practical consequence is that data parallelism is what you reach for when you want to use 64 cores, and task parallelism is what you reach for when you want to overlap a handful of independent latencies.

code

text · 13 lines
text
DATA PARALLEL  (width grows with N)
  input: 1,000,000 records
  split into 64 chunks -> 64 workers run score(chunk) -> merge counts
  more records  => more chunks => more usable cores

TASK PARALLEL  (width fixed at 3)
  render_page():
     A: fetch_profile()     |
     B: fetch_orders()      |  independent, run concurrently
     C: fetch_recs()        |
     D: merge(A, B, C)         waits for all three
  more users on the page  => still 3 concurrent tasks
  wall clock >= duration of the slowest of A/B/C, plus D

go deeper

for a junior

Define both with one clean example each, and state that data parallelism means the same operation on different slices of data.

for a middle

Add the scaling argument: data-parallel width grows with input size, task-parallel width is capped by the number of independent tasks and the critical path.

for a senior

Discuss how the two nest in a real pipeline, what each demands (associative combine and no shared accumulator vs a correct dependency graph), and which one you reach for depending on whether you lack CPU throughput or are hiding latency.

for a principal

Frame decomposition as an architectural commitment — the partition key, state placement, rebalancing and failure granularity follow from the choice — and note that task decomposition sets a hard ceiling you cannot buy your way past.

## Two ways to decompose work When you want work to happen simultaneously, you must first decide *what* to split. There are two answers, and they behave very differently. **Data parallelism (domain decomposition).** The same computation is applied to different pieces of data. You partition the input, run identical code on each partition, and combine results. All workers execute the same instructions on different addresses. Examples: multiplying a matrix by splitting it into row blocks; scanning a 20 GB log by splitting it into byte ranges; scoring a million rows through the same model; adjusting brightness on every pixel. **Task parallelism (functional decomposition).** Different computations run at the same time because they are mutually independent. Workers execute *different* code. Examples: while rendering a dashboard, one task queries the database, another calls a pricing service, a third loads static config; in a compiler, parsing one file while type-checking another; in a game loop, running physics and audio concurrently within a frame. An easy test: if you doubled the input size, would you get more opportunities to run in parallel? If yes, it is data parallelism. If the number of concurrent things stays the same, it is task parallelism. ## Scaling behaviour, which is the point of the distinction **Data parallelism scales with data.** With N items and P workers you can, in principle, use every worker as long as N ≫ P. Add machines and you can add partitions. This is why every large-scale processing framework — parallel collections, GPU kernels, sharded databases, map-style batch jobs — is built on data parallelism. Its limits are partition skew (one partition much bigger than the others), the cost of splitting and combining, and any step that must see all data at once. **Task parallelism has a hard ceiling.** If the job decomposes into 5 independent tasks, 6 cores gain you nothing over 5, and no amount of extra hardware helps. Worse, tasks usually have dependencies: the merge step waits for the fetches. The **critical path** — the longest chain of dependent tasks — is a lower bound on wall-clock time no matter how many workers you have. Task parallelism therefore targets *latency hiding* (especially overlapping I/O waits), not throughput scaling. A useful summary: data parallelism gives you *scalable* parallelism whose width is a function of input size; task parallelism gives you *fixed-width* parallelism whose width is a function of program structure. ## They are not exclusive Real systems nest them. A request handler runs three independent service calls (task parallel); one of those calls internally re-ranks 50,000 candidates by splitting them across cores (data parallel). A stream processor runs distinct stages concurrently while each stage is replicated across partitions of the key space. The right question in an interview is not "which one is this?" but "where does the parallel width come from at each level, and what caps it?" ## What each one demands of your code **Data parallel** requires that per-partition work be independent: no partition may read another's in-progress results, and no shared mutable accumulator may be updated without coordination. The combine step must be associative (and commutative if partitions can finish in any order) or you must preserve partition order explicitly. The classic mistake is a shared counter or shared collection mutated from all workers — it either corrupts data or serialises the whole computation on one lock, erasing the speedup. **Task parallel** requires a correct dependency graph: you must know which tasks can start before which others finish. Errors show up as either a missing edge (a task reads a result that is not ready — a race) or an over-conservative edge (needless serialisation). Task-parallel work is also frequently heterogeneous in duration, so the slowest task dominates and load balancing is about *which* tasks, not *how many* items. ## Where the confusion usually lands Two neighbouring terms are worth separating explicitly: - **Concurrency vs parallelism.** Concurrency is a structuring property — multiple logical activities in flight, possibly interleaved on one core. Parallelism is a hardware property — activities physically executing at the same instant. Both decompositions above are ways of *obtaining* parallelism; task decomposition is also useful with no parallelism at all, purely to overlap waiting. - **Data parallelism vs SIMD.** SIMD (single instruction, multiple data — vector instructions) is a hardware realisation of data parallelism inside one core. Data parallelism as a design concept is independent of whether it is realised by vector lanes, threads, or machines. ## Choosing between them Ask what you are short of. If you are short of **CPU throughput** on a large input, decompose the data — that is the only route that scales with the machine. If you are short of **wall-clock time on a small input dominated by independent waits** (network, disk, other services), decompose the tasks — you are hiding latency, not adding compute. If the input is large *and* the pipeline has stages, do both: task-parallel stages, data-parallel within a stage.

  • Your job splits into exactly four independent tasks and you are given a 32-core machine. What speedup can you expect, and what would you do about it?
    At most about 4x, and less if the tasks differ in duration, because the longest task bounds the wall clock. Extra cores are idle. To use them you must find data parallelism inside the tasks — partition the work each task does — or find more independent tasks by decomposing further. Adding hardware alone cannot exceed the task-graph width.
  • Can a workload be both data parallel and task parallel at the same time?
    Yes, and large systems usually are. Distinct pipeline stages run concurrently as tasks, while each stage internally splits its input across workers as data. The useful analysis is per level: identify where the parallel width comes from at each layer and what caps it — task count and critical path at the outer level, partition count and skew at the inner one.

A restaurant kitchen. Ten cooks each chopping their own crate of onions is data parallelism — add more onions and you can add more cooks. One cook on sauce, one on grill, one plating is task parallelism — a fourth cook has no distinct station to take, and the meal is not ready until the slowest station finishes.

saying these in an interview costs you the question

  • Using data parallelism and task parallelism as synonyms for multithreading in general.
  • Claiming task parallelism scales with more cores, when its width is fixed by the number of independent tasks.
  • Equating data parallelism with SIMD or vector instructions only; SIMD is one hardware realisation of the idea.
  • Confusing concurrency (structure, interleaving) with parallelism (simultaneous execution on real hardware).
  • Assuming data-parallel workers can freely share a mutable accumulator without coordination.

context

open as a page

Describe the fork-join model of parallel computation: how does a unit of work split itself, and what does the 'join' step guarantee to the code that runs after it?

level: juniorimportance: must knowfreq 50%

basics

~20 s

A task checks whether its input is small enough to do directly. If not, it splits the input into independent pieces, forks them so they can run in parallel, then joins - waits for each piece to finish - and combines their results. Join gives completion and result visibility.

open as a page

Explain the map, shuffle (regroup), and reduce phases of the map-reduce processing model: what does each phase do to the data, and why is the middle phase needed at all?

level: juniorimportance: must knowfreq 50%

basics

~20 s

Map transforms each input record independently into key-value pairs. Shuffle regroups those pairs so that all values for the same key land together on one reducer. Reduce folds each key's group into a result. Without the regroup step a reducer would only see part of each key.

open as a page

Explain pipeline parallelism as a way to organize concurrent work: what a stage is, how items move between stages, and what determines how many items the arrangement finishes per second once it runs steadily.

level: juniorimportance: must knowfreq 55%

basics

~20 s

Split a repeated job into ordered steps called stages, each with its own worker and a queue in front of it. Different items sit in different stages at the same time, so all workers run at once. Steady-state throughput is one item per slowest stage's time.

open as a page

In a recursive divide-and-conquer parallel algorithm, why do implementations stop splitting once a subproblem is below some size and run the remainder sequentially, and how would you choose that threshold?

level: middleimportance: must knowfreq 45%

basics

~20 s

Each split costs bookkeeping - creating, queueing and scheduling a task. Below some size that overhead exceeds the useful work, so the recursion would slow things down. Pick the threshold by measuring: choose the smallest size where parallel still beats sequential, comfortably above the break-even point.

open as a page

When you fold a collection with a binary operation in parallel instead of left-to-right, what properties must that operation have for the result to be correct, and what role does an identity element play?

level: middleimportance: must knowfreq 52%

basics

~20 s

The operation must be associative, so any bracketing of the same ordered elements gives the same answer, letting chunks combine in a tree. An identity element gives empty chunks a value to return, so partitioning is free. Commutativity is only needed if chunks may combine out of order.

open as a page

A three-step processing chain runs each step on its own worker with queues between them, taking 10 ms, 50 ms and 10 ms per item. What throughput and per-item latency do you expect, and what would you change to raise throughput?

level: middleimportance: must knowfreq 50%

basics

~20 s

Throughput is 1 per 50 ms (20/s), set by the slowest step; latency is at least 70 ms plus waiting. To improve it, attack only the 50 ms step: make it cheaper, split it into smaller sequential steps, or run several copies of it in parallel. Speeding up the 10 ms steps changes nothing.

open as a page

Amdahl's law describes the speedup limit when you parallelize a fixed-size workload. State the law, explain what the serial fraction is, and work out the best possible speedup for a job whose runtime is 5% serial.

level: middleimportance: must knowfreq 60%

basics

~20 s

Amdahl's law: for a fixed job with serial fraction s, speedup on N workers is 1 / (s + (1 - s)/N), capped at 1/s as N grows. At 5% serial the ceiling is 20x, however many workers you add.

open as a page

A parallel runtime spreads dynamically created tasks across a fixed set of worker threads. Explain how a work-stealing scheduler does this, and why it is usually preferred to having every worker pull from one shared task queue.

level: middleimportance: must knowfreq 42%

basics

~20 s

Each worker owns a local task queue and normally runs tasks it created itself; a worker whose queue is empty steals a task from another worker's queue. Compared with one shared queue, nearly every handoff stays thread-local, so there is no central contention point and no coordinator assigning work.

open as a page

How do you choose task granularity — the amount of work per parallel unit — when each unit carries a fixed scheduling overhead? Explain the cost model and what goes wrong at each extreme.

level: seniorimportance: must knowfreq 48%

basics

~30 s

Each unit costs a fixed overhead (submit, schedule, synchronise) on top of its useful work. Make units large enough that overhead is a small fraction of the work — typically at least tens of microseconds of work per unit — but small enough that you have several times more units than workers, so the tail after the last split is short. Too fine: overhead dominates. Too coarse: idle workers and a long tail.

open as a page

A batch computation was parallelised across all available cores and ended up slower than the single-threaded version, with correct results. Walk through the causes you would investigate and how you would distinguish between them.

level: seniorimportance: must knowfreq 46%

basics

~20 s

Look for: per-unit scheduling overhead exceeding the work, contention on a shared lock or counter, false sharing of cache lines between workers, saturation of a shared resource (memory bandwidth, one disk, one connection pool), too little total work to amortise startup, and oversubscription causing context switches. Distinguish by checking whether speedup falls as workers are added — that points to contention or bandwidth, not overhead.

open as a page

A divide-and-conquer computation runs on a fixed-size pool of worker threads. What goes wrong when its tasks perform blocking work - a network call, a file read, or waiting on a lock that another task holds - and how do you keep the pool healthy?

level: seniorimportance: must knowfreq 45%

basics

~20 s

Blocked workers still occupy their pool slot, so throughput collapses and, if the blocked tasks are waiting on work that is queued behind them, the pool deadlocks. Keep the compute pool for CPU-bound non-blocking work; run blocking work on a separate, larger pool or an async path.

open as a page

Between the stages of a concurrent processing chain you can place an unbounded queue or a fixed-capacity one. Which do you choose, how do you pick the capacity, and what happens to the neighbouring stages when that queue is full or empty?

level: seniorimportance: must knowfreq 48%

basics

~20 s

Always bounded. An unbounded queue turns a rate mismatch into unbounded memory and unbounded latency. Bounded means a full queue blocks the upstream stage and an empty queue stalls the downstream one — both visible, self-limiting signals. Size it to cover normal jitter only, since depth becomes latency.

open as a page

What makes a workload "embarrassingly parallel", and what specific properties disqualify a workload from that label?

level: middleimportance: should knowfreq 44%

basics

~20 s

Embarrassingly parallel means the work splits into independent pieces that need no communication or coordination while running — each piece reads only its own input and writes only its own output. Cross-piece dependencies, shared mutable state, order sensitivity, or contention on one shared resource disqualify it.

open as a page

You are splitting a large collection across worker threads. Compare contiguous-block partitioning, round-robin partitioning, and dynamic chunk handout, and say which you would choose when the per-item cost varies a lot.

level: middleimportance: should knowfreq 40%

basics

~20 s

Block gives each worker one contiguous range — best locality, worst load balance under uneven costs. Round-robin interleaves items across workers — balances a smooth cost gradient, hurts locality. Dynamic handout gives workers a new chunk when they finish — best balance under unpredictable costs, at the price of coordination on the shared cursor. With highly variable per-item cost, use dynamic (or many small static chunks).

open as a page

In a recursive divide-and-conquer task that splits work in two, why is the pattern 'fork the left half, compute the right half on the current thread, then join the left' preferred over forking both halves and joining both? And what goes wrong if you join a subtask immediately after forking it?

level: middleimportance: should knowfreq 38%

basics

~20 s

Forking both halves wastes the current thread, which then only waits; doing one half in place keeps it busy and halves the number of tasks. Forking then immediately joining is worse still - it serializes the two halves while paying full task overhead, giving sequential speed at parallel cost.

open as a page

Compare reducing n values by accumulating them one at a time into a single running total against combining them in a balanced binary tree. How many combine operations does each perform, how deep is each, and what does that mean for parallel execution?

level: middleimportance: should knowfreq 38%

basics

~20 s

Both perform n-1 combines - the total work is the same. The chain is n-1 steps deep because each waits for the previous; the balanced tree is about log2(n) levels deep and each level's combines are independent. Depth, not operation count, is what limits parallel speed.

open as a page

How do you measure the parallel efficiency of a program, which baseline should the measurement use, and what does it mean if you observe a speedup greater than the number of processors used?

level: middleimportance: should knowfreq 30%

basics

~20 s

Speedup S = T(baseline) / T(N); efficiency E = S/N, the fraction of each worker that does useful work. The baseline must be the best sequential implementation, not the parallel code on one worker. Speedup above N (superlinear) is usually a cache or memory-hierarchy effect, not an error.

open as a page

In a work-stealing scheduler, each worker keeps its tasks in a double-ended queue: the owner pushes and pops at one end (LIFO for itself) while thieves take from the other end (FIFO with respect to the owner's pushes). Why is it arranged that way rather than both ends behaving the same?

level: middleimportance: should knowfreq 33%

basics

~20 s

LIFO locally keeps execution depth-first: the newest task's data is still cache-hot and memory stays bounded. Thieves take the oldest task because, near the root of the computation, it is the largest subtree — so one steal buys a lot of work. Using opposite ends also keeps owner and thief off the same memory, so the local path needs almost no synchronization.

open as a page

In a map-reduce style job, what is a combiner (a local pre-aggregation step run on a worker's own output before data is exchanged), when is it safe to apply, and how would you handle an aggregation like an arithmetic mean where it is not directly applicable?

level: seniorimportance: should knowfreq 40%

basics

~20 s

A combiner folds a worker's own key-value output locally before the exchange, cutting the data crossing the network. It is safe when the reduce operation is associative and commutative and its output type can be fed back in as input. For a mean, pre-aggregate (sum, count) pairs and divide at the end.

open as a page

A parallel aggregation groups records by key and assigns each key to a worker with hash(key) modulo the worker count. Most workers finish quickly while one runs for hours. What is happening, what options do you have, and what guarantees about output ordering does this style of processing give you?

level: seniorimportance: should knowfreq 42%

basics

~20 s

Key skew: one key (or a few hashing together) holds a huge share of the records, and a key cannot be split across workers, so its worker becomes the straggler. Options: pre-aggregate locally, salt the hot key into sub-keys and aggregate twice, or isolate hot keys. No ordering is guaranteed - not of values within a group, nor of outputs.

open as a page

After a team reorganized a request handler into a chain of concurrent stages with queues between them, total requests per second went up but the time an individual request takes got worse. Explain why that is the expected outcome and how you would reason quantitatively about the per-item time.

level: seniorimportance: should knowfreq 45%

basics

~20 s

Pipelining raises throughput by keeping several items in flight, but one item still visits every stage in order and now also pays handoff and queue waiting. Per-item time = sum of stage service times + waiting. Little's law gives the wait: average time in system = items in system / throughput.

open as a page

Amdahl's law says parallel speedup is capped at 1 divided by the serial fraction, while Gustafson's law predicts speedup that keeps growing with processor count. Which assumption differs between the two models, and when does each one apply?

level: seniorimportance: should knowfreq 33%

basics

~20 s

Amdahl fixes the problem size and asks how much faster it finishes; Gustafson fixes the runtime and asks how much bigger a problem you can solve. Scaled speedup is S = N - s(N-1), which grows with N because the parallel work grows too.

open as a page

A recursive parallel sort splits an array in half, sorts both halves concurrently, and then merges the two sorted halves with an ordinary sequential merge. Even on a machine with many idle cores the speedup flattens out. How do you reason about the ceiling, and how would you restructure the decomposition to raise it?

level: principalimportance: should knowfreq 30%

basics

~20 s

The sequential merges form a chain: the final merge alone touches every element, and nothing can shorten it. That chain is the critical path of the task graph and bounds the finish time no matter how many cores exist. Raise the ceiling by parallelizing the merge itself.

open as a page

You must process a continuous stream where each item goes through parse, then a remote enrichment call, then a transform, then a database write. Would you build a staged chain with dedicated workers per stage, or a pool of workers that each carry one item through all four steps? How do you decide?

level: principalimportance: should knowfreq 38%

basics

~20 s

Decide by resource heterogeneity and serialization. If the steps need very different resources, or one step must be single-threaded, batched, or ordered, use stages sized independently. If items are homogeneous with no serial resource, a pool carrying each item end to end is simpler and keeps better locality.

open as a page

You measure a service's throughput at increasing concurrency levels and find it rises to a peak around 24 in-flight workers, then declines as you push higher. How would you use scalability modelling to decide between buying more hardware, partitioning the system, or redesigning the hot path?

level: principalimportance: should knowfreq 24%

basics

~20 s

A declining curve means coordination cost is growing faster than capacity, so more hardware will make it worse. Fit contention and coherence coefficients to the sweep: contention-dominated means split the bottleneck; coherence-dominated means partition or stop sharing. Meanwhile cap concurrency at the measured peak.

open as a page

A team runs their whole workload on a work-stealing thread pool and sees poor CPU utilization and erratic tail latency. What properties of a workload make a work-stealing scheduler perform badly, and what would you change?

level: principalimportance: should knowfreq 26%

basics

~20 s

Work stealing assumes short, non-blocking, splittable, order-insensitive tasks. Blocking tasks park workers while their queued work sits unreachable; tasks too fine make steal and bookkeeping overhead dominate; tasks too coarse leave nothing to steal; and its last-in-first-out local order is deliberately unfair, so latency-sensitive requests need a separate fair pool.

open as a page

The Universal Scalability Law extends Amdahl's model of parallel speedup with a second penalty term. What two coefficients does it model, and why can it predict throughput that actually decreases as you add capacity?

level: seniorimportance: nice to knowfreq 22%

basics

~20 s

It adds contention (serialization on shared resources, cost grows with N) and coherence (keeping shared state consistent, cost grows with N squared because it is pairwise). Once the quadratic coherence term outgrows the linear capacity gain, throughput peaks and then falls.

open as a page

When a task in a parallel runtime spawns a subtask, the runtime can either keep running the parent and leave the subtask to be stolen, or immediately run the subtask and leave the rest of the parent to be stolen. Describe both strategies and the consequences of choosing one over the other.

level: seniorimportance: nice to knowfreq 17%

basics

~20 s

Child stealing: the spawned subtask goes on the deque, the spawning thread keeps running the parent. Continuation stealing: the thread jumps straight into the subtask and puts the parent's remainder (its continuation) on the deque. Continuation stealing matches sequential order and bounds outstanding tasks, but needs runtime or compiler support to capture a continuation, so libraries usually do child stealing.

open as a page

You are designing a compute-heavy batch job that must run on one machine today and across a cluster later. How would you decide between decomposing it by data and decomposing it by task, and what does that choice commit you to long term?

level: principalimportance: nice to knowfreq 32%

basics

~20 s

Decompose by data if you need the job to scale with input size — parallel width then grows with the data and survives the move to a cluster. Decompose by task only where you need to overlap a fixed set of independent operations. The data choice commits you to a partition key, to state that is per-partition, and to a skew and rebalancing story.

open as a page