Queues with quotas, resource-aware admission across clusters, gang scheduling without starvation, and backfill on predicted runtimes. All while keeping scheduling event-driven, lock-free, and sub-millisecond.
This is part 4 of Inside Union, a five-part series on the engine behind Flyte 2.0. It follows part 1 on the durable AI runtime, part 2 on the replay log, and part 3 on leases and the scheduling engine, and continues with part 5 on observability.
“Make it work, make it right, make it fast.”
— Kent Beck
A slot counter is not enough
Queues with quotas, resource-aware admission across clusters, gang scheduling without starvation, and backfill on predicted runtimes. All while keeping scheduling event-driven, lock-free, and sub-millisecond. It was fast and fair. It also had a blind spot that anyone who has run GPUs on Kubernetes will recognize. A queue could allow a hundred concurrent actions, even if those actions were a hundred eight-GPU Ray jobs targeting a cluster with only sixteen GPUs. The pods pended. The leases sat in Sent. Everything behind them waited, and from the control plane, a full cluster looked exactly like a slow one. At the scale customers actually run, hundreds of thousands of pods across multiple clusters, that blind spot is expensive. The decisions themselves slow down, everyone's work waits behind them, and the error that comes back is a pending pod with no explanation.
The deeper problem is that a concurrency limit is only meaningful if you know what each action uses. A hundred single-CPU tasks and a hundred eight-GPU jobs may count the same, but they demand very different resources.
This post covers the requirements that matter when compute is shared and finite:
- good decisions on finite compute
- gang scheduling for distributed jobs, without starvation
- quotas and concurrency per team and workload
- priority between workloads that can change while they run
- stickiness for warm pools and data
- a real-time answer to "what is it waiting for"
The first post explained why we were in a position to do better. Every action arrives with its full infrastructure declaration, so the control plane knows what an action needs before any pod exists. The resource-aware scheduler turns that into admission control across clusters. It doesn't bin-pack nodes; the in-cluster scheduler still handles placement. Instead, it answers a simpler question: is there room for this on a cluster this queue can use right now?
The declaration travels with the action
In Flyte 2 you declare infrastructure next to the code that needs it:
Nothing here is pre-registered. The TaskEnvironment page describes the declaration, and this post is about what the scheduler does with it. The full specification travels with every action when it's enqueued: image, requests, limits, and accelerator, pod template, retry and timeout policy, queue. The control plane knows what an action needs before any pod exists, so it can compute a demand vector and schedule on it. Because the spec is data, a retry can change the infrastructure as well as re-run the code. And because names are stable, the platform accumulates a history per task that feeds back into scheduling rather than sitting in a dashboard.
Infrastructure here means the whole environment of a task. Each task can run in its own container, with its own image, resources, and accelerator, or share a container with its neighbours through a reusable warm pool. Warm pools are ephemeral. A pool is spun up just in time, when the first action that needs it arrives; the next actions reuse its containers instead of starting their own; and when the pool goes idle, it's freed and the capacity is reclaimed. You get the speed of a long-running service with the flexibility of infrastructure that only exists while it's needed.
Queues
A queue is a named, scoped backlog with a policy attached, and it's the unit of policy in the scheduler. Queues are tenant-wide today, with per-project and per-domain scoping coming. A run picks a queue at creation, a task can override it (the queues page shows the three levels: environment, task, and run), and the resolved name is written to the lease so it survives restarts.
Queues are dynamic. Updating one never kills in-flight work; the new limits gate the next admission or the next tick (one scheduling cycle).
The tick
Each shard has one scheduler that owns its scheduling decisions. It's a single-threaded, event-driven engine. It runs the instant an action is enqueued, a worker connects, or capacity is released, with a timer as the backstop for anything that slips between events. With unlimited capacity downstream, an action goes from enqueued to scheduled in under a microsecond, and the whole path, end to end, takes 10–20 ms. We still call one pass a tick, and one reads like this:
Dispatch itself, meaning persisting the assignment, building the lease, and pushing it down the worker's stream, happens on a separate pool fed by a channel. If the channel is full, the tick rolls its reservation back and the lease is picked up next time. Nothing in the tick waits on I/O, and nothing in it takes a lock.
Maintaining strict concurrency limits
Every cap above needs an "in flight" number, and we set goals for it before choosing a design: never over-count, tolerate a brief under-count, never drift or go negative, and be fast. Paired increment-and-decrement counters fail those goals in practice. A missed decrement leaks capacity forever, and a double decrement over-schedules. The system is designed to do neither while staying fast. It uses monotonic vector clocks and lock-free data structures, so the in-flight counts never leak and never over-decrement, and the scheduler reads them without taking a lock.
Resource caps on queues are maintained the same way. The same set of system primitives keeps them fast, with low overhead.
Demand and capacity
We compute a demand vector for what each action needs. Actions in Flyte can be heterogeneous: some are Ray tasks, others Spark, Dask, or distributed training. For each one, the demand is computed from its workload shape, and then used to maintain the caps on each queue and to never overshoot a cluster.

Unknown means open. When we don't know the shape of a cluster, we default to unlimited capacity, rather than holding work hostage to a gap in what we can see.
Maintaining resource caps and the scheduling algorithm
For each action the gate asks three questions, in order, each a handful of integer compares. Does the demand fit the queue's remaining `max_resources` budget? Is there a routable cluster whose remaining capacity, meaning capacity minus in-flight, covers the total? Does every pod of the action fit some node shape on that cluster? Among clusters that pass, the gate prefers the one with the most free GPU, then CPU, and reserves a worker slot.
A gang needs all of its pods up at once, so the scheduler has to decide what happens when the head of the queue doesn't fit. That decision is the scheduling policy.

`STRICT_FIFO` is the default, and it's head-of-line blocking by design. A gang needs all of its pods up together, so it's checked against a cluster's total remaining capacity. If it doesn't fit, it holds the clusters it could use and nothing behind it in the queue is placed there this tick, so freed capacity accumulates until the head fits. Without that hold, a steady stream of two-GPU actions keeps a six-GPU Ray job waiting indefinitely, which is precisely what people told us their schedulers did. `GREEDY_CAPACITY` is the other choice, for queues of independent single-pod work where throughput matters more than order, and the wait histogram makes starvation visible if it happens. Across queues, strict priority still applies: a higher-priority queue may take capacity a blocked gang is waiting for.
The fast lane
Warm environments are handled differently, because once the first action pays the cost of spinning up the infrastructure, a reusable container pool or a Ray cluster, we want to reuse that environment as much as possible. The scheduler derives a stable identity for each environment, places follow-on actions directly on the cluster where the environment already runs, and counts the environment's footprint once, however many actions run inside it.
Within a tick, warm-pool tasks and regular tasks are scheduled in separate lanes, interleaved one to one, so a burst of follow-on calls can't crowd out the rest of the queue.
Backfill, coming next
`STRICT_FIFO` is correct for gangs, and it idles capacity. While a gang waits for six GPUs, a CPU-only job behind it sits even though the cluster could run it now. Backfill, coming soon, lets that work run in the gap when it provably can't delay the gang. The algorithm is EASY backfill, borrowed with gratitude from Slurm, where it has run for two decades. It needs one input the platform already has: runtime history per task, which exists because every action has a stable name, rolled up to percentiles.
The rule is simple to state. A later action may run in the gap only if, based on its runtime history, it will be gone before the blocked gang could possibly start.

Borrowing from prior art
We built this by learning from many previous systems. The table below lists what inspired us, and where our approach differs.
Where it stands
- 0.60 ms per tick, 1,000 pending actions, 100 queues, 3 clusters
- +0.17 ms added by the resource gate in enforce mode, per 1,000 actions
- 76% fewer allocations per tick after caching action keys on the lease
- 3.7 ns per queue admission check, zero allocations
When head-of-line blocking or any other wait happens, the scheduler records the reason and passes it to the queue, and the console shows it. Every wait has a reason.


Put together, these are the controls a team actually touches:
- Queues with priority between them, so interactive and production work doesn't wait behind batch.
- Concurrency control per kind of workload, so the jobs that write to a database are capped at ten concurrent writers while the queue next to them allows ten thousand.
- Throttling on resources as well as counts, with a quota per queue on the CPU, memory, and GPUs in flight.
- Priority between workloads that changes live: a queue's priority, depth, concurrency, and quota all change without killing anything in flight.
- Isolation between queues, so a queue a million deep costs nothing to the latency of the small interactive run next to it.
All of it is set from the UI, with no code changes and no redeploys.
That's the promise from the start of the series: work in one queue never slows down another. Admission, scheduling, and dispatch each have their own budget, and none of them waits on the others.
Next up
Everything above decides where and when an action runs. The next post is about watching it happen: a million-action run as a live tree, "what is it waiting for" answered from the engine's own state, and the analytical layer underneath, with runtime history, error classes, and cost per action. If you'd rather see the whole path from the user's side first, the life of a run page walks it end to end. And if you'd rather run it than read about it: `pip install flyte` and `flyte start devbox`.
The Inside Union series
- The Durable AI Runtime for Flyte 2.0
- The Replay Log That Makes Flyte 2.0 Durable
- Leases, the Scheduling Engine Behind Flyte 2.0
- Resource-Aware Scheduling for Flyte 2.0 (this post)
- Observability at a Million Actions







