skip to content

Once the affected set of projects is known, how do build tools schedule the resulting build/test/lint tasks — what determines which tasks run in parallel versus which must wait, and how does this scale across multiple CI machines?

level: seniorimportance: must knowfreq 65%

answer

  1. task DAG derived from project graph + intra-project ordering
  2. eligible = predecessors done
  3. critical path caps parallelism benefit
  4. distributed execution = network workers not just cores
  5. historical-duration-based load balancing / sharding

basics

~20 s

Tasks form a chain based on which project needs which other project built first. Tasks with no unfinished dependency can run at the same time; tasks that depend on another task's output have to wait for it. Big pipelines split this work across several machines to go faster.

solid answer

~50 s

The affected projects' tasks (build, test, lint) form a task DAG derived from the project dependency graph plus intra-project task ordering rules (e.g., 'test depends on build'). The scheduler topologically sorts this DAG and dispatches every task whose dependencies have completed to any free worker, so independent branches of the graph run concurrently while dependent chains execute in order — parallelism is bounded by the DAG's width at each point, not just by available cores. At larger scale, task execution is distributed across multiple CI agents/machines rather than one box's cores, often via a remote execution service, with tasks scheduled to workers based on availability and sometimes historical duration to balance load. The main scheduling risk is a long dependency chain (a 'critical path') through a small number of heavily-depended-on projects, which caps how much wall-clock time parallelism alone can save regardless of worker count.

go deeper

for a junior

Should grasp that some tasks can run at the same time while others have to wait for a prerequisite to finish.

for a middle

Should describe topological scheduling — eligible tasks (predecessors done) get dispatched to free workers — and recognize that parallelism has a width limit determined by the DAG shape.

for a senior

Should explain the critical-path concept as the hard floor on wall-clock time and know that distributing tasks across multiple CI machines introduces network/environment concerns beyond simple core-count scaling.

for a principal

Should reason about restructuring the dependency graph or optimizing specific bottleneck projects to shorten the critical path, and discuss load-balancing strategies like historical-duration-based sharding across a worker fleet.

## From affected set to task graph Once the affected set of projects is fixed, the next problem is ordering and parallelizing the actual work — building, testing, and linting each affected project. This is governed by a **task graph**, which is a finer-grained structure than the project dependency graph used for affected detection, though it's derived from it. The task graph has two sources of ordering constraint: - The first is **inter-project**: if project A depends on project B, then A's build task can't start until B's build task (or at least B's build artifact) is available, because A needs to compile against B's output. - The second is **intra-project**: within a single project, a test task typically depends on that project's own build task completing, and a lint task might run independently of both. Combining these constraints across every affected project produces a single task-level **directed acyclic graph**, where nodes are individual tasks (not projects) and edges are 'must complete before.' ## Topological execution and DAG width Given this DAG, the scheduler's job is a topological execution: at any moment, every task whose predecessors have all completed is eligible to run, and eligible tasks are dispatched to whatever workers are free. This means parallelism isn't just 'run N tasks per core' — it's bounded by the **DAG's width**, i.e., how many tasks are simultaneously eligible at each point in the execution. - Early in a run, if many affected projects have no inter-dependencies, dozens of build tasks can start at once. - Later, if several projects all depend on one shared library, execution narrows to a bottleneck until that library's task finishes, then fans back out. ## The critical path The key trade-off scheduling has to manage is the **critical path**: the longest chain of sequentially-dependent tasks through the graph. No matter how many workers are available, total wall-clock time can't drop below the critical path's cumulative duration, because those tasks are strictly ordered. A monorepo where every project transitively depends on one 'core' library has a critical path that funnels through core's build/test every single run, capping the benefit of adding more parallel workers — this is a common real bottleneck teams hit and a reason for splitting overly-central libraries or optimizing their build/test time specifically, since shaving time off the critical path helps every downstream run, while shaving time off a parallel branch only helps if it was itself the local bottleneck. ## Scaling out to a fleet At CI-fleet scale, this scheduling problem extends beyond a single machine's cores to distributing tasks across many CI agents or a remote-execution cluster. Tools like Bazel's remote execution API, or cloud task-distribution services (Nx Cloud, BuildBuddy, or custom sharding built on top of GitHub Actions/Buildkite matrix jobs), take the eligible-task set and dispatch it to a pool of workers rather than local cores, so the effective parallelism scales with fleet size rather than one machine's CPU count. This introduces its own scheduling concerns: - workers need the right toolchain/environment; - task dispatch has network and queueing latency; - load balancing matters — a naive round-robin dispatch can leave one worker with a string of slow tasks while others sit idle, so production schedulers often use estimated task duration (from historical run data) to bin-pack tasks more evenly, a technique sometimes called 'historical timing-based sharding.' ## What the tools actually do A concrete real-world pattern: Nx Cloud's distributed task execution splits the affected task graph across a pool of agent machines, respecting the same dependency-ordering constraints, and uses prior run durations to balance the split so no single agent becomes the long pole. Bazel's remote execution similarly fans build/test actions out to a worker pool while enforcing the action graph's ordering via input/output artifact dependencies. In both cases, the underlying discipline is identical to single-machine scheduling — topological ordering plus eligible-task dispatch — just with network-distributed workers substituting for local threads, and the critical path remaining the hard floor on total pipeline time regardless of how many workers you add.

  • If a team notices CI wall-clock time isn't improving even after doubling the number of parallel workers, what's the most likely explanation?
    The pipeline has hit its critical path — the longest chain of sequentially-dependent tasks — and additional workers can only help with tasks that are currently sitting idle waiting for a free worker, not with tasks that are blocked on a dependency. The fix is usually to shorten the critical path itself: speeding up the slow task(s) on that chain, or restructuring the dependency graph to reduce how much funnels through it, rather than adding more workers.
  • Why might a lint task for a project run independently of that same project's test task, even though both belong to the same project?
    Lint typically only needs the source files, not a built artifact, so it has no ordering dependency on the build task the way tests usually do. Tools that model tasks at this fine a grain can schedule lint to start immediately in parallel with build, rather than serializing it behind build-then-test, squeezing extra parallelism out of the same affected set.
  • What's a risk specific to distributing tasks across multiple physical CI machines rather than running them on one box's cores?
    Environment consistency and caching locality: each worker needs the same toolchain/dependencies available, and artifacts produced by one task on one machine may need to be fetched over the network by a dependent task scheduled on a different machine, adding latency that doesn't exist when everything shares one filesystem. Poorly configured remote execution can end up network-bound rather than compute-bound.

It's like a kitchen prepping a multi-course meal: dishes with no shared ingredients get cooked simultaneously on different burners, but any dish that needs a stock made first has to wait for that stock — no matter how many cooks you hire, dinner can't be served faster than the longest chain of 'make stock, then reduce it, then finish the sauce.'

saying these in an interview costs you the question

  • thinks parallelism is unlimited if you just add more workers/cores
  • doesn't recognize the concept of a critical path capping wall-clock time
  • assumes every affected task can always run simultaneously with no ordering constraints
  • conflates the project dependency graph with the finer-grained task DAG

context