skip to content

A CI test suite that takes 30 minutes on one job is split across 10 parallel jobs, but wall-clock time only drops to 12 minutes. What explains the gap, and how would you choose the split count?

level: seniorimportance: should knowfreq 45%

answer

  1. fixed cost per job does not divide
  2. slowest shard sets the wall clock
  3. split by recorded timings, not file count
  4. one database does not scale with shards
  5. flake probability compounds across shards

basics

~20 s

Each parallel job repays a fixed cost — queue wait, checkout, dependency restore, image pull, service startup — and that cost does not divide. Add unbalanced shards and shared bottlenecks and the speedup flattens. Choose the split where the next shard's saving stops exceeding that fixed cost.

solid answer

~50 s

Wall clock for a split suite is roughly *fixed cost + longest shard*, not *total time ÷ shards*. The fixed cost is everything each job repeats before it runs a single test: waiting for a runner, cloning, restoring dependencies, pulling a container image, starting a database, running migrations. Ten shards pay that ten times in money and once in latency, and it sets a floor nothing can go below. On top of that, shards are usually unbalanced — splitting by file count rather than by recorded per-test durations leaves one shard running twice as long as the rest, and the slowest shard is your wall clock. Then there are shared bottlenecks: one database instance, a rate-limited external service, a licence server. To size it, I measure the fixed cost per job and the serial test time, then pick the point where the marginal shard saves less than it costs. Usually the bigger win is attacking the fixed cost and the slow tests first.

go deeper

for a junior

Know that every parallel job repeats the setup work — checkout, dependency install, service startup — so ten jobs are not ten times faster than one.

for a middle

Explain the model out loud: wall clock is fixed cost plus the slowest shard, and describe why timing-based splitting beats splitting by file count.

for a senior

Diagnose a real gap: separate overhead from imbalance from shared bottlenecks, quantify the fixed cost from actual runs, and order the fixes so overhead and slow tests are attacked before more shards are bought.

for a principal

Own the tradeoff across teams: shard counts sized against pool capacity and concurrent demand, a cost-per-merge view rather than a single-pipeline stopwatch, and a stance on the flakiness budget that high parallelism demands.

## The arithmetic people expect versus the arithmetic they get The naive model is `wall = total / N`. The real model is closer to: ``` wall ≈ queue_wait + setup_per_job + (test_time / N) + merge_step ``` with the additional constraint that the shard that finishes last, not the average shard, determines the result. Plug in the numbers from the scenario: if setup is four minutes and the tests are 26 minutes of actual execution, then ten perfectly balanced shards give `4 + 2.6 ≈ 6.6` minutes — and the observed 12 means something else is also wrong, almost always imbalance. The important consequence: as N grows, `test_time / N` shrinks toward nothing while `setup_per_job` does not move at all. Past a certain N you are buying compute and getting no latency back. ## Where the fixed cost actually comes from Enumerate it, because "overhead" is too vague to act on. Per job: waiting for a free runner in the pool; provisioning the machine or container; cloning the repository (a deep clone of a large repo is minutes on its own); restoring or re-downloading dependencies; pulling a container image; starting service containers such as a database or message broker; running schema migrations and seeding fixtures; and at the end, uploading results. Each of those is separately attackable, and the attack usually beats adding shards: shallow clones, a warm dependency cache, smaller images or pre-baked runner images, a pre-migrated database template. Cutting setup from four minutes to ninety seconds improves *every* shard count and reduces the floor itself. ## Imbalance is usually the actual culprit Splitting by file count assumes files cost the same, which they never do — one integration test file can outweigh three hundred unit tests. The standard fix is timing-based splitting: record per-test durations from previous runs, then bin-pack tests into shards of roughly equal predicted duration. Most test runners and CI splitters support this; the input is a stored timings file, and the failure mode to watch is a stale one, where newly added slow tests all land in whichever shard the fallback rule chooses. There is also an irreducible tail: a single test, or a single non-parallelisable file, that takes longer than your target. No split count helps. That test has to be made faster or moved out of the critical path. ## Shared bottlenecks and contention Ten shards hitting one database instance do not get ten times the database. The same applies to a rate-limited third-party sandbox, a shared licence server, a single artifact store on a saturated link, and the runner pool itself — if the pool has twelve executors and you request ten for every pull request, the second concurrent PR queues, and the *team's* median wait gets worse even though one pipeline got faster. That last effect is invisible in a single pipeline's metrics and very visible in engineers' complaints. ## Flakiness compounds If each shard independently has a small chance of hitting a flaky failure, the probability that *at least one* shard fails grows with the shard count: at 1% per shard, ten shards fail about 10% of the time, twenty about 18%. A pipeline that needs a retry every other run has a much worse effective latency than its green-path number suggests. High parallelism therefore raises the bar on test hygiene rather than lowering it. ## How to choose N in practice 1. Measure `setup_per_job` and total serial test time from real runs; do not estimate. 2. Compute the predicted wall clock for a few candidate N values and compare against the observed — a large discrepancy means imbalance, not overhead. 3. Choose N around the point where the marginal shard saves less time than the setup cost it adds; in most suites that lands well under ten. 4. Sanity-check against pool capacity and concurrency: what happens when five pull requests run at once? 5. Re-derive it periodically, because both terms drift. And state the ordering explicitly, because it is the part that separates a senior answer: reduce the fixed cost first, fix the slowest tests second, balance the shards third, and only then buy more parallelism. Parallelism is the one lever that costs money linearly while returning less and less.

  • How do you split shards so they finish at roughly the same time?
    Bin-pack by recorded per-test duration from previous runs rather than by file count, since a single integration file can outweigh hundreds of unit tests. Keep the timings file fresh and have a sane rule for tests with no recorded duration, otherwise every newly added slow test lands in the same shard and quietly recreates the imbalance.
  • Your runner pool has twelve executors and each pull request now asks for ten. What happens?
    One pipeline gets faster and the team gets slower. Concurrent pull requests queue behind each other, so median time-to-feedback rises even though the single-pipeline number improved. Parallelism has to be sized against pool capacity and expected concurrency, not against one pipeline in isolation.
  • Where would you look first if you wanted the biggest win, before touching the shard count?
    The fixed per-job cost, because it improves every shard count and lowers the floor itself: shallow clones, a warm dependency cache, smaller or pre-baked images, a pre-migrated database template. Then the slowest individual tests, which set an irreducible tail no split count can cross.
  • Why does high parallelism make flaky tests more painful?
    The pipeline fails if any shard fails, so independent per-shard flake probabilities compound. At roughly 1% per shard, ten shards fail about one run in ten. The green-path wall clock looks great while the effective time-to-merge includes frequent retries, which is often a net loss.

saying these in an interview costs you the question

  • Assumes wall clock is total time divided by shard count
  • Splits shards by file count and calls it balanced
  • Ignores that every shard repeats checkout and setup
  • Adds shards without checking runner pool capacity
  • Treats the average shard rather than the slowest as the result

context