skip to content

When does evaluating a pure transform over partitioned chat history in parallel actually pay?

level: seniorimportance: nice to knowfreq 33%

answer

  1. independent partitions first
  2. per-partition work versus fixed cost
  3. the slowest partition sets the finish
  4. an associative merge keeps workers apart
  5. the sequential fraction caps the speedup

basics

~20 s

It pays when partitions are independent, each large enough to dwarf the cost of splitting and merging, reasonably even in size, and combined by an associative merge. Skew and the sequential fraction, not worker count, set the floor on wall-clock time.

solid answer

~50 s

Four conditions, and all of them have to hold. **Independence**: a partition's result must not need data from another partition. **Grain**: per-partition work must dominate the fixed cost of splitting, scheduling and merging, or you spend more arranging the work than doing it. **Balance**: the run finishes when the *largest* partition finishes, so one room holding most of the history caps the gain no matter how many workers you add. **An associative merge**: partial results must combine in any grouping, which is what lets the combination itself be parallel. Underneath all four sits Amdahl's law — reading the input, merging and writing the result are sequential, and that fraction bounds the speedup however wide you go. The shape that satisfies all of this is a pure per-partition fold; a shared accumulator that every worker writes turns the transform back into guarded mutable state and the contention eats the gain.

code

pseudocode · 7 lines
pseudocode
// fits: each partition folds into its own partial result
partials = parallel_map(partitions_of(history), count_by_sender)
totals   = reduce(partials, merge_counts)   // merge is associative

// does not fit: one shared table written on every record
for each message in history in parallel
    counts[message.sender] = counts[message.sender] + 1

go deeper

for a junior

Remember the basic requirement: each partition must be computable on its own, and the results must combine afterwards without workers talking to each other.

for a middle

Explain why a shared accumulator destroys the fit, and why the fixed cost of splitting and merging means small partitions can be slower than no parallelism at all.

for a senior

Diagnose a disappointing run: measure skew and the sequential fraction, and say which of the two explains the flat speedup before anyone adds workers.

for a principal

Decide when a workload deserves this shape at all, given that reshaping a transform to be partitionable and associative is real engineering with its own maintenance cost.

Data-parallel evaluation is the fourth answer to concurrency in this comparison, and the one most tied to how the code is *shaped*. The claim the functional-leaning style makes is that a transform with no shared writes can be evaluated on many partitions at once with no coordination at all. That claim is true, and it is conditional. The conditions are what an interviewer is probing. ## The shape that makes coordination unnecessary Counting messages per sender across years of chat history is the worked example. Two shapes are available: - **A per-partition fold.** Each worker folds its own slice into its own partial result, touching nothing anyone else touches, and the partials are merged afterwards. No worker ever waits for another. - **A shared accumulator.** Every worker updates one table of counts as it goes. This is the imperative shape, and it reintroduces exactly the guarded-mutable-state problem on the hottest possible path: one shared write per input record. It is often *slower* than the single-threaded version, because the coordination is paid per item while the useful work per item is tiny. The first shape is why this family fits data parallelism: purity is not a moral position here, it is the property that makes the absence of coordination *safe* rather than merely hoped for. ## The four conditions 1. **Independence.** The partition's result must be computable from the partition. Counting per sender within a room qualifies. "Which sender was most active across all rooms" still qualifies, because it becomes per-partition counts plus a merge — but "the rank of each message in a global ordering" does not, because every partition would need the others. 2. **Grain.** Splitting, dispatching and merging have a fixed cost. If a partition holds a few hundred records, that cost is comparable to the work, and the parallel version loses to the plain loop. Bigger partitions, fewer of them, is usually the better shape once independence is established. 3. **Balance.** This is the condition people skip. Wall-clock time is set by the partition that finishes last, not by the average. If one room holds most of the history, `p` workers do not give a `p`-fold speedup; they give roughly the time of the biggest partition. The fix is to split by *size* rather than by room, which is only possible if the transform tolerates a room being spread across partitions — a per-sender count does, a per-room running state does not. 4. **An associative merge.** Associativity means the partials can be combined in any grouping, so the merge can itself be a tree of parallel combinations rather than a left-to-right pass. Combining count tables is associative; "the first message of each day" is only associative if you keep enough information in the partial to decide the tie. ## Amdahl's law, stated in the right direction Every such job has a part that does not parallelise: reading the input, building the partitions, merging the partials, writing the result. If that part is a fraction `s` of the total, the speedup is bounded by `1 / s` however many workers you use. Ten percent sequential caps you at ten times, regardless of hardware. The consequence in practice is that shrinking the sequential part is usually worth more than adding workers — and that a job which is 40% merging will disappoint everyone who was told parallelism would fix it. ## Putting it against the other styles | | Fits data-parallel evaluation | Why | |---|---|---| | Pure per-partition fold | Yes | No shared writes, so workers never coordinate | | Immutable input, mutable partials | Yes | The partials are private to one worker for their whole life | | Shared accumulator across workers | No | One coordination point per input record | | Per-room isolated units | Partly | Natural partitions, but the merge crosses units | The fit sentence worth remembering: **data parallelism rewards a style where the work is a function of a partition and the combination is associative, and punishes one where the work is an update to something shared.** That is a statement about how the code is arranged, which is why it belongs to a comparison of styles rather than to a discussion of how many workers to run.

  • One room holds most of the history. What does that do to the wall-clock time?
    It sets it. The job ends when the largest partition ends, so extra workers idle once the small partitions are done and the speedup flattens near the ratio of total work to that partition's work. Splitting by size rather than by room fixes it, if the transform tolerates a room spanning partitions.
  • Why does an associative merge matter more here than a commutative one?
    Associativity lets partials be combined in any grouping, so the merge can be a parallel tree instead of one sequential pass. Commutativity additionally frees the order, which you need only if you merge partials as they land; merging them in partition order requires associativity alone.

saying these in an interview costs you the question

  • Assumes more workers always shorten the run
  • Reasons from the average partition size and ignores skew
  • Keeps a shared accumulator and still calls the transform pure
  • Thinks an order-dependent merge parallelises as easily as an associative one
  • Treats splitting and merging as free next to the work itself