Advanced Scheduling on Multi-cluster/Cloud

The queue is
the control plane.

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.

The problem

A slot counter is not enough.

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.

01

One backfill eats production

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.

02

A multi-GPU job never starts

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.

03

The error is a pending pod

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.

Concurrency that holds under load

Concurrency control and priority, as queue properties.

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.

Queue knob
  • PriorityStrict ordering across queues
  • Run & action concurrencyHow much may be in flight
  • DepthAdmission and backpressure
  • max_resourcesA quota in CPU, memory, GPU
  • ClustersWhere this queue may dispatch
  • Scheduling algorithmWhat happens when the head won’t fit
  • StatusActive, draining, drained
Gangs scheduled whole

What happens when the head of the queue doesn't fit.

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.

STRICT_FIFO
Default

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.

GREEDY_CAPACITY
Throughput

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.

BACKFILL

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.

8 × H100  ·  gpu-h200-1
Already running Gang · 6 GPU training job Small work · 1–2 GPU Idle capacity
Queue, head first
ganggrpo-train6 GPU · ~5m
taskeval-shard-a2 GPU · ~2m
taskeval-shard-b2 GPU · ~2m
taskfeaturize1 GPU · ~3m

Gang scheduling

The scheduler knows what a job needs before a pod exists.

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?

  • Demand is computed once, at enqueue, and carried on the lease. For a Ray task it's the head plus every worker group at its replica count; for Spark, the driver plus executors; for clustered training, replicas × container.
  • All-or-nothing admission. A gang is checked against a cluster's total remaining capacity, and every pod is checked against a node shape that fits it.
  • Placement stays in the cluster. Union does admission across clusters; the in-cluster Kubernetes scheduler still does bin-packing onto nodes.
  • Unknown means open. A cluster with no capacity data passes the resource gate rather than holding work hostage to a reporting outage — and the unknown cases are counted and alerted on.
demand vector  ·  grpo-train
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(),
)
↓ resolved at enqueue
CPU32
Memory256 Gi
GPU32 × H100
Pods4 × 8-GPU
gang = true all four pods must be up at once
A quota per team, in GPUs

In the units that actually run out.

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.

In-flight against quota
Counters are monotonic reservations reconciled against a once-per-second scan — nothing is ever decremented, so a missed release can't leak capacity or drift negative.

Per-queue quota

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.

Org-wide limits

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.

One flag routes across clouds

Route work to the cluster you want — across clouds and across silicon.

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.

Pick a queue to see where its work can land
queues
clusters
gpu-h200-1
aws · us-east-2 · 32 × H200
gpu-b200-eu
nebius · eu-north1 · 16 × B200
trn2-pool
aws · us-west-2 · 64 × Trn2
onprem-dc1
bare metal · dc1 · 32 × H100

Workers dial out

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.

Availability-aware placement

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.

Warm environments stay put

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.

Under the hood

One owner, one pass, no locks.

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.

0.60 msper tick — 1,000 pending actions, 100 queues, 3 clusters
+0.17 msadded by the resource gate, per 1,000 actions
3.7 nsper queue admission check, zero allocations
3.6 msp50 dispatch at 50k actions with fairness on
one scheduling pass
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
Dispatch — persisting the assignment and pushing the lease down the worker's stream — happens on a separate pool. If it's saturated the tick rolls its reservation back and picks the lease up next time.
Every wait has a reason

The console shows you which one.

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.

WaitingQueue resource quota reached — max_resources for training-queue is fully committed4m 12s
WaitingNo routable cluster with 6 free GPUs — gang held at the head of the queue1m 38s
WaitingNo node shape fits an 8-GPU pod on gpu-b200-eu22s
WaitingRun concurrency reached — 100 of 100 runs in flight9s
RunningDispatched to gpu-h200-1 · lease held, heartbeating2m 05s

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.

In practice

Set it from the UI or the CLI. The workflow doesn't change.

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
--priority Low Strictly ordered behind every High and Medium queue. Production is always served first when it has work.
--run-concurrency The backlog egresses at ten runs at a time instead of saturating the cluster.
--depth Submissions past the depth are rejected immediately, so backpressure reaches the caller instead of a growing backlog.
# 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
--cluster Repeatable. Omit it and the queue dispatches to any healthy cluster it can reach.
--edit Opens the queue's settings. The cluster subset, priority, concurrency, and depth are all editable while the queue is live; work in flight is never killed.
multi-silicon A second queue routing to a Trainium cluster runs the same Python program — the accelerator class is queue configuration, not an integration.
import flyte

# Pin a whole run to a queue at submit time.
flyte.with_runcontext(queue="gpu-h200-fast").run(
    train_pipeline, epochs=40,
)
with_runcontext The run and every action under it inherit the queue's priority, concurrency, quota, and cluster routing.
resolved once The queue name is written to the lease, so the routing decision survives a control-plane restart.
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)
.override(queue=…) Re-routes one task inside a running workflow — different priority, different concurrency, different cluster — with no fork and no resubmit.
one rule The override queue must share the run's cluster pool, because the action's inputs, code, and secrets were uploaded where that pool's clusters can read them. A cross-pool override rejects fast with an actionable error.
Enterprise grade Flyte

Open source at the core.

Union is built on Flyte, the open-source AI runtime we create and maintain under the Linux Foundation AI & Data.

A Linux Foundation AI & Data Project
4000+ companies using Flyte today
18M+ Flyte SDK downloads

Get a scheduling review.

Bring your queue, priority, concurrency, and cluster-routing layout. We'll walk through what the scheduler would do with it.