Union.ai
Flyte
AI

Inside Union: Resource-Aware Scheduling for Flyte 2.0

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:

Copied to clipboard!
gpu = flyte.TaskEnvironment(
    name="trainer",
    image=flyte.Image.from_debian_base().with_pip_packages("torch"),
    resources=flyte.Resources(cpu=8, memory="64Gi", gpu="H100:8"),
    reusable=flyte.ReusePolicy(replicas=4, idle_ttl=timedelta(minutes=10)),
)

@gpu.task(retries=3, timeout=flyte.Timeout(deadline=timedelta(hours=2)))
async def train(shard: int) -> Checkpoint: ...

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.

Knob What it controls
Priority Strict ordering across queues. High is always served before medium before low. Priority orders dispatch; it never preempts running work.
Depth Admission. Submissions beyond the depth are rejected immediately, so backpressure reaches the caller instead of a growing backlog.
Run and action concurrency How many root runs and how many actions may be in flight. Scoped to a queue, this is how a team caps a kind of workload: a queue for jobs that write to a database can allow ten in flight while the queue next to it allows ten thousand.
Clusters Which clusters this queue may dispatch to: all, a named list, or none (parked).
`max_resources` An absolute cap on the summed CPU, memory, GPU, and ephemeral storage of in-flight work. A quota.
Scheduling algorithm What happens when the head of the queue doesn't fit: `STRICT_FIFO`, `GREEDY_CAPACITY`, or `BACKFILL` (coming soon).
Status active, draining (reject new, finish in-flight), drained.

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:

Copied to clipboard!
fetch up to n schedulable leases
reap anything past its deadline or max_queued_time
snapshot the connected workers (id, cluster, free slots)
for each shard:
  place finalize leases first, sticky to their worker
  for each active queue, priority descending:
    filter workers to the queue's clusters
    cap by run concurrency, per-run concurrency, action concurrency
    cap by max_resources − in-flight resources           ← resource gate
    for each action (oldest first, fast lane interleaved 1:1):
      find clusters with room and a node shape that fits   ← resource gate
      reserve a worker slot; hand off to the dispatch pool
      reserve counts and resources

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.

System What we borrowed What we do differently
Kueue The demand model of pod sets, flavor-shaped capacity, `StrictFIFO` and `GreedyCapacity` as a per-queue knob. Capacity is derived from real node pools rather than declared by an admin.
Omega All-or-nothing admission for gangs; small private scheduling views. One owner per shard instead of optimistic concurrency, since tenants are disjoint.
Twine Sharding the control plane by owner rather than by cluster; entitlement-shaped absolute caps. Multi-cloud and zero-trust: workers dial out, no inbound access to any cluster.
Volcano Admitting before pods exist; leaving placement to the in-cluster scheduler. Admission spans many clusters and clouds from one control plane.
Slurm EASY backfill, bounded work per cycle, age-based anti-starvation. Runtime estimates come from the platform's own per-task history, not user declarations.

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

  1. The Durable AI Runtime for Flyte 2.0
  2. The Replay Log That Makes Flyte 2.0 Durable
  3. Leases, the Scheduling Engine Behind Flyte 2.0
  4. Resource-Aware Scheduling for Flyte 2.0 (this post)
  5. Observability at a Million Actions
Sign up

30 day free trial

Try the devbox
No items found.