skip to content

How do you shard a long-running automated suite across parallel workers so that wall-clock time actually drops?

level: seniorimportance: nice to knowfreq 24%

answer

  1. Wall clock, not total machine time
  2. The longest shard sets the run
  3. Balance by measured duration
  4. Shards must not share mutable data
  5. One long case is a floor

basics

~20 s

Wall clock is set by the slowest shard, so balance shards by each case's measured duration from previous runs rather than by name or count, and make shards genuinely independent - no shared mutable data, fixed identifiers, ports or accounts.

solid answer

~50 s

Sharding helps only if the shards finish together and are genuinely independent. Wall clock is the longest shard, not the average, so splitting by file name or by case count leaves stragglers: balance instead by each case's recorded duration from previous runs, either by packing shards greedily longest-case-first or by having workers pull from a shared queue. Measure the critical path rather than total processor time - one very long case is a floor no shard count beats, so it must be split or moved to a cheaper level. Independence is the harder half: shards must not share mutable records, fixed identifiers, ports or accounts, or you get failures that appear only under a particular interleaving. Aggregate results and artefacts into a single report, keep any per-shard retries visible rather than silently absorbed, and remember that shards buy wall clock with environment cost - fix the feedback budget first, then buy the shards that fit it.

code

pseudocode · 15 lines
pseudocode
# balance by measured duration, longest case first
durations = load_durations_from_previous_runs()   # caseId -> seconds
cases     = sort(all_cases, by = durations, descending = true)
shards    = [ new_shard(i) for i in 1..8 ]

for case in cases:
    target = shard_with_smallest_total(shards)
    target.add(case, durations[case.id])

# the run takes the longest shard, not the average
wall_clock = max(shard.total for shard in shards)

# each shard gets its own data namespace so records cannot collide
for shard in shards:
    shard.identifierPrefix = "shard-" + shard.index + "-"

go deeper

for a junior

Know the basic shape: splitting a suite across several workers can shorten the wait, and the run is not finished until the slowest worker is. You are not expected to design the balancing at this level.

for a middle

Be able to explain why an even split by count is not an even split by time, and what a shard needs in order not to interfere with its siblings - its own data namespace, no fixed ports, no ordering assumptions.

for a senior

Demonstrate that you measure the critical path and the duration distribution before buying workers, that you balance from recorded durations, and that you aggregate results and keep per-shard retries visible in the report.

for a principal

Frame it as budget and cost. Decide the feedback budget per tier, weigh worker cost against pruning and relocating cases, and say plainly when the answer is a smaller suite rather than a wider split.

## Start from the budget, not the shard count A suite has a feedback budget: the time a team is willing to wait before the result changes what they do next. A change-gating suite might get ten minutes; a nightly pack might get a night. Sharding is one of several instruments for fitting a suite into its budget, and it is the one people reach for first because it requires no decisions about the cases themselves. That is exactly why it is worth being precise about when it works. ## Wall clock is the longest shard The number that matters is the critical path, not the total. Split 3,182 cases eight ways by file name and the shards will not be equal, because case durations are not uniformly distributed: a handful of long journeys sit somewhere in the alphabet. One real suite over a clinical-trial data capture form split this way produced shards ranging from 41 minutes to 96 minutes. Total processor time was unchanged and the dashboard reported a healthy eight-way split, but the run still took 96 minutes, because seven workers sat idle waiting for the eighth. The perfectly balanced split would have been about 47 minutes each. Two balancing strategies fix that. **Static balancing by measured duration.** Record each case's duration on every run, then assign cases to shards greedily, longest first, always into the shard with the smallest running total. This is the classic longest-processing-time heuristic and it gets close to optimal for this shape of problem. It needs a durable store of durations keyed by a stable case identifier - which is why renaming cases wholesale quietly degrades balance until the next run's data arrives. **Dynamic balancing by pull.** Put every case on a shared queue and let each worker take the next one when it becomes free. This self-corrects without historical data and copes with machines of different speeds, at the cost of a coordinator and of losing the ability to know in advance which case runs where. A common hybrid pulls in duration order, longest first, so the long tail starts early. ## The floor no sharding beats If one case takes 38 minutes, no shard count gets the run under 38 minutes. Before adding workers, look at the distribution: the top 2 percent of cases by duration often own a quarter of the run. That time usually goes into setup that could be shared, waiting that could be signalled rather than slept through, or a journey that is really four independent checks glued together. Splitting or relocating a handful of cases can beat doubling the worker count, and it costs nothing per run afterwards. ## Independence is the hard half Shards are only parallel if they cannot interfere. The recurring collisions are mundane: - **Shared mutable records.** Two shards creating a subject record with the same identifier will fail in a pattern that neither shard reproduces alone. Give each shard its own namespace - an identifier prefix, a schema, a tenant - or generate identifiers that cannot collide. - **Fixed resources.** A hard-coded port, a single service account with a session limit, one seeded login used by every case, a single output directory. - **Ordering assumptions.** A case that only passes because another case ran first was already fragile; sharding is simply the moment it is discovered, since the two land on different workers. - **Global state in the system under test.** A feature flag toggled by one case, a locale or clock setting changed globally, a cache cleared mid-run. A useful preparation step before sharding at all is to shuffle case order within a single worker: anything that fails under shuffling will fail under sharding too, and diagnosing it in one process is far easier. ## Reporting and cost Each shard produces its own results and artefacts, so publish them into one aggregated report keyed by case, not per worker - otherwise triage begins with finding which of eight reports holds the failure. Name artefacts with the shard and run so they cannot overwrite one another. If shards retry internally, keep the retry visible in the aggregate: a run that is green only because one shard passed on its third attempt is not the same result as a clean run, and hiding that turns a maintenance signal into noise. Finally, shards cost money and environments. Eight workers usually means eight application instances or eight data namespaces, and the marginal minute saved gets more expensive as the split gets finer. Past a point the honest answer is not more shards but fewer cases: prune duplication, move checks to a cheaper level, and reserve the wide split for the tier that genuinely has to exercise the whole system.

  • When is splitting or relocating a few cases better than adding more workers?
    Whenever the duration distribution is long-tailed. If the slowest case is 38 minutes, no shard count gets the run below that, and the top few percent of cases often own a quarter of the runtime. Splitting a bloated journey or moving its checks to a cheaper level cuts the floor permanently, while extra workers cost money on every run forever.
  • What is the tradeoff between assigning cases to shards up front and letting workers pull from a queue?
    Up-front assignment needs a store of historical durations and stable case identifiers, but the plan is reproducible and you know in advance where each case runs. Pulling from a queue self-corrects for uneven durations and unequal machines without any history, at the cost of a coordinator and of a run that is harder to reproduce case-for-case.
  • How would you check a suite is safe to shard before you shard it?
    Shuffle case order within a single worker and run it. Anything that depends on ordering, on a record another case created, or on global state left behind will fail there too, and diagnosing it in one process is far cheaper than chasing an interleaving across eight. Then audit for fixed ports, shared accounts and hard-coded identifiers.

saying these in an interview costs you the question

  • Reports total machine time as if it were wall clock
  • Splits by file name or case count and calls it balanced
  • Adds workers while one case dominates the runtime
  • Assumes separate processes cannot collide on shared data
  • Leaves per-shard retries out of the aggregated report
  • Ignores the environment cost of every extra shard

context