Design a Job Scheduler for ML Training and Batch Inference on a Shared Cluster
Company: Netflix
Role: Software Engineer
Category: System Design
Difficulty: medium
Interview Round: Technical Screen
Design a job scheduler for machine learning workloads. Engineers and teams submit ML jobs, such as model training runs and batch inference jobs, that need compute resources including GPUs. The scheduler queues the jobs, decides when and on which machines of a shared cluster each one runs, tracks each job through its lifecycle, and handles failures.
Only the problem area is given. Settle the scope, scale and scheduling policy with the interviewer before you design.
```hint What makes ML jobs different
Compare a multi-GPU distributed training job with a web request or a cron task. Think about how it must be placed, how long it runs, and what happens to it if one of its machines fails.
```
```hint Contention for scarce GPUs
Ask what should happen when the cluster is full, a high-priority job arrives, and the running jobs belong to other teams.
```
### Clarifying Questions
- Which job types are in scope: training, hyperparameter sweeps, batch inference, interactive notebooks? Are any jobs recurring on a schedule or chained into pipelines?
- Do jobs span multiple machines, as in distributed training, and must all of their workers start together?
- Is the cluster shared by many teams, and are there quotas or priorities between them?
- May running jobs be preempted, and do jobs checkpoint so that they can resume?
- What scale should the design target: machines and GPUs in the cluster, jobs submitted per day, typical and maximum job duration?
- Is the hardware heterogeneous (several GPU types or pools), and can users request a specific type?
### What a Strong Answer Covers
- Requirements that capture what is specific to ML jobs: GPU requests, all-or-nothing placement for distributed jobs, long runtimes, and checkpointing
- A job lifecycle with explicit states and a durable job store as the source of truth
- A scheduling policy covering priority, fairness or quotas between teams, preemption, and starvation of large jobs
- Placement that accounts for fragmentation, GPU type and network topology
- Failure handling for lost nodes, scheduler crashes, retries and duplicate launches
- Observability for users ("why is my job still pending?") and for operators (utilization, queue wait)
### Follow-up Questions
- A large distributed training job keeps waiting while small jobs take every GPU that frees up. How do you fix that without leaving GPUs idle?
- How would you support recurring pipelines in which a training job must wait for an upstream data job to finish?
- How would the design change if the cluster grew to several pools or regions?
- Would you build this on an existing container orchestrator or write your own scheduler, and what decides it?
Overview: Design a job scheduler for machine learning workloads such as training runs and batch inference on a shared compute cluster with GPUs. It tests requirements gathering, queueing and fairness between teams, gang placement for distributed jobs, preemption and checkpointing, and recovery from node and scheduler failures.
Read the full Netflix Software Engineer interview experience this question came from