Priority keeps a backfill from starving production, concurrency drains a backlog at the rate you set, and quotas keep one team from taking the whole GPU fleet. One flag routes work across clouds, regions, and accelerator classes, and you set it once in either the UI or the CLI.
Five things the scheduler decides. Follow a card to the detail.
Most schedulers gate on integers: free slots on a worker, depth of a queue, actions per run. That works until compute is shared and finite. A queue that allows a hundred concurrent actions will happily dispatch a hundred eight-GPU jobs onto a cluster with sixteen GPUs. The pods pend. Everything behind them waits. And from the control plane, a full cluster looks exactly like a slow one.
A backlog is submitted, saturates the cluster, and the production SLA misses — because concurrency was never a property of the queue, it was something a platform team hand-rolled in DAG code.
A steady stream of two-GPU work keeps the cluster busy, so the six-GPU training job at the head of the queue waits indefinitely. Nothing is broken. Nothing is scheduled either.
Work goes out to a cluster that can't hold it, and the only signal is a pod stuck in Pending with no explanation. Someone correlates logs to find out which resource ran out.
A queue is a named backlog with a policy attached, and it is the unit of policy in the scheduler. A run picks a queue at creation, a task can override it, and the resolved name is written to the lease so it survives a restart. Every knob below changes live — updating a queue never kills work in flight.
A gang — a Ray job, a Spark app, a multi-node training run — needs all its pods up at once. Whether the scheduler holds capacity for it or lets smaller work past is the queue's scheduling algorithm, and it decides whether that job starts on time or never starts at all.
Head-of-line blocking, on purpose. A gang that doesn't fit holds the clusters it could use, so freed capacity accumulates until it fits. Without that hold, a steady trickle of small work keeps a multi-GPU job waiting forever.
For queues of independent single-pod work where order matters less than throughput. Anything that fits is placed. The wait histogram makes starvation visible if it happens.
EASY backfill, as Slurm has run it for two decades. Later work runs in the gap only if, based on its own runtime history, it will be gone before the blocked gang could possibly start.
Across queues, strict priority still applies. A higher-priority queue may take capacity a blocked gang is waiting for — priority orders dispatch, and it never preempts work that is already running.
In Flyte 2 you declare infrastructure next to the code that needs it, and the full specification travels with every action when it's enqueued. So admission control can run in the control plane, early and cheaply, against a question Kubernetes can't answer on its own: is there room for this, on some cluster this queue may use, right now?
from flyte.clustered import (
ClusteredTaskEnvironment, TorchRun,
)
env = ClusteredTaskEnvironment(
name="trainer",
image=image,
resources=flyte.Resources(
cpu=8, memory="64Gi", gpu="H100:8",
),
replicas=4,
nproc_per_node=8,
runtime=TorchRun(),
)
Counting actions doesn't stop a team from taking the whole GPU fleet — ten actions can be eighty GPUs. So max_resources caps the summed CPU, memory, GPU, and ephemeral storage of a queue's in-flight work. Give each team a queue, give each queue a budget, and the fleet divides the way you decided rather than the way the submit order happened to fall.
An absolute cap on the resources a queue may have in flight. Paired with run and action concurrency, it's how a team caps a kind of work: a queue for jobs that write to a warehouse allows ten in flight while the queue next to it allows ten thousand.
A per-tenant safety envelope wrapping every queue — caps on concurrent runs, concurrent actions, request rates, and per-run fan-out. One noisy client can't destabilize shared infrastructure, and flyte get org lists every current value against its default.
A queue's cluster selector dispatches to any healthy cluster, or to named ones. Pin a high-priority queue to your reserved H200s, send Trainium work to a Trainium cluster, and burst to a neocloud when your own racks are full — with the same Python program and the same scheduler. The selector is mutable live, with no drain, so re-pinning a queue is one CLI call rather than an integration project.
A cluster worker opens an outbound stream to the control plane and holds leases. Union holds no credentials into your clusters and never opens an inbound connection — so a cluster on AWS, GCP, Azure, or a rack in your building looks identical from the control plane.
Among the clusters a queue may use, the gate picks one whose remaining capacity covers the whole job and whose node shapes fit every pod — preferring the cluster with the most free GPU, then CPU.
Once an action pays to spin up a reusable container pool or a Ray cluster, follow-on actions are placed directly on the cluster where that environment already runs, and its footprint is counted once however many actions run inside it.
Each shard has a single scheduler that holds the whole picture in memory — every pending action, every connected worker across every cluster, every queue and quota. It runs the instant something changes rather than on a poll, and nothing in the pass waits on I/O or takes a lock. One owner also means tenants are isolated by construction: a million-action backfill from one team can't slow a five-task run from another.
fetch up to n schedulable leases
reap anything past its deadline
snapshot connected workers (cluster, free slots)
for each active queue, priority descending:
filter workers to the queue's clusters
cap by run / per-run / action concurrency
cap by max_resources − in-flight ← resource gate
for each action (oldest first, fast lane 1:1):
find a cluster with room and a node
shape that fits every pod ← resource gate
reserve a slot; hand to dispatch pool
reserve counts and resources
Every decision the gate makes is recorded as a reason on the action, so "why isn't my run moving?" is answered by the same state machine that schedules the work — not by a monitoring pipeline that might disagree with it. Each reason is also a counter and a histogram, per queue and per cluster.
A gang that has been blocking its queue shows up as a gauge — which is how a team decides to move that queue to backfill. Preemptions and OOM kills are attributed to the exact action within seconds, because the cluster worker watches Kubernetes through informers rather than waiting for someone to correlate a log line.
The teams already running on these decisions, and the number each of them reported.
Priority, depth, concurrency, quota, and cluster routing are all queue configuration. None of them is workflow code, so none of them needs a redeploy — and a single task can be re-routed to a different queue at runtime without forking the workflow.
# A backfill queue that yields to production and drains at a rate you set.
flyte create queue "backfill-queue" \
--priority Low \
--run-concurrency 10 \
--action-concurrency 200 \
--depth 10000
# Pin a high-priority queue to named GPU clusters.
flyte create queue "gpu-h200-fast" \
--priority High \
--cluster gpu-h200-1 \
--cluster gpu-b200-eu
# Re-pin it live — no drain, no redeploy.
flyte update queue "gpu-h200-fast" --edit
import flyte
# Pin a whole run to a queue at submit time.
flyte.with_runcontext(queue="gpu-h200-fast").run(
train_pipeline, epochs=40,
)
import flyte
env = flyte.TaskEnvironment(name="pipeline")
@env.task
async def main() -> None:
# The fast path stays on the run's queue …
await score(batch)
# … while the long tail goes to the backfill queue.
await heavy_task.override(queue="backfill-queue")(batch)
Union is built on Flyte, the open-source AI runtime we create and maintain under the Linux Foundation AI & Data.
Where and when work runs is settled here. Whether it survives getting there, what it produces, and whose perimeter it sits inside are the other three.
Scheduling decides where work runs. Durability decides whether it survives getting there — leases, replay, and recovery from any failure.
Explore the runtime →The training, evaluation, and inference line that these queues feed — artifacts, lineage, and event-driven triggers between stations.
See the factory →Every cluster in the routing map is yours. Workers dial out, the control plane holds no credentials, and your data never leaves your perimeter.
See the architecture →Bring your queue, priority, concurrency, and cluster-routing layout. We'll walk through what the scheduler would do with it.