Service Autoscaling
Service autoscaling is an early-access feature. The syntax described here may evolve in future releases.
A workflow service normally runs as a
single, long-lived container. With an autoscaler block you can instead run a service as a pool
of replicas whose size grows and shrinks automatically based on metrics the service itself
exposes. This is useful for workloads whose demand varies over time — inference servers, queue
workers, or any service that benefits from running more (or fewer) copies as load changes.
The autoscaler decides when to add or remove a replica by evaluating PromQL expressions against the service’s own metrics. Replicas of an autoscaled service are reached through a single DNS name that always resolves to the currently-ready replicas, so clients do not need to know how many replicas exist at any moment.
When a service declares an autoscaler, Fuzzball:
- Pre-allocates capacity for up to
replicas.maxcopies of the service and immediately launchesreplicas.minof them (which may be0) when the service is first deployed — the minimum is baseline capacity honored at startup, not just a lower bound during scaling. - Scrapes the service’s metrics endpoint (when
metrics.enabledis set) at the configured interval. - Evaluates each
scale-uptrigger and thescale-downpolicy against those metrics. When ascale-uptrigger fires, the next pending replica is started; whenscale-downfires, a running replica is retired in two phases: it is first drained — removed from the service’s DNS record (and from any dependent dynamic configuration) while it keeps running — and only stopped after thescale-down.drain-periodelapses (default 30 seconds), so requests already in flight to it can complete. Adrainingevent marks the start of the retirement and the existingscale_downevent marks its completion. - Maintains a DNS record for the service that contains one address per ready replica, adding a replica’s address once it passes its readiness probe and removing it when the replica starts draining.
Example startup behavior: a service deployed with replicas: { min: 3, max: 10 } launches 3
replicas immediately when the service first starts, providing baseline capacity before any
load-based scaling occurs.
Each scaling action is rate-limited by its cooldown-period, and the replica count is always
bounded by replicas.min and replicas.max. Replicas that are draining no longer count toward
the live total, so a scale-down is refused whenever it would take the count of serving replicas
below replicas.min.
If a replica above replicas.min fails to start or dies, it is treated as optional capacity:
the replica is returned to the pending pool (so a later scale-up can retry it) and the workflow
and its other replicas keep running. Failures of the guaranteed baseline replicas (those at or
below replicas.min) still fail the workflow.
The provisioner pool a replica pool lands on must be large enough to actually host the replicas.
At submission time Fuzzball checks the selected provisioner definition’s capacity – how many
replicas of the requested size fit per node (flooring independently on cores, memory, and devices)
multiplied by the pool’s node count (the usable nodes for a static pool, or the definition’s
maxNodes for a dynamic pool that can still provision):
replicas.minmust fit. If the pool cannot host the guaranteed baseline, the workflow submission is rejected with an error naming the service, the definition, and the shortfall. Reducereplicas.minor select a larger pool.replicas.maxabove capacity is capped. If the baseline fits but the pool cannot host the fullreplicas.max, the submission is accepted, a warning is logged, and aservice_capacity_cappedevent records the effective maximum. The service runs at the capacity the pool allows.
At runtime, if a scale-up fires but the pool has no room for another replica, the scale-up is
skipped and a scale_up_capacity_exhausted event is emitted rather than leaving the replica
pending forever. fuzzball workflow why reports such a replica as blocked on pool capacity.
Similarly, if a scale-up fires while a stage the replica depends on is still completing – most
commonly the service’s image still being pulled and converted when the workflow starts – the
release is deferred and a scale_up_deferred event is emitted. The trigger’s cooldown still
applies, and a later firing starts the replica once its dependencies have completed.
This check is evaluated per service against the pool’s nominal capacity – it is not admission control across tenants. If several autoscaled services (or a service and other jobs) share the same provisioner definition, each is validated against the full pool independently, so they can pass individually while jointly exceeding the pool. Size shared pools for their combined peak demand.
A service may define both replicas and multinode. Each replica is then a gang-scheduled group of
multinode.nodes ranks that starts all together or not at all, and the pool scales in whole groups:
replicas.min and replicas.max count groups, a scale-up adds a whole group, and a scale-down
drains one. This is how a distributed server too large for a single node – an MPI application, or a
model server sharded across nodes – is served elastically.
Rank 0 is the replica’s address. It serves the endpoint, it is what the pool’s DNS records resolve
to, it is the rank whose readiness promotes the group into the pool, and it is where
autoscaler.metrics are scraped. Clients always reach rank 0, which coordinates its peers
internally. Within a group, ranks resolve each other by rank hostname (<service>-1, <service>-2,
…), scoped to their own replica.
Two consequences worth planning for:
- A group needs as many distinct nodes as it has ranks. Two ranks of the same replica are never
placed on one node, so a pool with fewer usable nodes than
multinode.nodescannot host even a single group, and a submission whosereplicas.mincannot be met is rejected with an error naming the node cost per replica. Separate groups may still share nodes where rank slots remain, so pool capacity is the pool’s total rank slots divided by the group size. When each rank consumes a whole node – the usual case for GPU-bound serving – that reduces tofloor(N / multinode.nodes)groups. - Any rank’s failure fails its whole group. The group is stopped, deregistered, and returned
to the pending pool for re-placement as a unit. Above
replicas.minthis emits areplica_failedevent and a later scale-up re-places the group; at or belowreplicas.minit is fatal to the workflow, as for a single-container replica. Every rank reports its own exit, so a rank that is killed outright, runs out of memory, or crashes takes the group down within a second or so, and the reported error names the rank, its node, and the exit status that ended it. A rank asked to stop politely (SIGTERM) exits cleanly and is not treated as a failure – that is how the platform stops ranks during a drain, stop, or cancel.
This detects a rank that stops, not one that stalls. A rank whose application hangs while its container keeps running still reports as healthy. A rank whose node stops answering the control plane is detected, but by its lease heartbeat going unanswered rather than by its own report, so that takes minutes rather than a second. If a partially-live group would produce wrong results rather than an outage, give the service its own application-level health signal on rank 0 – a readiness probe that exercises the whole group rather than rank 0 alone.
The default container metrics (container_cpu_*,container_memory_*) describe rank 0’s container only on a multi-node replica, not the group as a whole. A scaling policy that must reflect the whole group should be driven by a metric the workload itself exposes on rank 0 – which is typically where a distributed server already aggregates its group-wide counters.
The autoscaler evaluates PromQL expressions against metrics collected from your services. Fuzzball supports two types of metrics:
When you enable metrics.enabled: true on a service, Fuzzball scrapes a Prometheus-compatible
metrics endpoint from each replica at the configured interval. These are metrics that your service
application exposes — for example, vllm:num_requests_waiting from vLLM, or custom business metrics
from your application.
Configure the scrape with:
metrics.path— HTTP path to the metrics endpoint (default/metrics)metrics.port— TCP port where metrics are servedmetrics.interval— Scrape interval in seconds (default15)
All scraped metrics are tagged with fuzzball_service and fuzzball_replica labels, so your
PromQL queries can aggregate across the pool (e.g., sum(my_metric) by (fuzzball_service)) or
target specific replicas.
For every autoscaled service (whether or not it exposes its own metrics), Fuzzball automatically
collects synthetic container-level metrics from the underlying runtime. These metrics are available
immediately without requiring your service to expose a /metrics endpoint, making them ideal for
simple CPU- or memory-based scaling policies.
The default metrics follow the cAdvisor naming convention:
| Metric name | Type | Description |
|---|---|---|
container_cpu_usage_seconds_total | Counter | Cumulative CPU time in seconds. Use rate(container_cpu_usage_seconds_total[1m]) to get current CPU cores in use. |
container_memory_usage_bytes | Gauge | Current memory usage in bytes. |
container_memory_limit_bytes | Gauge | Memory limit in bytes (omitted if the container has no effective limit). Useful for ratio-based policies like container_memory_usage_bytes / container_memory_limit_bytes. |
PromQL functions that compute over consecutive samples —rate(),increase(),delta(), and similar — need a range window of at least twice the metrics collection interval to reliably contain the two samples they require. With the default 15-second interval, use a window of[1m]as in the examples on this page; a query likerate(container_cpu_usage_seconds_total[10s])can never return a result. Fuzzball rejects such queries when the workflow is submitted (see Validation rules). Loweringautoscaler.metrics.intervalpermits shorter windows.
These metrics are collected from the container runtime at the interval set by
autoscaler.metrics.interval (default 15 seconds) and are tagged with the same fuzzball_service
and fuzzball_replica labels. The interval applies even when autoscaler.metrics.enabled is not
set — configuring the interval on its own is the supported way to sample CPU and memory faster for
tighter scaling windows.
Example: Scale up when average CPU usage across all replicas exceeds 80% of one core for one minute:
autoscaler:
replicas:
min: 1
max: 10
scale-up:
triggers:
- metrics-query: rate(container_cpu_usage_seconds_total[1m]) >= bool 0.8
cooldown-period: 60
Example: Scale down when a replica’s memory usage falls below 50% of its limit:
autoscaler:
scale-down:
metrics-query: |
(container_memory_usage_bytes / container_memory_limit_bytes) <= bool 0.5
cooldown-period: 300
Container metrics are synthetic — they are derived from container runtime stats rather than scraped from an endpoint. This means they appear in your PromQL queries just like scraped metrics, but your service does not need to export them.
Replicas of a self-scaling service share a single autoscaler DNS name of the form:
<service>.autoscaler.<workflow-id>.fuzzball
Resolving this name returns one A record for every ready replica. This name is reachable from
inside the cluster only; to reach the pool from outside it, give the service a
network endpoint.
A pool’s endpoint is a single address for the whole pool, held for the life of the workflow, and
each request is forwarded to one of the replicas that is ready at that moment. Endpoints on a pool
must use type: subdomain and reference a named port.
A caller that would rather balance across the replicas itself can set per-replica: true on the
endpoint, which gives each live replica
its own address
alongside the pool’s. The two coexist: the pool URL serves callers that want Fuzzball to balance,
and the per-replica URLs serve callers that balance for themselves.
Unlike other services, which publish their declared port on the node as-is (this includes a
cross-service scaling source, since it is not itself a replica pool), each replica of a pool
serves its declared port on a randomly assigned host port (see the
workflow syntax reference) so that
same-port replicas can share a node. An IP address alone is therefore not enough to connect —
the port differs per replica. Two mechanisms carry the full host:port information:
- DNS
SRVrecords named_<port-name>._<protocol>.<service>.autoscaler.<workflow-id>.fuzzball, one per ready replica, each resolving to a replica’s hostname and its assigned port. SRV-capable clients and load balancers can consume these directly. - Dynamic configuration files, which expose per-port
ip:portarrays to a consumer service’s rendered config. This is the recommended pattern for fronting a replica pool with a proxy or load balancer that takes a static backend list.
The most common pattern is a service that scales its own replicas. Define replicas and one or
more scale-up triggers (with no services field) plus an optional scale-down policy:
version: v1
services:
inference:
image:
uri: docker://vllm/vllm-openai:latest
command: ["--model", "facebook/opt-125m", "--port", "8000"]
resource:
cpu:
cores: 8
memory:
size: 16GB
network:
ports:
- name: http
port: 8000
protocol: tcp
readiness-probe:
http-get:
path: /health
port: 8000
period-seconds: 5
autoscaler:
replicas:
min: 1
max: 8
metrics:
enabled: true
path: /metrics
port: 8000
interval: 15
scale-up:
triggers:
- metrics-query: vllm:num_requests_waiting >= bool 2
cooldown-period: 60
scale-down:
metrics-query: vllm:num_requests_running == bool 0
cooldown-period: 300
drain-period: 120
This service starts with one replica, scales up (one replica at a time, at most once per minute)
whenever two or more requests are waiting, and scales back down — no more often than every five
minutes — when a replica has no requests running. A retiring replica leaves DNS immediately but
keeps running for drain-period seconds (here two minutes, e.g. to let a long streaming
generation finish; default 30, maximum 600) before it is stopped.
Use the PromQLboolmodifier in comparisons (>= bool 2,== bool 0) so the query returns1when the condition holds and0when it does not. The autoscaler treats any non-zero result as a trigger to act.
Setting replicas.min: 0 lets a service idle with no running replicas and start its first replica
only when a trigger fires. This is well suited to bursty or on-demand workloads where you do not
want to hold capacity while the service is idle.
A service’s own scale-up triggers are PromQL queries over metrics scraped from its running
replicas, so at zero replicas there is nothing left to evaluate and the service cannot restart
itself. Two things can start that first replica:
- A request to the pool’s endpoint. A request arriving at the endpoint of a pool idling at zero
starts one replica and returns
503 Service Unavailablewith aRetry-Afterheader; retrying after the cold start reaches the new replica. See waking an idle pool. - A cross-service trigger. Another service in the same workflow scales this one, as described in cross-service scaling.
Both paths work every time the service returns to zero, not only the first time.
Instead of scaling its own replicas, a service may act as a source that scales one or more other services. This is useful for bringing a pool of workers up from zero in response to demand observed by a frontend.
To do this, give the source service’s scale-up triggers a services list naming the target
services. The targets define the replicas pool; the source does not:
version: v1
services:
frontend:
image:
uri: docker://mycompany/frontend:latest
network:
ports:
- name: metrics
port: 9090
protocol: tcp
autoscaler:
metrics:
enabled: true
path: /metrics
port: 9090
interval: 15
scale-up:
triggers:
- metrics-query: queue_depth >= bool 10
cooldown-period: 30
services: [worker]
worker:
image:
uri: docker://mycompany/worker:latest
autoscaler:
replicas:
min: 0
max: 20
scale-down:
metrics-query: worker_idle == bool 1
cooldown-period: 120
Here the frontend watches its queue depth and, when it grows, scales the worker pool up from
zero. Each worker owns its own scale-down policy, so it removes idle replicas independently of
the source.
A service in a workflow can render a configuration file from the live addresses of other services and have it injected into its container, then re-rendered automatically as those services’ replicas scale up and down. This lets a service (such as a load balancer or proxy) keep an up-to-date list of its backends without restarting.
Add a dynamic-config block naming the services to watch, the destination path inside the
container, and a script that emits the file contents. Each watched service is exposed to the
script as bash arrays (names uppercased, with any character that is not a valid identifier
replaced by _):
SERVICES_<NAME>— the ready replicas’ IP addresses.SERVICES_<NAME>_<PORT-NAME>— one array per named TCP port, holdingip:portpairs that carry each replica’s randomly assigned host port. Use these to address the replicas: the assigned host port differs per replica, so an IP plus the declared port is not reachable.
services:
proxy:
image:
uri: docker://haproxy:latest
dynamic-config:
path: /usr/local/etc/haproxy/backends.cfg
services: [worker]
script: |
for addr in "${SERVICES_WORKER_HTTP[@]}"; do
cat <<EOF
server ${addr} ${addr} check
EOF
done
Here worker declares a port named http; each rendered server line points at one replica’s
ip:port.
A retiring replica leaves the rendered arrays at the start of its drain period, while it is still running. A consumer that re-reads its configuration promptly therefore stops sending new work to the retiring backend well before it stops, and requests already dispatched to it can complete within the drain window.
Thescriptruns in a heavily restricted interpreter for safety: only bash builtins and heredocs are available. External commands (other than a no-argumentcatto support thecat <<EOFform) and any filesystem reads or writes are denied.
Fuzzball rejects a workflow at submission time if its autoscaler configuration is inconsistent. The most common rules are:
- A self-scaling service (no
scale-uptrigger targets other services) must definereplicas. - A network endpoint on a replica pool must set
type: subdomainand must reference a named port.type: pathaddresses a single backend and cannot select among a pool’s replicas. - An endpoint may only set
per-replica: trueon a service that definesreplicas, and only withtype: subdomain. replicas.maxmust be at least1and not less thanreplicas.min.- A cross-service source (any
scale-uptrigger setsservices) must not define its ownreplicasor ascale-downpolicy; each target service owns those. - Every target named in a
serviceslist must exist, must defineautoscaler.replicas, and must not define its ownscale-uptriggers (a pool is scaled up by exactly one source). - All of a service’s
scale-uptriggers must agree: either every trigger targets other services, or none do. The two modes cannot be mixed within one service. - A given service may be targeted by at most one
scale-uptrigger. scale-downmay never target other services.drain-periodis only valid onscale-down(not onscale-uptriggers) and may not exceed 600 seconds.- Every
metrics-querymust be valid PromQL. A range window read by a function that computes over consecutive samples — such as the[1m]inrate(container_cpu_usage_seconds_total[1m]), and likewiseincrease(),delta(),changes(), and friends — must be at least twice the service’s metrics collection interval (metrics.interval, default 15 seconds). A shorter window cannot reliably contain the two samples those functions need, so the query would return an empty result and the policy would never fire. Windows read by single-sample functions such asmax_over_time(), and subquery ranges like[10m:1m], are not subject to this rule. dynamic-configmay only reference services that exist in the workflow.