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?
answer
- where does the width come from?
- data width grows with input; task width is fixed
- the partition key is the architecture
- associative combine = same program on 1 or 100 machines
- over-partition + consistent hashing = cheap rebalance
basics
~20 sDecompose 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.
solid answer
~1 minAsk where the parallel width comes from. **Task decomposition** gives you a width fixed by the program's structure — five independent stages means at most five-way parallelism, and the critical path is a floor you cannot buy your way past. **Data decomposition** gives you width proportional to input, so it is the only route that keeps paying as the data and the cluster grow. For a compute-heavy batch job that must scale out, make data the primary axis and task the secondary one: task-parallel stages, each internally partitioned. The long-lived commitment is the **partition key**. It determines what state a worker can own locally (anything keyed the same way), what forces a shuffle (any aggregation across a different key), where skew appears (a hot key becomes a hot worker), the unit of retry and of failure, and how you rebalance when you add machines — which is why consistent hashing or many-more-partitions-than-workers matters from day one. Also decide early whether the combine is associative and commutative; if it is, the same code works on one machine or a hundred, and a reduce tree replaces a serial merge. If it isn't, you have baked in a sequential bottleneck that a cluster will make worse rather than better.
code
text · 10 linesSERIAL MERGE (hidden scalability wall)
K partitions -> each emits raw results -> one node sorts/dedups everything
network volume ~ N, merge time ~ N on a single node
adding machines speeds the map half only; the merge term stays
ASSOCIATIVE REDUCTION (scales out unchanged)
each partition pre-aggregates locally -> emits only its partial result
partials combined pairwise in a tree -> log(K) rounds
network volume ~ K * |partial|, combine time ~ log K
same code path whether K = 4 on one box or K = 400 on a clustergo deeper
Recognise that splitting by data scales with the amount of data while splitting by task is limited to the number of independent tasks.
Argue the width and critical-path point clearly and note that the two compose: parallel stages with partitioned work inside each.
Reason about the partition key: state locality, skew, retry granularity, and whether the combine is associative enough to replace a serial merge with a reduce tree.
Treat the decomposition as an architectural commitment — partition key, rebalancing strategy, idempotent retry, shuffle boundaries, observability of per-partition skew — and design so the single-machine and cluster versions are the same program.
## Frame the decision as "where does the width come from?" Every parallel design has a *width*: how many things can genuinely run at once. The two decompositions supply it differently. - **Task decomposition** supplies width from program structure. If the job is "extract, enrich, score, write" you have four things, and the width is four minus whatever the dependency edges serialise. It does not grow with input. Its floor is the **critical path** — the longest chain of dependent tasks — which no amount of hardware reduces. - **Data decomposition** supplies width from the input. N records split into K partitions gives width K, and K can be raised as machines are added. For a job that must survive growth in both data volume and machine count, that asymmetry decides it: **data is the primary axis**. Task decomposition remains valuable as a secondary axis — it overlaps stages, hides I/O latency, and lets heterogeneous work (CPU-bound scoring, network-bound fetching) proceed simultaneously — but it should not be what you rely on for scale. ## The decision inputs 1. **Does width need to grow with input?** Yes → data. This is usually the whole answer for batch. 2. **Is the per-item work homogeneous?** Homogeneous work partitions cleanly; heterogeneous stages argue for a task-parallel outer structure with data parallelism inside each stage. 3. **What state does the work touch?** If processing an item requires lookups against state, the partition key should match the state's key so the state can live with the worker. If it cannot, every item pays a remote lookup and the design is network-bound rather than compute-bound. 4. **What must be aggregated, and is the aggregation associative?** Associative (and commutative) combines let you do local partial aggregation then a tree merge — the single most important property for scaling out, because it turns an O(N) serial merge into O(log K) rounds and cuts the data crossing the network. 5. **What are the ordering requirements?** Global ordering forces either a single consumer or a resequencing stage; per-key ordering is compatible with partitioning by that key, which is why per-key ordering is the usual production compromise. 6. **What does failure cost?** The partition is normally also the unit of retry. Big partitions mean expensive retries; small partitions mean cheap retries and better balance, at higher coordination cost. ## What the data-decomposition choice commits you to **The partition key is the architecture.** Once chosen, it determines: - *Locality of state* — anything keyed the same way can be held locally, uncoordinated. Anything keyed differently requires a shuffle or a remote call. - *Skew exposure* — a key with far more data than others becomes a worker that finishes far later. Mitigations (salting the hot key, pre-aggregating before the shuffle, a dedicated path) must be designed in, not bolted on. - *Rebalancing* — going from 8 workers to 40 must not require reprocessing everything. Two standard defences: create many more partitions than workers from the start (so growth is reassignment, not re-partitioning), and use consistent hashing so adding a node moves a small fraction of the key space. - *Retry and idempotency* — retrying a partition must be safe, which means the write path has to be idempotent or transactional per partition. - *Observability* — per-partition duration and record counts are what tell you skew is happening, so they need to exist from the start. **The combine shape is the second commitment.** If you write a serial merge (sort everything at the end, deduplicate on one node), you have created a fixed serial fraction. On one machine it is invisible; on forty it becomes the dominant term and the job stops improving. Designing the combine as an associative reduction — with local pre-aggregation so the volume crossing the network drops by orders of magnitude — is what makes the single-machine version and the cluster version the same program. ## What task decomposition commits you to A dependency graph. Its usefulness is bounded by width and critical path, and it is generally *cheaper* to change than a partition key — you can add or split a stage without reprocessing data. That asymmetry is a good reason to fix the partitioning question early and let the stage structure evolve. The risk with task decomposition is treating it as a scalability plan. Adding machines to a five-stage pipeline does not help unless each stage is internally replicated — which is data parallelism again, now with the extra cost of moving data between stages. ## The hybrid that most real systems land on - Outer level: stages, connected by bounded queues, running concurrently (task parallelism) — good for overlapping heterogeneous latencies and for isolating failure. - Inner level: each stage replicated across partitions of the key space (data parallelism) — this is where scale comes from. - Between stages: the shuffle boundary, which is the expensive part, so the design minimises how many times the key changes. ## How to present this in an interview Don't lead with a technology. Lead with: what is the width, where does it come from, what is the partition key, is the combine associative, where is the skew, what is the retry unit, and what changes when you add machines. Then note the one-machine-today constraint: pick a design whose *shape* is unchanged by the move — partition-local state, associative combine, many more partitions than workers — so the cluster version is a deployment change and not a rewrite.
- What makes the partition key so hard to change later?Because state placement, locality, retry units and the shuffle boundary are all derived from it. Changing it means re-partitioning existing state, invalidating the assumption that a worker can serve lookups locally, and reprocessing or migrating data. Stage structure, by contrast, can usually be edited without touching stored data, which is why the partitioning decision deserves the up-front rigour.
- How do you keep a design that runs on one machine today from needing a rewrite when it moves to a cluster?Keep the shape invariant: partition from the start with many more partitions than workers, make all state partition-local, make the combine associative and commutative with local pre-aggregation, make partition processing idempotent so retries are safe, and never introduce a single-node merge that sees all the data. Then scaling out is a change in how many workers claim partitions, not a change in the algorithm.
- When is task decomposition still the right primary choice?When the job's cost is dominated by a small number of independent waits rather than by volume — for example a request that must call three services and merge — or when stages are genuinely heterogeneous in resource type and you want them overlapped and isolated. It is a latency-hiding and structuring tool. It is the wrong primary choice whenever you need parallel width to grow with the input.
saying these in an interview costs you the question
- Planning to scale a fixed set of pipeline stages by adding machines, when task width is capped by structure.
- Choosing a partition key from convenience without checking skew, state locality, or the retry unit.
- Designing a single-node merge step that sees all results, which becomes the dominant serial fraction at scale.
- Creating exactly one partition per worker, leaving no room to rebalance when machines are added.
- Assuming a combine is safely parallel without confirming it is associative — and commutative if partitions can finish in any order.