The Forward Deployed

OpenAI Interview: Design a Video Generation Pipeline on Scarce GPUs

A full solution to the OpenAI video generation question: an async job API, the job state machine, pull scheduling with leases, fencing tokens, checkpoints and preemption, admission with honest ETAs, worker pools per model version, delivery, and failure handling.

By Reviewed

Part of the OpenAI system design question bank. The question is representative of the round. The analysis and solution are this site's own.

Problem statement

A user submits a prompt and gets a job ID at once. They watch progress, can cancel, and eventually download a video. Behind the API, a scheduler assigns jobs to GPU workers, each running one video at a time. Workers may be preempted, may crash, and may lose their network. Accepted jobs must never be lost.

Clarifying questions

  • How long does one video take? For practice: about 4 minutes of GPU time on average, up to 15 for long videos.
  • How many submissions? 500,000 a day, with peaks of 10 per second.
  • One video per GPU worker? Yes, and a worker holds one model version loaded.
  • What kinds of GPUs? A mix of reserved GPUs and cheaper preemptible ones that can vanish with a short warning.
  • What latency targets? Job creation acknowledged in under 500 ms. Scheduling within 30 seconds when capacity exists.
  • Priorities? Paid and free tiers.

What makes this hard

The work is long, expensive, and runs on machines you cannot trust to stay up. A four-minute job on a preemptible GPU has a real chance of losing its machine partway. Restarting from zero wastes GPU minutes, which are the most expensive resource in the system.

Meanwhile, demand often exceeds supply. The GPU pool is smaller than peak demand, and it shrinks when preemptible capacity is reclaimed. Users will wait, but they must know how long, and the queue must stay fair.

And distributed failure is subtle. A worker that loses its network looks dead, so its job is given to another worker. Then the first worker comes back and tries to write its result. Without care, two workers both "finish" the same job, and one overwrites the other.

So the driving tension is GPU utilization versus reliability. Using cheap, preemptible capacity and keeping every GPU busy saves money; every such choice makes failure more common, so recovery must be cheap and correct.

flowchart LR
  C([Client]):::user -->|submit| API[Job API]:::svc --> DB[(Jobs)]:::store
  W1[GPU worker]:::svc -->|lease a job| DB
  W2[GPU worker]:::svc -->|lease a job| DB
  W1 --> OS[(Checkpoints + videos)]:::store
  C -->|status, download| API
  classDef user fill:#e6efec,stroke:#315e55,color:#171717;
  classDef svc fill:#f4f1e8,stroke:#315e55,color:#171717;
  classDef store fill:#fdf3dc,stroke:#c4492d,color:#171717;
Key idea. The metadata system is small. GPU time is the cost. Design so a lost worker costs seconds of work, never a lost job and never a double result.

Key concepts

Asynchronous jobs

A request that takes minutes should not hold an HTTP connection open. The API accepts the job, returns an ID, and the client polls or subscribes for status.

Leases

A lease is a time-limited claim on a job. A worker holding a lease must renew it with heartbeats. If it stops, the lease expires and the job becomes available again. Leases turn "is the worker dead?" from an unanswerable question into a timeout.

Fencing tokens

A fencing token is a number that increases every time a job is leased. Every write from a worker includes its token, and the store rejects writes with an old token. A worker that lost its lease can no longer change the job, even if it wakes up and tries.

Checkpoints

A checkpoint saves the job's intermediate state, so a new worker can resume from it. It bounds the work lost to any failure to the time since the last checkpoint.

Little's law

The number of jobs in progress equals the arrival rate times the average duration. It sizes the GPU pool.

Key idea. Async API, leased work, fenced writes, checkpointed progress, and Little's law for capacity.

  1. Requirements

Before reading on. List the requirements. Then name the one guarantee you would put in the contract with users.

1.1 Functional requirements

  • Submit a job with prompt, model version, length, and resolution; get a job ID.
  • Get status, progress, and an estimated time; download the video when done.
  • Cancel a queued or running job.
  • Assign each job to one GPU worker with the right model loaded.
  • Recover from worker loss, preemption, and network failures without losing accepted jobs.

1.2 Non-functional requirements

  • Durability of accepted jobs. An accepted job always finishes, fails with a reason, or is cancelled.
  • Exactly one result. A job never produces two competing results.
  • Scheduling latency under 30 seconds when capacity exists.
  • Bounded lost work. Under 30 seconds of GPU work lost per preemption.
  • Honest ETAs. Users see a wait estimate that tracks reality.

1.3 The constraint versus the property

No lost jobs, one result each, is the property. GPU capacity is the constraint. It is scarce, changing, and partly preemptible, and it sets the design of scheduling, checkpointing, and admission.

  1. Back-of-the-envelope estimation

2.1 Concurrency and GPUs

Peak arrival: 10 jobs per second. Average GPU time: 240 seconds. Little's law: 10 × 240 = 2,400 jobs running at peak, so about 2,400 busy GPU workers plus a warm buffer of 10%, about 2,650. The daily average is lower: 500,000 / 86,400 ≈ 5.8 per second, or about 1,400 busy workers.

2.2 Control plane

  • Job creations: 10 per second at peak.
  • Heartbeats: 2,400 workers every 10 seconds = 240 per second.
  • Progress updates: sent with heartbeats.
  • Status polls: if 50,000 users poll every 5 seconds, 10,000 reads per second, served from a cache.

A single well-indexed relational database handles the writes. Status reads go to a cache.

2.3 Storage

  • Final videos: 500,000 × 20 MB = 10 TB a day.
  • Checkpoints: several hundred MB each, written every 30 seconds while running: 2,400 × 300 MB / 30 s ≈ 24 GB/s at peak. That is a lot of write bandwidth; keep checkpoints on fast object storage in the same zone, keep only the latest two per job, and delete them when the job finishes.

2.4 Lost work

Without checkpoints, a preemption at a random point of a 4-minute job loses 2 minutes on average. With checkpoints every 30 seconds, it loses 15 seconds on average. If 5% of jobs are preempted, checkpoints save about 25,000 × 105 seconds ≈ 730 GPU-hours a day.

Key idea. 2,400 GPUs at peak, a small control plane, 10 TB a day of output, and checkpoints that save hundreds of GPU-hours a day.

  1. API design

Before reading on. A client's submit request times out. It retries. How do you avoid generating the video twice?

The client sends an idempotency key with the submit. The server stores it with the job ID. A retry with the same key returns the existing job instead of creating a new one.

POST /v1/videos
Idempotency-Key: <uuid>
{prompt, model_version, seconds, resolution, priority?}
-> 202 {job_id, status: "queued", eta_seconds}
-> 429 {reason: "too_many_queued"}          per-user cap
-> 503 {reason: "capacity", retry_after}    admission rejected (free tier under load)

GET  /v1/videos/:id
-> {status, progress: 0.0-1.0, eta_seconds, error?, result_url?}

POST /v1/videos/:id/cancel  -> {status}
GET  /v1/videos/:id/events  -> SSE stream of status and progress

Internal worker API:

POST /internal/lease      {worker_id, model_version, gpu_type}
  -> {job_id, attempt, params, checkpoint_ref?, lease_expires_at} | 204 no work
POST /internal/heartbeat  {job_id, attempt, progress}
  -> {lease_expires_at, cancel: bool} | 409 lease lost
POST /internal/checkpoint {job_id, attempt, checkpoint_ref, step}
POST /internal/complete   {job_id, attempt, result_ref, metrics}
POST /internal/fail       {job_id, attempt, error, retryable}

Every worker call carries the attempt number, which is the fencing token.

  1. Data model

jobs
  id, user_id, tier, priority, idempotency_key,
  status: queued | leased | running | uploading | succeeded | failed | cancelled,
  model_version, params_json,
  attempt int,                       -- fencing token, +1 on every lease
  worker_id, lease_expires_at,
  checkpoint_ref, checkpoint_step,
  progress, result_ref, error,
  created_at, started_at, finished_at
  index (status, model_version, priority, created_at)   -- scheduler pick
  index (status, lease_expires_at)                      -- sweeper
  unique (user_id, idempotency_key)

workers
  id, model_version, gpu_type, preemptible, status, last_heartbeat_at

  1. High-level design

5.1 A synchronous endpoint

The client calls POST /generate and waits minutes for the video. Proxies time out, the client has nothing to resume if the connection drops, and a server crash loses the work with no record it ever existed.

5.2 Fix 1: an async job API and a durable queue

Accept the job into a database, return 202 with a job ID, and let workers process it. The job survives any API server crash.

5.3 Fix 2: workers pull with leases

Idle workers ask for work matching their loaded model. The scheduler leases the best matching job with one conditional update. Workers heartbeat to keep the lease. A sweeper returns expired leases to the queue.

flowchart LR
  W[Idle worker<br/>model v3 loaded]:::svc -->|lease request| S[Scheduler]:::new
  S -->|UPDATE ... SET status=leased, attempt=attempt+1<br/>WHERE status=queued ... LIMIT 1| DB[(Jobs)]:::store
  W -->|heartbeat every 10 s| S
  SW[Sweeper]:::new -->|expired leases back to queued| DB
  classDef svc fill:#f4f1e8,stroke:#315e55,color:#171717;
  classDef store fill:#fdf3dc,stroke:#c4492d,color:#171717;
  classDef new fill:#ffffff,stroke:#c4492d,stroke-width:2px,stroke-dasharray:5 3,color:#171717;

5.4 Fix 3: fencing on every write

Every worker write includes its attempt number. The database applies it only WHERE id = ? AND attempt = ?. A zombie worker's writes match nothing.

5.5 Fix 4: checkpoints and resume

Workers save generation state to object storage every 30 seconds and on a preemption warning. A new lease includes the latest checkpoint, and the new worker resumes from it.

5.6 Fix 5: admission and ETAs

The API estimates each new job's wait from the queue ahead of it and the recent completion rate, returns it, and rejects free-tier jobs when the estimate passes a threshold.

5.7 The composed design

sequenceDiagram
  autonumber
  actor U as User
  participant A as Job API
  participant D as Jobs DB
  participant W as GPU worker
  participant O as Object storage
  U->>A: POST /videos (idempotency key)
  A->>A: admission: ETA under threshold?
  A->>D: insert job (queued)
  A-->>U: 202 {job_id, eta}
  W->>A: lease {model v3}
  A->>D: conditional update: leased, attempt=2
  A-->>W: job, attempt 2, checkpoint (if any)
  loop every 10 s
    W->>A: heartbeat {attempt 2, progress}
    A-->>W: lease extended, cancel=false
  end
  W->>O: checkpoint every 30 s
  W->>O: upload video
  W->>A: complete {attempt 2, result_ref}
  A->>D: UPDATE ... WHERE attempt = 2
  U->>A: GET /videos/:id
  A-->>U: succeeded, signed result_url
Key idea. Async submit, pull with leases, fence every write, checkpoint to bound loss, and admit with an honest ETA.

  1. Deep dives

6.1 Leases and fencing, precisely

Before reading on. Worker A holds job 42 with attempt 7. A's network drops for 45 seconds. Walk through what happens, including when A comes back.
  1. A's heartbeats stop reaching the scheduler. At 30 seconds, the lease expires.
  2. The sweeper sets job 42 back to queued. Attempt stays 7 for now.
  3. Worker B leases job 42. The conditional update sets attempt to 8 and worker to B. B resumes from the latest checkpoint.
  4. A's network returns. A sends a heartbeat with attempt 7. The scheduler sees attempt 8 on the job and answers 409, lease lost. A stops work and discards its local state.
  5. If A had instead tried to complete with attempt 7, the update WHERE attempt = 7 matches no row. The write is rejected. B's result is the only one.

Also fence object storage writes: A writes checkpoints and results under a path that includes the attempt number, such as jobs/42/attempt-7/. B's checkpoints go to attempt-8. The job row points at the path from the current attempt, so A's late uploads are orphaned files that a cleanup job removes.

What separates answers: ownership of a job

WeakAssumes a timed-out worker is dead

Reassigns the job, and lets the returning worker overwrite the result.

GoodLeases with heartbeats

Uses expiring leases and a sweeper to reassign jobs from silent workers.

StrongLeases plus fencing everywhere

Increments an attempt number on each lease, rejects stale writes in the database and namespaces storage writes by attempt, so only the current owner can finish the job.

6.2 Checkpoints and preemption

Before reading on. Checkpoints cost storage bandwidth. How often should a job checkpoint?

Trade the cost of writing a checkpoint against the expected work lost. If a checkpoint takes 2 seconds to write and preemption is rare, every 60 seconds may be enough. If preemptible GPUs are common, every 30 seconds keeps average loss near 15 seconds. Checkpoint also on the preemption warning, which cloud providers send shortly before reclaiming the machine: the loss then drops to almost nothing for warned preemptions.

A checkpoint must be consistent. Write it to a temporary path, then record it on the job only after the upload succeeds. A worker that dies mid-upload leaves the previous checkpoint as the valid one.

A checkpoint is only useful to a worker with the same model version and compatible hardware. Store the model version and step count with it.

6.3 Admission control and ETAs

Before reading on. A third of the preemptible GPUs are reclaimed in ten minutes. What changes for users?

The completion rate drops, so the ETA estimate for new jobs rises: ETA ≈ (jobs ahead in queue ÷ recent completions per second) + expected run time. Show it on submit and update it in status.

When the estimate passes a threshold, such as 30 minutes, the API rejects new free-tier jobs with a clear message, and keeps accepting paid jobs. Per-user caps on queued jobs, such as 5, stop one user from filling the queue. Say explicitly that queue wait may exceed the normal target during shortages; the system's promise is honesty, not infinite capacity.

6.4 Scheduling policy

Order the queue by priority, then age. Pure priority starves free-tier jobs during long busy periods, so add aging: a job's effective priority rises the longer it waits. Match jobs to workers by model version and minimum GPU memory; long high-resolution videos may need larger GPUs.

Keep workers loaded with one model version. Loading weights takes minutes, so switching versions per job would waste most of the GPU time. The autoscaler shifts workers between versions based on queue depth per version, and drains a version's queue before retiring it.

6.5 Cancellation and progress

Cancel sets a flag on the job. A queued job moves straight to cancelled. A running job's worker sees cancel: true in its next heartbeat response, stops within 10 seconds, releases the GPU, and deletes its checkpoints. The job becomes cancelled with the attempt check applied.

Progress rides on heartbeats: steps done over total steps. Write it to a cache, not the database, and let clients poll every few seconds or subscribe with SSE.

6.6 Delivery

The worker uploads the finished video to object storage under its attempt path, then calls complete. The API returns a short-lived signed URL, served through a CDN. Videos expire after a retention period unless the user saves them.

6.7 Failures

FailureEffectRecovery
Worker crashHeartbeats stopLease expires, job resumes from checkpoint on another worker
Preemption with warningWorker checkpoints and exitsJob resumes with almost no lost work
Network partitionWorker keeps running blindLease expires; fencing rejects its later writes
Poison jobCrashes every workerAttempt cap, then failed with the error
Scheduler crashNo new leases brieflyState in the database; another instance continues
Upload failureVideo not storedRetry from local disk before the lease expires
Database failoverSeconds of errorsWorkers retry heartbeats; leases are sized to survive it

6.8 The scheduler's pick query

Before reading on. 2,000 idle workers ask for jobs at the same time. How do you stop them all from fighting over the same row?

A naive pick, "select the oldest queued job, then update it," makes every worker select the same row, and all but one update fails. At scale, that contention wastes most of the scheduler's time.

Use a locking read that skips rows other transactions hold:

WITH next AS (
  SELECT id FROM jobs
  WHERE status = 'queued' AND model_version = :v
  ORDER BY effective_priority DESC, created_at
  LIMIT 1
  FOR UPDATE SKIP LOCKED
)
UPDATE jobs SET status = 'leased', attempt = attempt + 1,
       worker_id = :w, lease_expires_at = now() + interval '30 seconds'
FROM next WHERE jobs.id = next.id
RETURNING jobs.*;

Each worker takes a different row without waiting. The index on (status, model_version, effective_priority, created_at) makes the pick a short index scan.

At higher rates, move the ready queue into a dedicated queue per model version, such as a Redis sorted set by priority, and keep the database as the record of truth for state transitions. The lease update in the database stays the authoritative step; the queue is an accelerator.

6.9 Cost per video

For practice: 4 GPU-minutes per video. If a GPU-hour costs C, one video costs about C / 15 in GPU time. Preemptible GPUs might cost 40% of reserved ones. With checkpoints every 30 seconds and a 5% preemption rate, preemption adds about 15 seconds × 5% of lost work, well under 1%. So moving most of the fleet to preemptible capacity cuts GPU cost by more than half, and checkpointing is what makes that safe.

Keep a reserved core of GPUs for paid traffic and use preemptible capacity for the rest. When preemptible capacity is reclaimed, admission control raises ETAs for free users first.

What separates answers: scale and cost

WeakSingle-row contention

Every worker competes for the same queued row; costs are not discussed.

GoodSkip-locked picks

Uses skip-locked reads or a separate queue to hand out distinct jobs.

StrongCost-aware capacity

Quantifies cost per video, shows how checkpoints make preemptible GPUs safe, and splits capacity into a reserved core and a cheaper elastic tier.

  1. Variants

7.1 Multi-GPU jobs

Long or high-resolution videos may need several GPUs at once. Lease a gang of workers together, or schedule onto multi-GPU machines as one worker. Partial gangs waste GPUs, so reserve capacity for gang jobs explicitly.

7.2 Several regions

Run a queue per region and route submissions to the region with the shortest ETA, subject to data residency. Checkpoints stay in-region, so resuming a job means staying in its region.

7.3 Previews

Offer a fast low-resolution preview first. Run it on smaller GPUs, return it quickly, and let the user decide whether to spend the full job.

7.4 At ten times the traffic

At 100 submissions per second, peak concurrency reaches about 24,000 GPUs, which is more than one region can hold. Split the fleet across regions with a queue per region, route jobs by ETA, and keep checkpoints in-region. The scheduler's pick moves to per-region, per-model queues. Storage for outputs reaches about 100 TB a day, so default retention shortens and older videos move to cheaper storage.

  1. The transferable pattern

This is a leased job queue with fenced, checkpointed workers. The same design runs CI pipelines, video transcoding, batch ML inference, and data processing on spot instances. The recipe carries over: async submit with idempotency, pull scheduling with leases, fencing tokens on every write, checkpoints sized to the failure rate, and admission that tells users the truth about waits.

Review: the 30-second answer

  • Async API with idempotency keys. 202 and a job ID; retries never duplicate.
  • Pull with leases. Workers heartbeat; a sweeper requeues silent jobs.
  • Fence every write with the attempt number. Database updates and storage paths alike.
  • Checkpoint every 30 seconds and on preemption warnings. Lost work stays small.
  • Admit with an honest ETA. Shed free traffic when the wait grows; Little's law sizes the pool.

Quiz

+Why do workers pull jobs instead of the scheduler pushing them?

The scheduler would need an accurate, live list of available workers, which changes constantly with preemption and crashes. Pulling lets each idle worker ask for work that matches its loaded model.

+What does the attempt number protect against?

A worker whose lease expired, but which is still running, writing a result after another worker has taken the job. Its writes carry the old attempt number and are rejected.

+How often should a job checkpoint?

Often enough that expected lost work is small relative to checkpoint cost. With frequent preemptions, about every 30 seconds keeps average loss near 15 seconds. Also checkpoint on preemption warnings.

+How many GPUs are busy at peak, and how did you get the number?

Little's law: 10 jobs per second times 240 seconds per job equals 2,400 jobs in progress, so about 2,400 busy GPUs, plus a warm buffer.

+Why keep each worker on one model version?

Loading model weights takes minutes. Switching versions per job would spend most of each GPU's time loading instead of generating.

+What does FOR UPDATE SKIP LOCKED do for the scheduler?

Each worker's pick skips rows that another transaction already locked, so concurrent workers take different jobs without waiting on each other.

+Why do checkpoints make preemptible GPUs safe to use?

They cap the work lost when a GPU is reclaimed to seconds, so the discount on preemptible capacity far outweighs the small amount of repeated work.

Sources and further reading

NextMining an Unlabeled Corpus