BigFleet: A Fleet-Level Infrastructure Autoscaler
A design for a fleet-level infrastructure autoscaler that manages 100 million nodes across 20,000 Kubernetes clusters, accepting pluggable capacity providers and treating all capacity — bare metal, cloud, reserved, spot — as a single fungible pool priced by cost.
BigFleet is an implementation of the Infrastructure Autoscaler described in Fleet-Scale Kubernetes: An Operating Model for Homogeneous Clusters with Decoupled Capacity. It receives ClusterCapacityNeeds protobuf messages from per-cluster operators, diffs them against its own provisioned inventory, and provisions or reclaims nodes through pluggable CapacityProvider backends.
BigFleet is not a scheduler. It does not decide where pods run. It does not simulate kube-scheduler. It makes nodes appear in clusters that need them and removes nodes from clusters that don't. The scheduler handles everything after the node is Ready.
BigFleet is not a fleet manager. It does not deploy workloads, manage cluster lifecycle, or handle upgrades. It manages capacity — the pool of machines that clusters draw from.
This paper assumes familiarity with Fleet-Scale Kubernetes: An Operating Model for Homogeneous Clusters with Decoupled Capacity, which argues for the separation between scheduling (stays in each cluster) and capacity management (lives in a fleet-level autoscaler). BigFleet is a concrete implementation of that autoscaler.
BigFleet targets 100 million nodes across 20,000 Kubernetes clusters, spanning multiple cloud providers and bare-metal datacenters. No published system has solved this problem end-to-end. The building blocks exist in isolation — but they've never been composed into a single coherent architecture.
To ground what 100M nodes actually means: the largest publicly documented single-cluster Kubernetes deployment is Google's experimental 130,000-node GKE cluster, which required replacing etcd with a Spanner-based key-value store sharded across 20+ instances and disabling the cluster autoscaler entirely. The largest production single-cluster limits across the major public clouds are AWS EKS at 100,000 nodes (announced September 2025), Google GKE at 65,000 nodes, and Azure AKS at 5,000 nodes. Meta's Twine, the most ambitious published fleet manager, operates a single control plane across roughly 1 million machines within Meta's homogeneous infrastructure. BigFleet aims higher — two orders of magnitude more machines across fundamentally heterogeneous infrastructure.
The economic stakes are unforgiving. Cast AI's 2025 benchmark of 2,100+ organizations reports average CPU utilization of 10% and memory utilization of 23% across Kubernetes clusters. Datadog's State of Cloud Costs 2024 attributes 83% of container costs to idle resources — 54% from overprovisioned infrastructure, 29% from oversized workload requests. At 100M nodes, even small improvements in fleet-wide coordination translate into billions of dollars. A single percent of idle capacity, using AWS m6i.xlarge at $140/month, is roughly $1.68 billion per year of waste at on-demand rates. The industry operates at 10–20% idle capacity routinely; Komodor frames this as a "$100 billion waste problem."
Meanwhile, AI/ML training has reshaped what "capacity" even means. Meta trained Llama 4 on 100,000+ H100 GPUs spanning 5 buildings with fibre runs up to 3 km between GPU pairs. xAI's Colossus cluster reached 200,000 GPUs in 214 days. Anthropic has announced access to up to 1 million TPU chips from Google. These workloads have strict topology requirements — tensor parallelism must stay within NVLink domains (sub-microsecond latency, 900 GB/s per GPU), while inter-node fabric is 10–100× slower. A 10,000-GPU training job needs rack-level, same-spine, or same-building co-location that routinely exceeds single-cluster limits. No Kubernetes abstraction represents these constraints for cross-cluster scheduling.
The rest of this paper walks through how the industry has approached pieces of this problem, then describes BigFleet's architecture.
Fleet management at massive scale has been solved in pieces by different organizations. Each system made design decisions under specific constraints; BigFleet borrows from several and deliberately diverges from others. This section walks through the systems that inform BigFleet's architecture.
Borg has run Google's production infrastructure since the early 2000s. By the time the 2015 EuroSys paper [1] was published, the authors described "a decade of operational experience." It remains the most influential cluster management system ever published, and its design decisions are foundational to everything that came after — including Kubernetes, which Google built as its open-source successor.
Architecture. Borg organizes machines into cells of roughly 10,000 machines, each an independent fault domain within a datacenter building. The BorgMaster controlling each cell is replicated five times using Paxos, with a single elected leader. Failover typically takes ~10 seconds but can reach a minute for large cells as in-memory state is reconstructed from the Paxos-replicated change log. The scheduler runs as a separate process operating on a cached copy of cell state — an early form of optimistic concurrency. Several cells see arrival rates exceeding 10,000 tasks per minute.
Static stability. This is Borg's most important principle: running tasks continue operating even if all BorgMaster replicas fail. Per-machine Borglet agents keep tasks alive independently. New submissions and rescheduling require the master, but the data plane is decoupled from the control plane. BorgMaster achieves 99.99% availability.
Priority and preemption. Four priority bands — monitoring, production, batch, best-effort — with production-band tasks immune to mutual preemption (preventing cascading evictions). Quota is expressed as resource vectors at a given priority and enforced at admission time, not at scheduling time. Lower-priority quota is deliberately oversold. Every user has infinite quota at priority zero.
Mixed workloads and efficiency. 98% of machines in shared cells run both production and non-production tasks. Segregating them would require 20–30% more machines in the median cell. Roughly 20% of workload runs in "reclaimed resources" — capacity reserved by production tasks but not currently consumed.
Fauxmaster. A simulation tool containing a full copy of BorgMaster code with stubbed Borglet interfaces. Reads checkpoint files to replay real interactions, debug failures, perform capacity planning, and validate configuration changes before applying them.
What BigFleet takes. Cells as bounded fault domains (BigFleet's shards). Static stability as a foundational property — the data plane survives control plane failures. Priority-based preemption with the explicit no-mutual-preemption-at-equal-priority rule. Mixed-workload scheduling to capture the 20–30% efficiency gain. Simulation as an operational requirement, not a nice-to-have.
What BigFleet doesn't take. Borg is per-cell. It doesn't coordinate across cells. At Google's scale this works because cells are huge relative to most workloads. At BigFleet's target scale, cells would have to be impossibly large or the system would be per-cell like Borg — and per-cell autoscaling is what we're trying to escape.
The 2013 Omega paper [2] described an experimental platform for evaluating shared-state scheduling as an alternative to Borg's monolithic scheduler. Omega was never a production replacement for Borg — its value was as a research vehicle, and many of its ideas were subsequently folded back into Borg itself.
Shared state, optimistic concurrency. Each scheduler gets a private, frequently-updated copy of the full cell state. It makes placement decisions independently. Decisions are committed atomically to a master copy; conflicts are detected at commit time. No two-phase locking, no offer-based resource handoff.
Measured conflict rates. At moderate load with 3 parallel batch schedulers, conflict fractions stayed low (~0.1) while moving the scheduler saturation point from 4 seconds to 15 seconds of per-task decision time — a 3× scalability improvement. Fine-grained incremental transactions (rejecting only placements that would cause overcommit) had much lower conflict rates than coarse-grained conflict detection or all-or-nothing transactions. At high load, when service job scheduling times hit 10 seconds, conflict fractions exceeded 1.0 and SLOs were violated. This quantifies when optimistic concurrency breaks down.
What BigFleet takes. The declarative materialised-view model. Each shard maintains the full current state of all its clusters' needs and the full current state of its inventory. Worker cycles operate on snapshots. There's no locking on the hot path — the shard owns its state. BigFleet doesn't need parallel schedulers within a shard (the hot path is in-memory and fast), but the declarative approach is directly inherited.
Twine has managed Meta's fleet since well before the 2020 OSDI paper [3] that describes it. The paper reports on operational data from January 2019 through October 2020, documenting a consolidation effort already in progress. Twine is the most ambitious published fleet manager, operating a single control plane across 1 million+ machines in a region — 100× Borg's per-cell scope.
Three-level architecture. The Resource Broker (one per datacenter) tracks machine-to-entitlement assignments and hardware availability. A regional allocator fetches state from all Resource Brokers to make allocation decisions. The Twine Scheduler manages container lifecycles and is sharded — each shard handles ~170,000 machines. TaskControl lets applications collaborate on their own lifecycle (e.g., ZooKeeper's TaskController instructs Twine to restart followers first, leader last, during rolling upgrades).
Entitlements. The capacity abstraction. Grant a business unit a quota expressed as a count of machines of certain types, bound to a host profile (hardware and OS customisations). Jobs bind to entitlements, not physical clusters — enabling transparent migration across datacenters. Largest entitlement: ~60,000 machines. Largest job: ~15,000 tasks. Before entitlements, Meta's fleet was fragmented into hundreds of private pools with significant stranded capacity. The unified twshared pool grew from 15% of the fleet in January 2019 to 56% by October 2020.
ReBalancer. An asynchronous, continuously-running constraint solver that performs global optimizations — better balancing CPU, power, and network utilization — without blocking the critical scheduling path.
Delos. Twine's state storage, a replicated log-based storage system that replaced ZooKeeper for Meta's control plane services.
What BigFleet takes. Separating capacity allocation from workload scheduling. Twine's Resource Broker is the closest published analog to BigFleet's shard — it owns inventory and makes allocation decisions, but workload lifecycle is elsewhere. Continuous rebalancing via constraint solver, running off the critical path. The insight that capacity should be an abstract entitlement rather than a physical cluster.
What BigFleet doesn't take. Twine operates within Meta's homogeneous infrastructure. It doesn't manage multi-provider capacity. BigFleet's CapacityProvider interface generalises what Twine's Resource Broker does natively.
The 2021 SOSP paper [4] describes a system that was already running hundreds of applications at Meta. Shard Manager manages nearly 100 million shard replicas across 1 million+ servers and sits atop Twine in Meta's infrastructure stack.
Mini-SM pattern. The control plane itself is divided into 139+ mini Shard Managers, each handling a partition of servers and shards. Each application is divided into partitions of typically thousands of servers and hundreds of thousands of shard replicas. Application Managers sit above the mini-SMs, coordinating cross-partition decisions through a Partition Registry.
Constraint solver over hand-crafted heuristics. Migrating from imperative placement code to a constraint solver reduced allocator code to ~20% of the original. Hard constraints cap concurrent shard moves and enforce capacity bounds. Soft goals — expressible in a DSL — cover regional placement preferences, replica spread across fault domains, and multi-resource load balancing. Local search with group-based sampling.
What BigFleet takes. The critical insight that at 100M-entity scale, the control plane itself must be sharded. BigFleet's shards-with-a-global-coordinator structure is directly modelled on this. The preference for declarative constraint specification over imperative placement code — BigFleet's Phase 1/Phase 2/Phase 3 decision engine is structured so that policy changes don't require rewriting placement logic.
Alibaba's scheduling infrastructure [5] has evolved through multiple generations over more than a decade, now serving all Alibaba business units across dozens of datacenters with millions of containers and tens of millions of CPU cores.
Hierarchical scheduling. Level-0 provides global resource coordination, managing co-location clusters and handling exception detection. Level-1 has workload-specific schedulers: Sigma for long-running latency-sensitive services, Fuxi for batch (hundreds of thousands of parallel tasks). Level-2 implements business-specific logic for individual services.
Co-location results. CPU utilization went from ~10% (pure online workloads) to 50%+ with co-location, reaching 65% during Double 11 2021. Saved hundreds of thousands of CPUs. Task scheduling throughput: 20,000 tasks/second. Their co-location technology was open-sourced as Koordinator in April 2022.
Priority and isolation. LSE (exclusive cores), LSR (reserved cores), LS (standard latency-sensitive), BE (best effort, preemptable). Kernel-level isolation through CPU scheduling priority (BVT), memory collection and OOM priority, network priority, and I/O isolation.
What BigFleet takes. The validation that hierarchical scheduling with a global coordinator scales to millions of containers. The economic case for aggressive co-location — a 4–6× utilization improvement justifies substantial investment in the preemption and isolation machinery.
The 2023 SoCC paper [6] describes a system ByteDance had already deployed in production. Gödel achieves 5,000+ pods per second — 500× the default Kubernetes scheduler's ~10 pods/second in a 5,000-node cluster — and was open-sourced in April 2024.
Dispatcher → Scheduler → Binder pipeline. The Dispatcher (single instance) handles queuing and distributes applications across scheduler instances. Multiple Scheduler shards perform expensive filtering and scoring in parallel, each on its assigned node partition. The Binder (single instance) performs conflict detection and either commits placements or rejects conflicts for rescheduling. Expensive computation parallelised; conflict resolution serial.
Production scale. Clusters with 20,000+ nodes and over 1 million pods. Combined with Katalyst (their NUMA-aware node-level resource manager), average CPU utilization increased from 30% to 60%. Open-sourced April 2024.
What BigFleet takes. The pattern of parallelising expensive work while keeping conflict resolution serial. BigFleet applies this between shards (parallel independent decisions) and the global coordinator (serial rebalancing and cross-shard preemption). The demonstration that this pattern is fast enough in practice — 5,000 decisions/second through a single serial resolver.
AWS has not published a unified fleet manager, but patterns from the AWS Builders' Library and public talks inform BigFleet's reliability design.
Cell-based architecture [7]. Partition a service into independent replicas, each serving a subset of clients. A Cell Router routes requests; each cell contains everything needed to operate independently. AWS typically over-provisions by 50% with 3 AZs, so the service continues without new instances even if an AZ fails. Static stability in practice.
Shuffle sharding [8]. With 8 workers and 2-worker shards, there are C(8,2) = 28 combinations — the probability of two customers sharing an identical shard is ~3.6%. A 52-element fleet with 4-element shards yields over 300,000 combinations. Used in Route 53, AWS Shield, and many internal services.
AWS Hyperplane. The internal network function virtualisation platform powering NLB, NAT Gateway, PrivateLink, and Lambda VPC access uses a push-via-S3 architecture: the control plane writes configuration to S3, the data plane periodically downloads it. This avoids overwhelming the control plane when the data plane fleet is 100× larger.
What BigFleet takes. Cell isolation as the primary blast-radius limiter (BigFleet's shards are cells). Static stability as non-negotiable. The push-via-durable-storage pattern inspires BigFleet's durable store for machine assignments — shards read their assignments on startup rather than requiring the global coordinator to be up.
Cluster API (CAPI) [9] provides declarative APIs for cluster lifecycle management. Its abstractions directly inform BigFleet's capacity model.
Machine vs MachinePool. Machines are immutable — once created, never updated, only deleted and replaced. MachineDeployments handle rolling updates and scaling through MachineSet management. MachinePools delegate entirely to cloud-native scaling groups (AWS ASGs, Azure VMSS, GCP MIGs), treating the underlying machines as a black box. This is the fixed-versus-elastic distinction as a first-class API concern.
Three-provider pattern. Infrastructure Provider (provisions VMs and networks), Bootstrap Provider (installs Kubernetes), Control Plane Provider (manages API server components). Implementations exist for AWS, Azure, GCP, vSphere, bare metal (metal3.io), and many others.
What BigFleet takes. The Machine/MachinePool distinction — BigFleet's fixed and elastic capacity directly mirror this. The pluggable infrastructure provider pattern is the direct ancestor of BigFleet's CapacityProvider.
Karpenter was launched by AWS in November 2021 and reached v1.0 GA in August 2024 [10]. It represents the state of the art in per-cluster autoscaling.
Minimal provider interface. The CloudProvider interface defines exactly 9 methods: Create, Delete, Get, List, GetInstanceTypes, IsDrifted, RepairPolicies, Name, GetSupportedNodeClasses. Sufficient for full node lifecycle management. AWS (GA) and Azure (Preview) implementations.
Two-phase scheduling. Phase 1: try to place pending pods on existing nodes, checking taints, ports, volumes, resources, topology. Phase 2: create NodeClaim templates for new nodes, filter instance types by requirements, pick the cheapest. Pods queued by size (larger first for bin-packing). Passes 60 compatible instance type options to the Fleet API.
Offering-based cost optimization. Each instance type has offerings in multiple zones and capacity types (spot, on-demand, reserved). Karpenter prioritises Reserved → Spot → On-Demand, caching spot unavailability for 3 minutes. Consolidation: replace nodes with cheaper instance types; migrate pods from underutilised nodes; remove empty nodes.
Per-cluster limitation. Karpenter runs as a controller inside a single cluster with no cross-cluster awareness. Each of 20,000 clusters needs its own Karpenter deployment making independent decisions.
What BigFleet takes. Karpenter's minimal CloudProvider interface as a model for BigFleet's CapacityProvider. BigFleet's interface is deliberately smaller — six methods rather than Karpenter's nine — because BigFleet's simpler model eliminates several of Karpenter's concerns: no separate offering enumeration (List returns machines in every state, including speculative quota), no drift detection, no capability negotiation. Offering-based cost optimization is retained but expressed as pure cost functions rather than a hardcoded priority waterfall. The insight that simulation before provisioning prevents waste carries over.
What BigFleet doesn't take. The per-cluster limitation. The entire reason BigFleet exists is to solve the coordination problem Karpenter cannot see.
Five gaps define BigFleet's contribution:
No system unifies 100M nodes across multi-cluster, multi-provider infrastructure. Borg is per-cell (~10K machines). Twine is per-region within Meta's homogeneous infrastructure. Shard Manager reaches 100M shard replicas but manages application-level sharding, not infrastructure capacity. Google's 130K-node GKE experiment is single-cluster.
Per-cluster autoscalers cannot see fleet-wide demand. Karpenter and Cluster Autoscaler optimize locally within each cluster. At 20,000 clusters, each independently maintains buffer capacity and provisions new nodes while neighbors have headroom.
Multi-cluster fleet managers don't do infrastructure allocation. Karmada, KubeAdmiral, and predecessors distribute workloads across existing clusters. They assume clusters already exist with sufficient capacity. Infrastructure provisioning — creating clusters, sizing them, selecting instance types, managing bare-metal-to-cloud ratios — is outside their scope.
Priority-based cross-cluster preemption does not exist. Borg preempts within a cell. Twine allocates within a region. No published system reclaims capacity from a low-priority workload in cluster A to serve a high-priority workload in cluster B.
Hybrid fixed+elastic capacity with unified cost optimization is ad hoc. Bare metal, reserved instances, on-demand, and spot form a four-tier economic model. No system architecturally unifies capacity selection across these tiers and across heterogeneous providers (AWS + GCP + Azure + bare metal) into coherent fleet-wide decisions.
Infrastructure comes in two fundamentally different forms:
Fixed capacity — bare metal, reserved instances, dedicated hosts. Machines that exist whether you use them or not. You paid for them. The marginal cost of assigning one to a cluster is zero. The decision is allocation — which cluster gets this machine?
Elastic capacity — cloud on-demand, spot instances. Machines that don't exist until you ask for them. They cost money per second. The decision is procurement — should I buy this machine, and from whom?
BigFleet treats both identically. A CapacityProvider reports what it has available and what it costs. BigFleet picks the cheapest option that satisfies the request. The "waterfall" of bare metal → reserved → spot → on-demand isn't encoded in the architecture — it emerges naturally from the prices:
| Capacity type | Typical cost/hr | Availability | Interruption risk |
|---|---|---|---|
| Bare metal (owned) | $0.00 | Fixed pool | None |
| Reserved instance | $2.10 | Guaranteed | None |
| Spot | $1.80 | Variable | 5–15% hourly |
| On-demand | $6.00 | High | None |
The autoscaler doesn't need a tier system. It needs a cost function:
effective_cost = price + (interruption_probability × interruption_penalty)
A spot instance at $1.80/hr with 10% hourly interruption probability and a $5.00 interruption penalty (restart cost, wasted work) has an effective cost of $2.30/hr. The autoscaler compares this to on-demand at $6.00/hr and picks spot. If the workload is a long-running training job where interruption penalty is $50.00, the effective cost of spot becomes $6.80/hr and on-demand wins. No special cases. Just math.
BigFleet's entire purpose is to enable homogeneous clusters. If we shard BigFleet by location — "the Amsterdam shard manages Amsterdam machines" — we recreate the specialisation problem at the shard layer. The GPU shard becomes a bottleneck. The CPU shard has idle capacity. One shard fails and a specific capability is lost.
Instead, shards are horizontal slices across all infrastructure. Each shard has machines in multiple locations, of multiple types, with cloud burst capacity in multiple regions. No shard is the "GPU shard" or the "Amsterdam shard."
But perfect homogeneity is a gradient, not a binary. If Amsterdam has 1,000 GPUs and there are 200 shards, round-robin gives each shard 5 Amsterdam GPUs — nowhere near enough for a 64-GPU same-rack training job. Homogeneity at the machine level collapses at high shard counts for scarce resources.
The constraint is: the shard count must not exceed the point where topology-constrained requests become unsatisfiable. If the largest topology request is 64 same-rack GPUs and the smallest rack has 128 GPUs, you can assign whole racks to shards — meaning the shard count is bounded by the number of GPU racks, not the number of GPUs. If you have 20 GPU racks across the fleet, you can't have more than 20 shards that can satisfy same-rack GPU requests.
In practice, this isn't as restrictive as it sounds. The shard count scales with fleet size (Section 16), and fleet size correlates with resource pool size. A fleet small enough to have only 200 GPUs also has few enough clusters to run on 2–5 shards, where each shard has plenty. A fleet large enough to need 200 shards has tens of thousands of GPUs across hundreds of racks.
A shard is a controller that manages a set of clusters and a slice of the global machine inventory. It:
ClusterCapacityNeeds from its assigned clustersCapacityProvider backends to create or destroy machinesThe shard count is determined by two constraints in tension:
Controller blast radius — each shard should manage few enough clusters that a shard failure is survivable. ~100 clusters per shard is a reasonable target.
Resource pool depth — each shard must have enough of each scarce resource type (GPUs, high-memory nodes, specific hardware) to satisfy topology-constrained requests. This sets an upper bound on shard count.
At 100M nodes, both constraints are easily satisfied: 200 shards × 100 clusters × 500K machines, with tens of thousands of GPUs per shard. At smaller fleet sizes, the resource pool constraint dominates — you may need fewer, larger shards to keep topology requests satisfiable.
Machine assignment to shards: Fixed machines are distributed across shards at the topology-domain level, not the individual-machine level. Entire racks or zone-slices are assigned to shards, so each shard has contiguous topology domains. This preserves the ability to satisfy Same operator constraints (e.g., "all 64 nodes on the same rack"). Cloud machines are assigned to the shard that provisioned them.
Cluster assignment to shards: Clusters are assigned to shards on first contact — when a cluster's operator sends its first roll-up, BigFleet assigns it to the least-loaded shard. The assignment is permanent. Clusters don't move between shards during normal operation. The only exception is a shard split (see Section 16).
Every machine BigFleet knows about is owned by exactly one shard and lives in one of three states. A machine is represented the same way whether it's a bare-metal server in Amsterdam, a running cloud instance, or a quota slot that doesn't exist yet:
type Machine struct {
ID MachineID // BigFleet's internal ID (always present)
Host *HostRef // nil = speculative, non-nil = bound to a real host
Cluster ClusterID // 0 = unbound, non-zero = configured for this cluster
Shard ShardID // which shard owns this machine
Profile Profile // instance type, zone, resources
PricePerHour float64
}
type HostRef struct {
Provider string // "aws-eu-west-1", "bare-metal-amsterdam"
Ref string // AWS instance ID, bare-metal serial, etc.
}
Two fields encode the stable states, and the three meaningful combinations are the three steady states:
| State | Host | Cluster | Meaning |
|---|---|---|---|
| Speculative | nil | 0 | A quota slot — promise of capacity, no physical resource yet |
| Idle | set | 0 | Real hardware, no cluster membership, no kubelet configured |
| Configured | set | set | Real hardware, kubelet configured for this cluster, joined and ready |
The fourth combination — Host = nil, Cluster != 0 — is impossible. The core invariant: a machine must be bound to a real host before it can be configured for a cluster.
These are the stable states a machine rests in. The transitions between them are themselves visible as transitional states, because transitions take real time — seconds for Delete, minutes for Create, and potentially hours for Drain. A machine in a transitional state is not available for allocation; it's in motion toward a stable state.
| Transitional state | Meaning |
|---|---|
| Creating | Provider is running up hardware (cloud: EC2 instance starting; bare-metal: rack being commissioned) |
| Configuring | Bootstrap blob applied; kubelet starting and joining cluster |
| Draining | Workloads being evicted gracefully; kubelet preparing to unregister |
| Deleting | Hardware being torn down (cloud instance terminating) |
| Failed | A transition couldn't complete within its timeout; needs shard intervention |
The full state set is eight values: three stable (Speculative, Idle, Configured), four transitional, and one terminal-pending-cleanup (Failed). BigFleet's decision engine filters by stable states when picking capacity for allocation — machines in transitional states don't count as available.
Machines move between stable states via explicit operations, each of which takes the machine through an intermediate transitional state:
┌───────────────┐
│ │ Delete (cloud)
▼ │
┌─────────────┐ Create ┌──────┐ │
│ Speculative │──────────▶│ Idle │───┘
└─────────────┘ └──────┘
│ ▲
Bootstrap│ │ Drain + unregister
▼ │
┌─────────────┐
│ Configured │
└─────────────┘
Speculative → Idle (Create). BigFleet calls CapacityProvider.Create. The machine enters Creating. For cloud, this takes 30–90 seconds while the instance boots. For bare metal already commissioned, this is near-instantaneous; for bare metal that needs commissioning (firmware update, OS install, burn-in), this can take hours. On success, the returned host reference populates Machine.Host and the state becomes Idle.
Idle → Configured (Configure). BigFleet calls the cluster operator's GenerateBootstrap callback (Section 13) with the requirements that need to be satisfied. The operator returns a bootstrap blob. BigFleet calls CapacityProvider.Configure with the blob. The machine enters Configuring while kubelet boots, joins the cluster, and self-labels. On success, Machine.Cluster is set and the state becomes Configured.
Configured → Idle (Drain). BigFleet calls CapacityProvider.Drain. The machine enters Draining while workloads are evicted gracefully (respecting PDBs, grace period scaled by priority gap). The kubelet unregisters. Drain time is bounded by the slowest-to-terminate workload — seconds for stateless, minutes for stateful, hours for long-running training jobs with strict PDBs and small priority gaps. On success, Machine.Cluster returns to 0 and the state becomes Idle.
Idle → Speculative (Delete, cloud only). For cloud machines, BigFleet calls CapacityProvider.Delete. The machine enters Deleting while the instance terminates. On success, Machine.Host goes back to nil and the state becomes Speculative. For bare metal, this transition doesn't happen — idle bare-metal machines just sit in the free pool (they're already paid for; there's nothing to release).
Idle is the hub. Cross-cluster machine transfers always route through it. A machine in cluster A is drained to idle, then reconfigured for cluster B via a new bootstrap — the physical resource never leaves the fleet, only its cluster binding changes.
The four lifecycle RPCs (Create, Configure, Drain, Delete) return immediately with an acknowledgement. They don't block for completion, because completion can take hours. The provider accepts the request, begins the transition, and exposes progress through state changes visible via List and Get.
BigFleet observes state transitions rather than waiting on RPC returns. Each worker cycle re-reads the relevant state and makes decisions based on what's currently stable versus what's in transition.
Idempotency. All transition RPCs are idempotent. Calling Drain on a machine already in Draining returns the current state without starting a second drain. Calling Delete on a machine already Deleting is a no-op. This is essential because BigFleet shards can restart mid-transition and need to resume observation without re-triggering the operation.
Timeouts. Each transitional state has a maximum duration. Exceeding it transitions the machine to Failed. The shard takes corrective action: for Creating → Failed, the provider cleans up via Delete and the speculative record is either retried or removed from inventory if the provider signals the quota slot is bad. For Configuring → Failed, the shard requests a fresh bootstrap blob and retries Configure, or gives up and returns the machine to Idle for different use. For Draining → Failed, the drain has exceeded its grace period; the provider escalates to forced termination (kubelet restart, then instance shutdown). For Deleting → Failed, operator intervention is usually needed.
Slow drains under priority pressure. The hardest case is a long-running training job being drained for a high-priority preemptor. If the drain takes hours because the workload has a strict PDB and the priority gap is small, the preemptor is blocked on this machine. The shard emits a shortfall to the coordinator — "drain requested N minutes ago, still in progress, priority-X preemptor waiting" — and the coordinator can either find capacity on another shard, escalate the drain (reduce grace period on subsequent cycles), or accept the wait. The drain itself continues regardless; cancellation isn't supported because partially-drained workloads have already lost state.
Ownership can be reassigned between shards at any state, but the cost differs:
| State | Drain needed? | Cadence | Mechanism |
|---|---|---|---|
| Speculative | No | Seconds | Pure bookkeeping |
| Idle | No | Seconds | Bookkeeping; no cluster binding to break |
| Configured | Yes — drain to idle first | Rare, under priority pressure | Drain workloads, unregister from cluster, then transfer |
Idle machines move between shards as cheaply as speculative ones — no cluster binding exists, so nothing has to be undone. Configured machines move only via the drain-to-idle path, and that's gated by how expensive the workload drain is (priority-scaled grace period, PDBs respected).
Cross-cluster transfers within a single shard also route through idle. If shard 7 owns a machine configured for cluster A and wants to reassign it to cluster B on the same shard, the sequence is: drain from A, idle, bootstrap for B, configured in B. The machine never leaves the shard, but its cluster binding cycles. This is how bare-metal GPUs can "float" between clusters over time without being reprovisioned.
For cloud machines that become idle, BigFleet chooses between holding them idle (still paying for them) or calling Delete to free the quota. The choice depends on expected near-term demand:
For bare-metal machines, the decision doesn't arise — bare metal has no Delete path, and holding idle is free.
Why this matters. Shard imbalance self-corrects. When workloads complete naturally, their machines transition to idle and become eligible for rebalancing. The coordinator never has to fight running workloads — it either waits for them to release capacity, or (under priority pressure) triggers a drain. Suboptimal shard assignments have a bounded cost: they persist only until the affected workloads complete or get preempted.
No distributed locking on the hot path. Each shard makes independent decisions about its own machines. Cross-shard coordination happens only for:
BigFleet has two tiers: a global coordinator that handles fleet-wide concerns, and shard controllers that handle the hot path. CapacityProviders sit below both — they're not part of BigFleet's architecture, they're pluggable backends.
The global coordinator does not make provisioning decisions. It:
CapacityProvider backendsDecision latency budget: seconds. This tier optimizes globally but infrequently.
The hot path. Every ClusterCapacityNeeds roll-up is processed here. A shard controller:
1. Receives a roll-up from a cluster operator (protobuf, ~2KB, every 10s)
2. Diffs against its inventory: what does this cluster have vs. what does it need?
3. For each unmet need, queries available capacity: free machines in inventory (fixed or elastic slots), or calls the CapacityProvider directly
4. Selects the cheapest available option using effective_cost
5. For fixed machines: assigns the machine to the cluster (in-memory, immediate)
6. For cloud machines: calls the CapacityProvider's Create method
7. Reports the assignment back to the cluster operator (which writes UpcomingNode CRDs)
Decision latency budget: <500μs per cluster. At 100 clusters per shard, a full evaluation cycle completes in <50ms.
BigFleet accepts any capacity backend that implements the CapacityProvider gRPC interface. Cloud providers, bare-metal systems (MAAS, Tinkerbell, Ironic), private cloud platforms — all plug in through the same contract.
service CapacityProvider {
// Lifecycle — all transitions are asynchronous. These RPCs return
// immediately with an acknowledgement; the actual transition is observed
// via state changes in the Machine record.
rpc Create(CreateRequest) returns (TransitionAck); // speculative → creating → idle
rpc Configure(ConfigureRequest) returns (TransitionAck); // idle → configuring → configured
rpc Drain(DrainRequest) returns (TransitionAck); // configured → draining → idle
rpc Delete(MachineRef) returns (TransitionAck); // idle → deleting → speculative
// Inventory — List returns machines in any state, filtered by caller
rpc Get(MachineRef) returns (Machine);
rpc List(ListFilter) returns (MachineList);
}
message ListFilter {
repeated MachineState states = 1; // any subset of {SPECULATIVE, CREATING, IDLE,
// CONFIGURING, CONFIGURED, DRAINING, DELETING, FAILED}
string zone = 2;
string instance_type = 3;
int32 max_results = 4;
}
message TransitionAck {
string operation_id = 1; // for idempotent retry and provider-side tracking
Machine machine = 2; // snapshot at acceptance — typically in a transitional state
}
message Machine {
string id = 1;
MachineState state = 2; // one of eight states
HostRef host = 3; // empty when speculative
string instance_type = 4;
string zone = 5;
string capacity_type = 6; // "on-demand" | "spot" | "reserved" | "bare-metal"
double price_per_hour = 7;
double interruption_probability = 8; // 0.0–1.0, hourly (forecast for speculative, observed for real)
Resources resources = 9;
map<string, string> labels = 10;
int64 transition_started_unix = 11; // when the current transitional state began (0 if stable)
string last_error = 12; // populated if state == FAILED
}
The four lifecycle methods map directly to the state transitions in Section 5: Create moves Speculative → Creating → Idle, Configure moves Idle → Configuring → Configured, Drain reverses that, and Delete moves Idle → Deleting → Speculative. Cloud providers implement all four. Bare-metal providers implement Create and Drain but typically not Delete (bare metal doesn't get "terminated" — it stays in inventory).
All four lifecycle RPCs are asynchronous. They accept the transition request and return immediately with a TransitionAck containing an operation_id and a snapshot of the machine (typically in the corresponding transitional state — Creating, Configuring, Draining, or Deleting). The actual work proceeds in the background. BigFleet observes progress by reading state via subsequent Get or List calls.
This asynchrony is essential because transitions take real time: seconds for most cloud operations, minutes for bare-metal commissioning, and hours for drains of long-running workloads with strict PodDisruptionBudgets. Blocking RPC calls for those durations would break every sensible deadline and lose state on any reconnection.
Idempotency. All four transition RPCs are idempotent. Calling Drain on a machine already in Draining returns the current operation_id without starting a second drain. Calling Delete on a machine already Deleting is a no-op. Providers key idempotency off of (machine_id, target_state): if a transition toward the target state is already in progress, reuse it.
Timeouts and FAILED state. Each transitional state has a provider-defined timeout. On timeout, the machine transitions to Failed with last_error populated. The shard takes corrective action based on which transition failed:
Creating → Failed: provider cleans up via Delete; shard retries on a different speculative slot or marks the quota slot bad.Configuring → Failed: shard requests a fresh bootstrap blob and retries Configure, or returns the machine to Idle.Draining → Failed: grace period exceeded; provider escalates to forced termination.Deleting → Failed: operator intervention usually needed; the machine is held in Failed for investigation.List returns machines in any state — speculative slots (available quota), idle machines (running but unconfigured), configured machines (running and bound to a cluster), and machines mid-transition. The caller filters by state:
List(states=[SPECULATIVE]) to populate the speculative pool from the provider's available quota.List(states=[IDLE, CONFIGURED]) periodically to reconcile their stable inventory.List(states=[CREATING, CONFIGURING, DRAINING, DELETING])) to track in-flight operations and detect ones that have been stuck too long.The unified return type works because BigFleet's Machine model already discriminates speculative from real by the presence of the host reference. A provider's List response is a snapshot of its entire inventory in whatever state each machine currently occupies.
Karpenter's interface is designed for a single cluster. BigFleet's interface differs in three ways:
Machines carry interruption_probability directly. Karpenter treats spot as a binary — available or not. BigFleet needs risk-adjusted cost to make fleet-level decisions. The provider reports the probability on every Machine; the autoscaler computes effective cost.
No scheduling simulation. Karpenter's CloudProvider is called from within a scheduling simulation that evaluates topology constraints pod-by-pod. BigFleet's CapacityProvider is called with a simple request: "I need 8 machines matching these requirements in this zone." No pods, no topology simulation, no scheduler state.
Batch-oriented. BigFleet sends batched Create requests (64 machines at once, not 1). The provider handles the cloud API batching internally. This is critical for cloud providers where API rate limits are per-call, not per-machine.
A bare-metal CapacityProvider is simpler than a cloud provider. List returns the provider's full inventory across all three states: racks that are commissioned but unallocated appear as speculative machines; racks that have been claimed but aren't cluster-joined appear as idle; racks actively serving a cluster appear as configured. Create assigns a speculative machine from the free pool (sub-second). Delete typically isn't implemented — bare metal stays in inventory, transitioning back to speculative when no longer in use. The interruption_probability is 0.0. The price_per_hour is 0.0 (or amortized cost, depending on how the operator wants to model it).
For bare-metal machines that require commissioning (firmware update, OS install, burn-in), the provider can advertise availability: "commissioning" with an estimated_ready_time. BigFleet treats this identically to a cloud node that's launching — the UpcomingNode CRD tracks the progress.
The shard does not react to individual roll-ups. Roll-ups arrive asynchronously from cluster operators — one every 10 seconds per cluster, staggered. The shard ingests each one and updates an in-memory materialised view of all outstanding needs across all its clusters. A separate worker loop continuously evaluates this view against the shard's full inventory and makes assignment decisions.
No roll-up triggers an immediate action. No cluster "causes" a preemption. The worker simply looks at the total picture — everything everyone needs, everything the shard has — and assigns capacity from the top of the priority ordering downward. When it runs out, the bottom doesn't get served yet. That's preemption, but it's a consequence of the ordering, not an event.
This is how Borg works: the scheduler continuously scans the pending queue and matches jobs to machines. It doesn't react to individual job submissions.
The shard maintains two data structures:
The needs table. Every CapacityNeed from every cluster's latest roll-up, flattened into a single priority-sorted list. When a new roll-up arrives from cluster X, it replaces all of cluster X's entries (full replacement semantics, as defined in the operating model).
type Need struct {
Cluster ClusterID
Requirements []NodeSelectorRequirement
Resources Resources
Priority int32
Count int32
Spread []TopologySpread
InterruptionPenalty float64 // from CapacityRequest, $/hr cost of interruption
ReclamationPenalty float64 // from CapacityRequest, $ cost of reclaiming this capacity
}
// Always sorted: highest priority first, then by arrival time.
type NeedsTable []Need
The inventory. Every machine the shard owns, with its current state (see Section 5 for the full state model).
type Machine struct {
ID MachineID
Host *HostRef // nil = speculative, set = real
Cluster ClusterID // 0 = speculative or idle; set = configured
Profile Profile // instance type, zone, resources
CapacityType string // bare-metal, reserved, spot, on-demand
PricePerHour float64
IdleSince time.Time // when this machine became idle (for hold/release policy)
}
The worker runs continuously. Each cycle, it walks the needs table top-down and tries to satisfy each need from available capacity. The cycle produces a set of actions — assignments, provisions, and reclamations — that are executed asynchronously.
func (s *Shard) runCycle() []Action {
var actions []Action
// Snapshot current state.
needs := s.needsTable.Snapshot() // sorted by priority desc
inventory := s.inventory.Snapshot()
// Phase 1: Walk needs top-down. Assign capacity.
for _, need := range needs {
have := inventory.CountConfigured(need.Cluster, need.Profile())
want := int(need.Count)
if have >= want {
continue // Already satisfied.
}
deficit := want - have
// Prefer idle machines (cheapest: one Bootstrap call, no Create).
idle := inventory.FindIdle(need.Profile(), deficit)
if len(idle) > 0 {
actions = append(actions, bootstrap(idle, need.Cluster))
inventory.MarkBootstrapping(idle, need.Cluster)
deficit -= len(idle)
}
if deficit <= 0 {
continue
}
// Fall back to speculative (Create + Bootstrap).
speculative := inventory.FindSpeculative(need.Profile(), deficit)
if len(speculative) > 0 {
best := cheapestBy(speculative, effectiveCost(need))
actions = append(actions, provisionAndBootstrap(best, need.Cluster))
deficit -= len(best)
continue
}
// No idle, no speculative quota in this shard.
// Phase 2 will handle via preemption or escalate to coordinator.
}
// Phase 2: Check for priority inversions.
// Are there machines assigned to lower-priority needs that should
// be serving higher-priority unsatisfied needs?
actions = append(actions, s.resolveInversions(needs, inventory)...)
// Phase 3: Reclaim excess.
// Are there machines assigned to clusters that no longer need them?
actions = append(actions, s.reclaimExcess(needs, inventory)...)
return actions
}
Phase 1 tries idle machines first, speculative only as a fallback. This ordering matters for three reasons:
Speed. Idle → Configured is one callback and a kubelet restart, ~30 seconds to ready. Speculative → Idle → Configured adds a Create call (30–90s for cloud) before the Bootstrap step. Idle beats speculative by roughly a factor of three on latency-to-ready.
Cost. An idle cloud instance is already costing money. Using it for the next CR amortises that cost over productive work instead of wasting it. Meanwhile, using a speculative slot starts a fresh billing cycle on a new instance. Holding idle capacity while provisioning fresh speculative capacity for a different need is pathological — two machines billed for what one could serve.
Rate limits. Every Create call consumes a token against the cloud provider's API rate limit. AWS's RunInstances budget is small (~2/sec sustained per account-region). Preferring idle machines that need no Create call preserves rate-limit headroom for genuinely new demand.
Within each tier, the existing cost functions apply. Among idle candidates, the tiebreaker is reclamation penalty (prefer reassigning machines that have less accumulated operational value). Among speculative candidates, the tiebreaker is effective cost (price + interruption_probability × interruption_penalty).
After Phase 1, some high-priority needs may be unsatisfied. Meanwhile, some machines are serving lower-priority needs. Phase 2 identifies these inversions and corrects them — but the correction is smarter than "pick the lowest priority victim."
Victim selection is multi-dimensional. When the shard needs to reclaim capacity, it scores every candidate machine across four dimensions:
func (s *Shard) scoreVictim(machine Machine, preemptorPriority int32) float64 {
priorityGap := float64(preemptorPriority - machine.AssignedPriority)
return (priorityGap * weightPriority) +
(1.0 / machine.EstimatedDrainTime.Seconds() * weightSpeed) +
(1.0 / machine.InterruptionPenalty * weightPenalty) +
(1.0 / machine.ReclamationPenalty * weightReclaim)
}
Higher score = better preemption candidate. The four factors:
Priority gap. A machine at priority 100000 is a better victim than one at priority 900000 when the preemptor is at 1000000. This is the baseline — lower-priority victims are preferred.
Estimated drain time. A node running a stateless serving pod with a 5-second graceful termination drains in seconds. A node running a stateful workload with a 30-minute termination grace period and strict PodDisruptionBudgets takes half an hour. When the preemptor is urgent, the shard prefers the node it can actually free up fastest.
Interruption penalty. The victim's own CapacityRequest declared how costly interruption is. A workload with interruptionPenalty: 0.50 is signalling "I'm cheap to interrupt — I checkpoint, I restart fast, it's fine." A workload with interruptionPenalty: 50.00 is signalling "interrupting me destroys hours of work." The shard prefers victims that said they don't mind.
Reclamation penalty. Interruption penalty captures the cost of the workload being interrupted. Reclamation penalty captures the cost of losing this specific machine — operational value that has accumulated at its current location. A GPU that's been burned in and thermally characterised, a node whose performance profile has been learned by the workload's scheduler, a machine where accumulated caches or warm-up time represent substantial investment — all of these raise the cost of picking this machine as a victim, even when the workload would technically survive the interruption.
Reclamation penalty and interruption penalty are distinct concerns. A training job with good checkpointing might have a low interruption penalty (it can resume fast) but a high reclamation penalty (it's profiled against these specific GPUs and getting different ones would cost re-characterisation). A stateless serving workload might have low values for both. The shard prefers victims that are cheap across both dimensions.
The weights between these dimensions shift based on how saturated the fleet is. Under normal conditions, priority dominates — the lowest-priority victim is preempted gracefully. Under extreme pressure (no free capacity, no elastic capacity, critical workload blocked), drain speed dominates — the shard wants the node it can get right now.
Drain aggressiveness scales with the priority differential. BigFleet's reclaim action includes a gracePeriod that the operator passes through to the kubelet. The grace period is computed from how far apart the preemptor and victim are in priority:
func drainGracePeriod(preemptorPriority, victimPriority int32) time.Duration {
gap := preemptorPriority - victimPriority
switch {
case gap > 900000:
// Extreme gap (e.g., 1000000 vs 0). Minimal grace.
return 10 * time.Second
case gap > 500000:
// Large gap. Reduced grace.
return 30 * time.Second
case gap > 100000:
// Moderate gap. Short grace.
return 2 * time.Minute
default:
// Small gap. Full graceful shutdown.
return 10 * time.Minute
}
}
This is a gradient, not a binary. A best-effort research job (priority 0) yields almost instantly to a production training run (priority 1000000). Two production workloads with a small priority gap (900000 vs 1000000) get the full graceful treatment — PDBs respected, termination periods honoured, no rush.
At the extreme end — fleet completely saturated, no bare metal, no cloud capacity, critical workload blocked — the shard picks the victim with the lowest interruption penalty, the fastest drain time, and the lowest priority, and reclaims it with minimal grace. The workload said it was cheap to interrupt, and the fleet needs the capacity now.
What the operator does with the grace period. The operator receives a reclaim instruction with a grace period and signals the kubelet accordingly. The kubelet still handles pod eviction — but the grace period constrains how long it waits. A 10-second grace period means "evict now, don't wait for long-running termination handlers." A 10-minute grace period means "drain gracefully, respect PDBs, let pods finish cleanly." BigFleet doesn't bypass Kubernetes safety mechanisms — it adjusts how much time they get.
There's no "preemption event." There's a continuous evaluation where the shard looks at everything and assigns optimally. But from the outside, the story looks like this:
t=0. Cluster A is running a priority-500000 batch pipeline across 64 GPU nodes. The shard's needs table has these entries. The inventory shows 64 GPU machines assigned to Cluster A. Everything is balanced.
t=7s. Cluster B's operator sends a roll-up: 64 GPU nodes needed at priority 1000000. The shard updates the needs table. It doesn't do anything yet.
t=10s. The worker cycle runs. It walks the needs table top-down:
1. Cluster B's need (priority 1000000): needs 64 GPU nodes, has 0. Deficit: 64. No free GPU machines. No elastic GPU capacity available. Unsatisfied.
2. Cluster A's need (priority 500000): needs 64 GPU nodes, has 64. Satisfied.
Phase 2 detects the inversion: 64 GPU machines are serving priority 500000 while a priority 1000000 need is unsatisfied. It emits a preempt action.
t=10s onward — the mechanics:
1. The shard scores Cluster A's 64 GPU machines: low priority (500000 vs 1000000 = large gap), low interruption penalty (batch pipeline with checkpointing), estimated drain time ~30 seconds (stateless workers). High preemption score.
2. The priority gap is 500000. The drain grace period is 30 seconds.
3. BigFleet tells Cluster A's operator: "nodes gpu-east-0142 through gpu-east-0205 are being reclaimed, grace period 30s."
4. The operator signals each node with a 30-second grace period. The kubelet evicts pods, the batch workers checkpoint, and the nodes drain.
5. As each node drains, the machine transitions to idle in the shard's inventory.
6. The next worker cycle sees free GPU machines and Cluster B's unsatisfied need. It assigns them to Cluster B.
7. Cluster B's operator writes UpcomingNode CRDs. The nodes register with Cluster B. kube-scheduler places the training pods.
8. Cluster A's pods are now Pending. Cluster A's operator creates new CapacityRequests. The next roll-up includes them at priority 500000. The shard provisions replacement capacity — spot instances, on-demand, whatever is cheapest and available.
The key insight: no component "decided to preempt Cluster A for Cluster B." The worker walked the priority list, assigned what it could, noticed an inversion, and corrected it. If Cluster B's need had arrived 30 seconds later when a new batch of bare-metal machines came online, there would have been no preemption — the free machines would have been assigned in Phase 1. The outcome depends on the state of the world at evaluation time, not on the order of arrival.
Priority is the sole ordering. The worker walks needs top-down by priority. Same priority, same treatment. No cluster is special. No manual overrides.
Preemption is a last resort. Phase 1 tries free capacity and elastic provisioning first. Phase 2 only fires when Phase 1 leaves high-priority needs unsatisfied. In practice, the natural slack in a fleet — bare-metal machines between assignments, recently reclaimed capacity — means most cycles have free capacity and Phase 2 produces zero actions.
Aggressiveness is proportional to urgency. A small priority gap (1000000 vs 900000) gets a full graceful drain — PDBs respected, termination periods honoured, minutes of grace. A large gap (1000000 vs 0) gets a short drain. The victim's own interruption penalty factors in — workloads that declared themselves cheap to interrupt are preferred over workloads that declared interruption costly. Nobody gets yanked without warning; the question is how much warning.
No preemption at equal priority. If two needs have the same priority number, neither preempts the other. Both compete for free and elastic capacity. If there isn't enough, both wait. The worker doesn't pick favourites.
The blast radius is one cluster. A preemption event displaces one cluster's lower-priority workloads. Those workloads go Pending, generate CapacityRequests, and get served on the next cycle — at their priority level, with whatever capacity is available.
Within a shard, inversion resolution is a local operation — the shard sees all its clusters' priorities and its full inventory. Cross-shard preemption is harder and rarer.
When a shard's worker cycle leaves high-priority needs unsatisfied after both Phase 1 and Phase 2, it emits a shortfall to the global coordinator (Section 9). The global coordinator's response follows a priority order that mirrors Phase 1's own preference within a shard:
1. Reassign idle machines from other shards. If another shard has idle inventory matching the profile, move it. Seconds — bookkeeping plus state transfer, no drain needed, no new provisioning.
2. Reassign speculative quota from other shards. If another shard holds unused quota for the profile, shift the quota allocation. Seconds — pure bookkeeping — but the receiving shard still has to Create + Bootstrap, which takes 30–90 seconds.
3. Cross-shard preemption. Drain configured machines from a lower-priority workload on another shard, then transfer. Minutes; disruptive. Only under genuine priority pressure.
Cross-shard preemption happens when all other options are exhausted: no idle machines available anywhere, no speculative quota available anywhere, and the shortfall priority is high enough to warrant workload disruption. The global coordinator queries all shards for their lowest-priority configured machines matching the shortfall profile. Shard 47 reports 64 GPU nodes at priority 100000 (best-effort research). The global coordinator instructs Shard 47 to drain those machines to idle and transfer ownership to the requesting shard, which then bootstraps them for the high-priority workload.
This is intentionally expensive — it requires global coordination and workload disruption. But it's also intentionally rare. Most shortfalls resolve through idle-machine or speculative-quota reassignment because shards are roughly balanced. Cross-shard preemption only triggers when the entire fleet is saturated at a specific resource type. If this happens frequently, it's a signal to buy more capacity, not to optimise the preemption path.
When provisioning elastic capacity (Phase 1, step 2), the shard picks the cheapest available option:
func (m *Machine) EffectiveCost(interruptionPenalty float64) float64 {
return m.PricePerHour + (m.InterruptionProbability * interruptionPenalty)
}
The interruptionPenalty is specified in the CapacityRequest by whoever created it. A training job that can't tolerate interruption sets a high penalty. A batch pipeline with checkpointing sets a low one. BigFleet doesn't interpret it — it just plugs it into the cost function.
The previous section introduces the field but doesn't say how to pick a value. Since interruptionPenalty directly determines whether your workload runs on spot or on-demand, the creator of the CapacityRequest needs a mental model for what the number actually represents.
The question the penalty answers is: if a node running this workload gets yanked, how many dollars of work does that destroy?
That's the whole field. Not "how important is this workload," not "how much do I hate interruptions." Dollars of wasted compute. The reason the formula works is that interruption_probability × interruption_penalty gives you the expected hourly cost of being interrupted, which can be compared directly against the price differential of using non-interruptible capacity.
The general formula:
penalty = (hours_of_work_lost_on_interruption) × (hourly_cost_of_all_affected_resources)
Two terms, and the second one is subtle — the blast radius matters. A node dying in a 64-GPU training cluster doesn't just waste that node's time; it wastes the whole cluster's time, because the other 63 nodes were participating in collective operations and have to roll back too.
Worked examples:
Stateless web server, one replica of a pooled service. If the node dies, Kubernetes restarts the pod on another node in 20 seconds. Some in-flight requests retry. If you're running 100 replicas, you've lost 1% capacity for 20 seconds. Work lost: negligible. Penalty: ~$0.50. Spot is always cheaper.
Stateful database follower with Raft replication. If a follower dies, no user impact. If the leader dies, election takes 5–10 seconds and clients retry. Work lost: tens of seconds of one node's time. Penalty: ~$5. Spot still fine — replication factor absorbs it.
Short batch job, 30 minutes, no checkpointing. Worst case, the node dies at minute 29 and you redo the whole job. At $6/hr, that's $3 of wasted compute. Penalty: $6 (upper bound of a run). Spot often still wins: 10% × $6 = $0.60/hr expected interruption cost vs $4.20/hr saved on the instance.
Long training job, 64 GPUs, checkpoints every 30 minutes. A node dies, the whole cluster rolls back to the last checkpoint, you lose up to 30 minutes of 64 GPUs. 64 × $31/hr × 0.5hr = ~$1,000 penalty. On spot at 5% interruption probability that's $50/hr of expected cost per node — which might still beat on-demand depending on the savings, but the math is tight.
Long training job, 10,000 GPUs, no intermediate checkpointing. Day 3 of a 7-day run. A single node dying restarts the entire job. You've spent 10,000 × $31 × 72 = $22M of compute so far, and that's the penalty. Expected cost of interruption on spot: 5% × $22M = $1.1M/hr per node. Nobody should run this on spot. Set penalty high enough (say $1e12) that the cost function always picks on-demand or reserved.
Patterns by workload class:
| Workload | Penalty formula | Typical range |
|---|---|---|
| Stateless replica (pooled service) | recovery_time × this_node_cost | $0.50–$2 |
| Stateful pod with replication | failover_impact × (replicas × cost) | $5–$20 |
| Short batch, no checkpoint | full_runtime × cluster_cost | equal to hourly cost |
| Long batch, with checkpoints | checkpoint_interval × cluster_cost | scales with cluster size |
| Long training, no checkpoints | time_since_start × cluster_cost | grows unboundedly |
The last row is why checkpointing exists in training frameworks. Without it, the interruption penalty grows linearly with how long the job has been running, so an hour-29 interruption costs 29× more than an hour-1 interruption. Checkpointing caps the penalty at checkpoint_interval × cluster_cost regardless of how long the job runs. That cap is usually the single most important operational parameter for a large training run, because it directly determines whether spot capacity is economically usable.
Priority and interruption_penalty are independent. A priority-1000000 training job has a high penalty because interruption destroys real dollars of work. A priority-1000000 stateless serving workload has a low penalty because it restarts in seconds — even though it's equally important. BigFleet makes no assumption that high priority means high penalty; the CapacityRequest must specify both.
When in doubt, overestimate. The penalty is a safety valve, not a cost lever. If you set it too low and your workload was more interruption-sensitive than you thought, you pay in wasted compute. If you set it too high, you pay a small premium in slightly-preferred on-demand over spot. The asymmetry argues for erring high — especially for workloads where you haven't measured actual interruption cost yet.
interruptionPenalty and reclamationPenalty sound similar but answer different questions.
Interruption penalty answers: if this workload gets interrupted and has to restart on different capacity, what's the cost? This captures wasted compute, lost progress, failover impact — the cost of the work being disrupted.
Reclamation penalty answers: if BigFleet takes these specific machines away from this workload, what's the cost beyond just restarting the work? This captures operational value that has accumulated at the current location — value that's tied to the specific machines, not the workload generally.
The classic example is GPU burn-in. A training job spins up on 64 H100s and spends the first few hours burning them in: running stress tests, characterising thermal behaviour, identifying slow or unreliable units. That knowledge becomes part of how the workload schedules — which GPUs get tensor-parallel ranks, which get pipeline-parallel, which are flagged for retry on failure. If BigFleet reclaims those 64 GPUs and provisions 64 different ones, the workload doesn't just lose the ongoing training — it loses the characterisation and has to redo burn-in. The interruption cost and the reclamation cost are both real, and they're different costs.
The general formula:
reclamation_penalty = (time_to_recreate_operational_value) × (workload_running_cost) + (risk_premium_for_worse_replacement)
Worked examples:
Stateless web server. Nothing about the specific machine matters. Workload schedules fresh wherever it lands. Reclamation penalty: $0.
Stateful database pod with attached persistent volume. The machine isn't special — the PV is. K8s handles this via volume attachment portability. Reclamation penalty: $0 (the workload's state isn't coupled to the machine identity).
Training job with burned-in GPU characterisation, early in training. Losing the burn-in means 2–4 hours of re-characterisation on 64 GPUs. 64 × $31/hr × 3hr = ~$5,950. Plus risk premium for possibly getting worse GPUs on replacement — maybe another $2,000. Reclamation penalty: ~$8,000.
Training job with 72 hours of runtime, specific GPU scheduling tuned for this hardware. The operational value accumulated over three days of observation — knowing which GPUs slow down under sustained load, which cards show occasional ECC events, which nodes correlate with higher throughput. Replacement means losing all of this plus re-tuning. Reclamation penalty: tens of thousands of dollars, sometimes more than the interruption penalty itself.
Long-running stateful service with in-memory caches warmed against specific-host disk layouts. The caches represent accumulated work that's not captured in a checkpoint. Replacement means re-warming, which could take hours of degraded performance. Reclamation penalty: scales with cache warmup time × cost of degraded service.
Workloads that genuinely want to stay put should set this high. A training job that's invested in its current hardware tells BigFleet "don't move me unless you absolutely have to" by setting reclamation_penalty to something like $10,000 or more. The Phase 2 victim scoring will prefer virtually any other candidate.
Don't use reclamation_penalty as a general "this is important" lever. Priority is for that. Reclamation penalty is specifically about the cost of losing this machine versus losing any equivalent machine. If the workload would be equally happy on any machine matching its profile, reclamation penalty is zero. If it's tied to specific hardware, it's non-zero. Most workloads are in the first category.
Reclamation penalty decays on machine change. If a machine is reclaimed and replaced, the workload's reclamation penalty on the new machine starts at zero and grows as operational value accumulates. This is implicit in the field's definition — reclamation penalty measures value tied to this specific machine, so it resets when the machine changes. Clients that track this (good training frameworks already do, for their own reasons) can update the CapacityRequest over time.
After assignments and inversions are resolved, the worker checks for machines configured for clusters that no longer need them. A cluster's roll-up said it needs 10 GPU nodes but has 64 configured — the job finished, pods were deleted, CRs were garbage collected. The worker reclaims the 54 excess machines.
Reclaimed machines become idle, not speculative. The drain-and-unregister transition (Configured → Idle) returns the machine to the idle pool, ready for the next need. This is important: if another CR arrives within minutes that could use these machines, allocating from idle is much cheaper than re-provisioning from speculative.
Picking which 54 to reclaim uses reclamation-cost thinking as a secondary factor after price. The primary sort is "most expensive per hour first" — return on-demand before bare metal, because holding on-demand costs money and holding bare metal doesn't. When two candidates have similar per-hour cost, reclamation penalty breaks the tie: prefer to release the machine with the lower reclamation penalty, because the other one has more accumulated operational value. Free up the cheap-to-replace capacity first, hold on to the hard-to-replace capacity even if it's a bit more expensive.
Idle → Speculative happens lazily, per-provider. Idle machines that stay idle eventually get released back to speculative (the Idle → Speculative transition) to free their cloud billing. Policy varies by capacity type:
| Capacity type | Hold duration | Rationale |
|---|---|---|
| Bare metal | Forever | No cost to hold; no "release" path exists |
| Reserved instances | Forever | Commitment is paid regardless; release frees nothing useful |
| On-demand | A few minutes | Idle on-demand costs real money per second |
| Spot | ~1 minute | Spot is cheapest but interruptible; release fast if not needed |
The exact timeouts are tuning knobs per provider. They trade off re-provisioning cost against idle holding cost. A shard under steady demand can hold longer; a shard with bursty demand benefits from releasing faster.
Each shard maintains an in-memory inventory of its assigned machines. The minimal per-machine record:
| Field | Size | Description |
|---|---|---|
| machine_id | 8B | BigFleet's internal identifier |
| host_ref | 0–24B | Provider + resource ID; empty means speculative |
| cluster_id | 2B | Configured for this cluster (0 = speculative or idle) |
| shard_id | 2B | Owning shard |
| instance_type | 2B | Index into instance type catalog |
| zone | 1B | Zone identifier |
| capacity_type | 1B | bare-metal / reserved / spot / on-demand |
| state | 1B | speculative / creating / idle / configuring / configured / draining / deleting / failed |
| cpu_capacity | 2B | Allocatable CPU (millicores / 100) |
| memory_capacity | 4B | Allocatable memory (MiB) |
| gpu_count | 1B | Number of GPUs |
| price_cents | 2B | Cost per hour in cents |
| last_heartbeat | 4B | Unix timestamp |
| transition_started | 4B | Unix timestamp when current transitional state began (0 if stable) |
Total: ~30–55 bytes per machine depending on state (speculative records omit host_ref and last_heartbeat). At 500K machines per shard with a mix, roughly 20MB — fits in L3 cache.
The state field tracks all eight machine states from Section 5: three stable (speculative, idle, configured), four transitional (creating, configuring, draining, deleting), and one failure state. The host_ref field is the speculative-vs-real discriminator: nil means the machine hasn't been provisioned (speculative, or creating), set means real hardware exists. The transition_started timestamp lets the shard detect stuck transitions — if a machine has been in configuring for longer than the expected timeout, the shard takes corrective action.
The global coordinator maintains a summary per shard, not per machine:
| Field | Size | Description |
|---|---|---|
| shard_id | 2B | Shard identifier |
| total_machines | 4B | Total machine count |
| free_machines | 4B | Unassigned machine count |
| per_type_counts | ~64B | Counts by instance type |
| per_zone_counts | ~32B | Counts by zone |
| utilization | 8B | Aggregate CPU/memory utilization |
| shortfalls | ~128B | Unsatisfied needs by profile |
Total: ~240 bytes per shard. At 200 shards, global state is 48KB. Still fits in L1 cache.
Utilization tells the global coordinator what shards have. It doesn't tell it what they need. Without explicit demand signalling, rebalancing is either reactive (wait for shards to fail to satisfy requests) or guesswork (move capacity based on stale utilization). Neither scales.
So each shard's report to the global coordinator includes a shortfall list — needs that Phase 1 and Phase 2 couldn't satisfy in the last cycle:
type Shortfall struct {
Profile Profile // instance type, zone, resources
Priority int32
Count int32 // how many machines needed
Age int32 // cycles this has been unsatisfied
InterruptionPenalty float64 // from the original CapacityRequest
}
A shortfall is emitted when a shard cannot satisfy a need from its own idle or speculative inventory and cannot resolve it through in-shard preemption. It means "this need is blocked and can't be resolved locally."
The global coordinator aggregates shortfalls across shards and generates rebalancing actions. The priority order matches Phase 1's intra-shard preference — idle before speculative, drain-first only as a last resort:
1. Reassign idle machines. If Shard 12 has a GPU shortfall and Shard 47 has 100 idle GPUs matching the profile, reassign them. Seconds — bookkeeping plus state transfer. The receiving shard then bootstraps them for its cluster (one callback, ~30 seconds to ready). No Create call, no cloud API rate-limit consumption.
2. Reassign speculative machines. If Shard 12 has a GPU shortfall and Shard 47 has speculative GPU quota it isn't drawing against, reassign the quota. Seconds for the reassignment itself, but the receiving shard must then Create + Bootstrap (30–90 seconds). Burns a RunInstances token.
3. Cross-shard preemption. Only if the shortfall priority is high enough and no idle or speculative capacity exists anywhere. Involves workload disruption on the donating shard — the preemption mechanism of Section 8, but cross-shard. Routes through idle: drain on donor, transfer ownership, bootstrap on receiver.
Shortfalls also inform longer-term decisions:
Shortfalls are bounded and aged. A shard tracks at most 100 shortfalls (the highest-priority ones). A shortfall that remains unsatisfied for more than 5 cycles is escalated to fleet-level telemetry for human attention — something is wrong that BigFleet can't fix automatically.
Machines don't heartbeat to BigFleet directly. Clusters report their node state to their operators. Operators include node health in their ClusterCapacityNeeds roll-ups (or a separate health report). BigFleet reconciles its inventory against cluster reports on a 30-second cycle.
If a machine stops appearing in cluster reports, BigFleet marks it as unhealthy after 3 missed reports (90 seconds). For cloud machines, it verifies with the provider (Get). For bare-metal, it checks the provider's health endpoint. If confirmed dead, it provisions a replacement.
BigFleet's inventory is the source of truth for what it has provisioned. It knows because it created every machine (or was told about every fixed machine at startup). The cluster's node list is a secondary signal used for health verification, not for inventory tracking.
This is important: BigFleet never discovers machines. It either provisioned them or was told about them during shard initialization. This eliminates the entire class of "phantom node" bugs where the inventory tracker disagrees with reality.
Cloud quota is not architecturally special. It's represented as speculative machines — machine records with no host reference — distributed across shards the same way any other capacity is distributed.
When BigFleet is configured with a cloud account (say, AWS account X in eu-west-1 with quota for 1,000 p5.48xlarge instances), the provider reports this availability through List — returning 1,000 speculative Machine records, each with a p5.48xlarge profile, eu-west-1 zone, and no host reference. The coordinator distributes these speculative machines across shards. Shard 12 gets 50. Shard 47 gets 50. Each shard's CapacityProvider binding is restricted to the machines assigned to that shard.
This is mechanically identical to bare-metal machine assignment. Instead of "Shard 12 owns physical machines in racks 7-14," it's "Shard 12 owns 50 speculative p5.48xlarge records backed by AWS account X." Both are capacity the shard can draw from without coordinating with anyone.
When a shard decides to allocate one of its speculative machines to a workload, it calls the CapacityProvider's Create, which actually runs an EC2 instance. The returned instance ID populates the machine's Host field — the machine transitions from speculative to real. No new record is created; the same MachineID now refers to a real running instance.
When the workload releases the machine and the shard decides to reclaim it, the CapacityProvider's Delete terminates the EC2 instance. The machine's Host field goes back to empty — the machine transitions from real back to speculative, available for the next allocation.
Cloud quota isn't static. AWS might raise or lower your limit. A CapacityReservation might become available. An RI might be purchased. Subsequent List calls return an updated view — more speculative records when availability increases, fewer when it decreases (but never fewer real ones, even if quota shrinks below current usage).
If the provider reports that a quota slot is no longer available but a real machine is already running against it, the running machine is unaffected. Only speculative (unbound) slots disappear. This matches how cloud quotas actually work — shrinking a quota doesn't terminate running instances.
Each shard calls its CapacityProvider directly to transition speculative machines to real (and back). Rate limiting, request batching, circuit breakers, retries — all of this is the provider's responsibility, not BigFleet's architecture. A well-written AWS provider uses the AWS SDK (which handles rate limiting), batches RunInstances calls where possible, and maintains local caches of instance type availability.
Provider API rate limits at 100M-node scale:
| Provider | Create rate (sustained) | Burst | Practical throughput |
|---|---|---|---|
| AWS EC2 | ~2 RunInstances/sec | 1,000 tokens | ~8,200 instances/hr per account-region |
| Azure | 25 VM creates/sec | 1,500/min | ~90,000 VMs/hr per subscription-region |
| GCP | Per-project quota | Varies | Requires explicit quota increase |
A single account per region isn't enough at scale. At AWS's rate, a region-wide demand spike of 50,000 instances would take 6 hours to satisfy from one account. The solution is multiple accounts — the CapacityProvider implementation can fan requests across 10 AWS accounts in the same region, giving ~82,000 instances/hour.
But this multi-account fan-out is a provider implementation detail, not BigFleet architecture. From BigFleet's perspective, each provider exposes capacity as a pool of speculative machines, and BigFleet doesn't care whether the provider is backed by one account or ten.
The most important reliability property: clusters continue operating when BigFleet is down. BigFleet manages capacity; it doesn't manage workloads. If every BigFleet component fails simultaneously:
This is Borg's most important design principle: "already-running tasks continue to run even if the Borgmaster goes down." BigFleet inherits this property by design — it operates below the scheduling layer.
The global coordinator runs as 3 replicas using Raft for leader election and log replication. One replica is leader, two are followers.
Steady state. The leader handles all writes — shard membership changes, cluster-to-shard assignments, quota distribution, shortfall-driven rebalancing decisions. Every write is replicated to the followers via Raft before being committed. Followers don't serve writes; they stay current so they can take over.
Leader failure. When the leader crashes or is partitioned, followers' heartbeat timers expire (150–300ms typical). A follower transitions to candidate state, increments its term, and requests votes from the other replicas. With 3 replicas a candidate needs 2 votes (majority), which it gets from itself plus one follower. Election typically completes in under a second.
During the gap. No new writes happen for the duration of the election — no shard splits, no rebalancing, no new cluster assignments. But the shards don't stop. They keep receiving roll-ups, making decisions, and provisioning from their existing capacity and quota allocations. In-shard preemption continues to work. Only cross-shard rebalancing and preemption pause, and only for the duration of the election.
This is the static stability property. The data plane (shards) operates autonomously while the control plane (global coordinator) is having its election. Running workloads are unaffected. New provisioning from each shard's existing allocations continues. Only fleet-wide coordination is briefly unavailable.
After election. The new leader reads its log and reconstructs state. It reconnects to shards and resumes accepting shortfall reports. Rebalancing decisions pick up where they left off — any shortfall that was unresolved before the failover is still reported on the next cycle, so nothing is lost.
What the coordinator stores. The coordinator's state is small and simple:
Total: ~500MB at 100M-node scale. Write rate at steady state: ~10/sec (Section 13). This is tiny relative to any real database.
The canonical choice is embedded Raft with local disk storage, using a library like etcd/raft or hashicorp/raft on top of BoltDB or similar. Three reasons:
1. No bootstrapping problem. If BigFleet depended on an external etcd running in a Kubernetes cluster that BigFleet itself manages, you'd have a circular dependency. Embedded Raft avoids this — the coordinator is its own consensus group with no external state dependency.
2. The state is small enough. 500MB of state and 10 writes/sec doesn't justify the operational cost of a full database. etcd or CockroachDB would work, but they're overkill for this workload.
3. Static stability requires minimal dependencies. Every external service on the control plane path is something that can fail independently. Embedded Raft means the coordinator's availability depends on exactly one thing — its own replicas.
External etcd is a valid alternative for organisations that already operate etcd at scale and prefer reusing operational expertise over minimising dependencies. The choice is operational, not architectural — BigFleet's design doesn't care which approach is used, as long as the state is durable, replicated, and accessible with low latency from the coordinator replicas.
Split-brain prevention. If the old leader is merely partitioned rather than crashed, it might think it's still leader. Raft prevents it from committing anything — it can't get a majority for writes on its side of the partition. More importantly, the new leader operates at a higher Raft term. Every global coordinator instruction to a shard carries an epoch (the Raft term). Shards ignore instructions with a stale epoch. Combined with the fencing tokens described below, there's no way for a zombie leader to issue conflicting rebalancing or preemption commands.
Shards don't track which replica is the leader. They talk to any of the three replicas, tagging every message with the highest coordinator term they've seen. Followers either proxy to the leader or return a redirect with the leader's address. This is the standard Raft-aware client pattern — the same model used by etcd and CockroachDB clients.
Crucially, the direction of communication matters:
Leader → shard (instructions). Every instruction (rebalance machines, reassign quota, cross-shard preemption) carries the coordinator term. The shard compares this to the highest term it has seen. If the instruction's term is stale — meaning it came from an old leader that doesn't realise it was replaced — the shard rejects it. If the term matches or advances, the shard accepts it and updates its high-water mark.
Shard → coordinator (shortfall reports). Shards send shortfall reports to any replica. If they hit a follower, the follower forwards to the current leader. If they hit the leader, it's handled directly. If they hit a replica that's in the middle of an election, the report retries with exponential backoff. A shortfall delayed by a second during failover doesn't matter — the next cycle's shortfall will contain the same information.
Shard ↔ CapacityProvider. These interactions don't involve the coordinator at all. Shards operate on their existing allocations (machines and quota slices) independently. This is why the data plane keeps running during coordinator failover — shards don't need the coordinator to service roll-ups or provision capacity from their own allocations.
Zombie leader scenarios. If the old leader is partitioned but still running, it may keep trying to send instructions. Three things prevent damage: (1) Raft prevents it from committing new state changes on its side of the partition, so it has no new instructions to send. (2) Any instructions it does try to send carry its old term, which shards reject. (3) Shards that briefly talked to the old leader before the partition don't get corrupted state — the term number they saw is just their high-water mark, and the new leader's higher term will supersede it on the next instruction.
The global coordinator's durable state is small. At 100M nodes across 200 shards, the full state is roughly:
| Data | Size | Notes |
|---|---|---|
| Shard membership | ~10KB | 200 shards, ~50 bytes each |
| Cluster-to-shard assignments | ~200KB | 20,000 clusters, ~10 bytes each |
| Machine-to-shard assignments | ~500MB | 100M machines, ~5 bytes each (just shard pointer) |
| Quota allocations | ~50KB | Per-provider, per-shard quota slices |
| Provider registry | ~1KB | List of configured CapacityProviders |
The machine-to-shard assignment table is the largest. At 500MB, it still fits comfortably in memory. Writes to it are rare — only when machines join the fleet, leave the fleet, or are rebalanced between shards.
The canonical design uses embedded Raft for replication and persistence. The three coordinator replicas form a Raft group. State is held in-memory and periodically snapshotted to local disk. Every write goes through Raft's consensus protocol before being applied. Libraries like etcd/raft (used by etcd, TiKV, CockroachDB) or hashicorp/raft (used by Consul, Nomad) handle the mechanics.
Embedded Raft is preferred over an external store because:
That said, operators who already run etcd at scale (e.g., for their Kubernetes clusters) may prefer to point BigFleet at an existing etcd cluster rather than operate another stateful service. This is a valid alternative. The coordinator's state model (key-value, small dataset, low write rate) fits etcd's sweet spot. The tradeoff is an additional external dependency in exchange for reusing existing operational knowledge.
When a shard controller fails:
1. Its clusters stop receiving new capacity (but keep running — static stability)
2. The global coordinator detects the failure (missed heartbeat, 30s)
3. The global coordinator restarts the shard controller (or promotes a standby)
4. The new controller loads the shard's machine assignments from the durable store
5. Roll-ups resume. Normal operation resumes.
Clusters are permanently assigned to shards, so there's no redistribution. The shard comes back, picks up where it left off. The durable store (the source of truth for machine-to-shard and cluster-to-shard assignments) survives controller restarts.
If a shard controller is permanently lost (hardware failure, data loss), the global coordinator performs a shard split in reverse: it assigns the orphaned clusters to other shards and redistributes the orphaned machines. This is a rare, manual-intervention-level event — not an automated recovery path.
Every instruction from the global coordinator to a shard carries a fencing token: (coordinator_term, sequence_number). The term is the current Raft term; it increments whenever a new leader is elected. Shards track the highest term they've seen and reject any instruction with a lower term. This prevents a zombie leader from issuing conflicting commands.
Every provisioning action from a shard to a CapacityProvider carries a shard-level fencing token: (shard_id, shard_epoch, sequence_number). The shard epoch increments on every shard controller restart. CapacityProviders reject writes with stale epochs, preventing a zombie shard from making conflicting provisioning decisions.
BigFleet does not maintain a percentage-based warm capacity buffer. At 100M nodes, even 1% would be 1 million idle nodes — economically indefensible.
Instead, provisioning delay is absorbed through three mechanisms that don't require idle capacity:
Bare metal IS the buffer. Idle bare-metal machines in the shard's inventory cost nothing to hold — they're already paid for. The natural slack in a fleet (machines freed by completed jobs, machines between assignments) provides zero-cost absorption. This is one of the core economic advantages of the fixed+elastic model.
Elastic provisioning is fast enough. Cloud instances launch in 30–90 seconds. For most workloads, a sub-minute wait between the CapacityRequest and a Ready node is acceptable. The worker loop runs every 10 seconds; a provisioning request goes out on the next cycle; the node arrives 1–2 cycles later. No buffer needed.
For latency-critical demand, the cluster pre-provisions. If a team knows a training job will need 256 GPUs at 2am, the CapacityRequest can be created ahead of time. BigFleet provisions the nodes before the pods exist. This is a scheduling decision made by the workload owner, not a fleet-level buffer policy.
The only scenario where a buffer genuinely helps is sustained burst demand — many clusters simultaneously needing new capacity faster than cloud APIs can provision. This is handled by the CapacityProvider's internal queuing and batching, and by fanning across multiple cloud accounts, not by pre-provisioning idle machines.
| Component | Per-unit size | Units | Total | Fits in |
|---|---|---|---|---|
| Machine record | ~40B avg | 500K per shard | ~20MB | L3 cache |
| Shard summary | 120B | 200 shards | 24KB | L1 cache |
| Cluster roll-up | ~2KB | 100 per shard | 200KB | L2 cache |
| Global fleet state | 120B × 200 | 1 | 24KB | L1 cache |
| Operation | Volume | Frequency | Throughput |
|---|---|---|---|
| Cluster roll-ups (per shard) | 100 clusters | Every 10s | 10/sec |
| Per-cluster diff | 1 cluster | Per roll-up | <500μs |
| Full shard evaluation | 100 clusters | Every 10s | <50ms |
| Shard → CapacityProvider | ~10 provisioning requests | Per cycle | 1/sec |
| Global coordinator aggregation | 200 shard reports | Every 30s | ~7/sec |
BigFleet scales from a single process to the full two-tier hierarchy as the fleet grows (see Section 16 for the concrete scaling table). The key transitions:
1 → 2 shards is the first coordination event. Before this, BigFleet is a single process with zero distributed systems complexity. After this, you need a global coordinator (which can be a simple leader-elected process) and a durable store for machine assignments.
~20 shards is when multiple cloud accounts per region start to matter. Below this, a single cloud account can keep up with provisioning demand. Above this, the CapacityProvider implementation should fan requests across multiple accounts to increase throughput.
~200 shards is the full scale. The global coordinator manages shard membership, shortfall-driven rebalancing, and fleet-wide preemption. Shard controllers handle the hot path. CapacityProviders manage their own cloud API rate limits and multi-account fan-out.
A set of questions that come up in every serious review: how does a provisioned node actually end up in the right cluster, how does scale-down happen without BigFleet watching pod events, where does BigFleet run, and does it replace existing rightsizing tooling.
BigFleet itself does not join nodes to clusters. It transitions machines through three states (speculative → idle → configured, Section 5), and the cluster binding happens during the idle → configured transition via a callback to the cluster operator.
The callback interface. Each cluster's operator exposes a GenerateBootstrap endpoint that BigFleet can invoke:
service ClusterOperator {
rpc GenerateBootstrap(BootstrapRequest) returns (BootstrapBlob);
}
message BootstrapRequest {
ClusterID cluster = 1;
repeated NodeSelectorRequirement requirements = 2; // labels this node must satisfy
}
message BootstrapBlob {
bytes user_data = 1; // cloud-init, ignition, PXE config, etc.
int32 ttl_seconds = 2; // join token expiry
}
When BigFleet needs to configure an idle machine for a cluster, it calls GenerateBootstrap with the requirements the new node must satisfy. The operator generates a bootstrap blob tailored to the request: kubelet version, feature gates, node labels, a fresh join token, the cluster's CA certificate and apiserver endpoint. The operator's code (or template engine) knows how to produce a kubelet configuration that results in a node with the requested capabilities.
idle machine on shard
│
│ needs to configure for cluster B
▼
BigFleet ──(GenerateBootstrap, requirements)──► Cluster B Operator
│
│ generate blob
▼
BigFleet ◄─────────────(BootstrapBlob)────────── Cluster B Operator
│
│ apply user-data to the machine
▼
kubelet boots, self-labels, joins cluster B → Configured
Why pushback-based rather than pre-registered. An earlier design had operators pre-register bootstrap blobs with BigFleet. That required operators to enumerate every possible node configuration up front — a cross-product of kubelet versions, feature gates, label variants — which doesn't scale. The callback model inverts this: BigFleet describes what it needs, and the operator generates what satisfies it. Operators don't need to predict every future workload's requirements; they just need to respond to requests.
How requirements flow through the system. A user's pod carries nodeAffinity requiring kubernetes.io/kubelet-version >= 1.33 or feature.node.kubernetes.io/dra = "true". The pod goes Pending, the cluster's CapacityRequest controller creates a CR with those requirements, the roll-up carries the CR's requirements up to BigFleet, and BigFleet's GenerateBootstrap call passes them to the operator. The operator generates a blob that produces a node satisfying those requirements. The node self-labels when it boots, and the cluster's scheduler places the Pending pod on it.
BigFleet never calls the Kubernetes API directly. It doesn't issue join tokens. It doesn't know the apiserver's IP. It receives an opaque blob from the operator and passes it as user-data to the CapacityProvider. Cluster bootstrap stays entirely within the operator's purview.
Rebootstrap on cluster change. When BigFleet moves a machine from cluster A to cluster B (drain → idle → bootstrap for B), it calls cluster B's GenerateBootstrap, not A's. The machine gets a fresh blob with cluster B's context. This is how cross-cluster machine transfer works in practice — the physical machine never leaves the fleet, only its cluster binding cycles.
Cluster capability bounds. The cluster operator knows its apiserver version and skew policy. If BigFleet asks for a kubelet version the operator can't support (outside the skew window), the operator returns an error on GenerateBootstrap. BigFleet treats this as an unsatisfiable requirement — the machine stays idle, the CR stays pending, a shortfall is emitted. The user learns through the normal "my pod is Pending" signal that their cluster can't support the feature they asked for.
Bare metal vs. cloud. For cloud providers, the bootstrap blob is user-data (cloud-init or ignition). For bare metal, it's PXE/iPXE configuration or the equivalent. The CapacityProvider.Create call passes the blob to whatever mechanism the provider uses to configure new nodes. From BigFleet's perspective, the blob is opaque bytes — it's the operator's job to produce something the provider can consume.
BigFleet does not watch pod events. It doesn't watch the Kubernetes API at all. Scale-down flows through the same contract as scale-up:
1. A pod completes (job finishes, deployment scales down, user deletes). The pod is deleted.
2. The CapacityRequest owned by that pod is garbage-collected via ownerRef.
3. The cluster's operator computes the next roll-up — the CR is no longer in the list.
4. The new ClusterCapacityNeeds message arrives at the shard with a lower count for that profile.
5. The shard's Phase 3 (reclaim excess) sees the gap: inventory says cluster X has 64 GPU nodes, roll-up says it needs 0. Reclaim all 64.
6. The shard issues a reclaim instruction to the operator. The operator signals graceful node shutdown. The kubelet drains.
7. Machines transition from configured back to idle. For bare metal, they sit in the idle pool for reuse. For cloud, the shard decides whether to hold them idle briefly or call Delete to return them to speculative — based on near-term demand forecast (Section 5).
BigFleet reacts to the roll-up, not to pod events. This is by design — it keeps BigFleet's interface narrow and cluster-agnostic. The operator is the bridge between Kubernetes-native events (pods, PDBs, ownerRefs) and BigFleet's simple roll-up protocol.
If the cluster's operator is down, scale-down doesn't happen. The last roll-up remains the source of truth until a new one arrives. This is safe — worst case, BigFleet holds capacity that the cluster no longer needs. It doesn't cause incorrect behaviour; it causes temporary over-provisioning.
BigFleet's components run in different places with different blast radii:
Global coordinator: One logical service (3 replicas for HA, Raft consensus). Runs in a management cluster or a dedicated cluster separate from the clusters it manages. A bug or outage in the global coordinator means shard membership can't change and fleet-wide preemption can't happen. Running workloads are unaffected (static stability). Blast radius: fleet-wide control plane, but no data plane impact.
Shard controllers: 200 at full scale. Run wherever — same management cluster as the global coordinator, separate shard clusters, or even in the clusters they serve (though that creates circular dependency). A shard controller failure affects only the clusters assigned to it — their existing capacity keeps running, but no new capacity is provisioned until the shard recovers. Blast radius: ~100 clusters.
Global blast scenarios. The worst case is a buggy global coordinator deployment that issues wrong instructions to all shards simultaneously — e.g., reclaiming all machines across the fleet. Three protections: (1) shards validate global coordinator instructions against sanity limits (no single reclamation action can affect more than N% of a shard's capacity), (2) shards require fencing tokens on every global coordinator instruction, (3) the global coordinator deploys to canary cells first — a new version is tested against 5% of shards before rolling out.
Static stability is the key safety property. Even if every BigFleet component fails simultaneously, clusters keep running. No running pod is affected. No scheduling decision is affected. Only new provisioning and reclamation stops.
The hot path is entirely in-memory. Roll-ups arrive, update the needs table, trigger worker cycles. Nothing hits durable storage on the 10-second cycle.
What hits durable storage:
Total durable store write rate at steady state: ~10 writes/second across the fleet. Nothing that stresses etcd, Spanner, CockroachDB, or any reasonable durable store.
Request storms: the worst case is 20,000 clusters simultaneously sending roll-ups — 20,000 concurrent gRPC requests. Distributed across 200 shards, that's 100 concurrent requests per shard, each processing in <50ms. A modern gRPC server handles 10,000+ concurrent RPCs per core. No storm at this scale.
The storm that matters is elastic provisioning demand — 20,000 clusters simultaneously needing new cloud capacity. This is handled by CapacityProviders, not by BigFleet's durable store. Cloud provider API rate limits are the binding constraint (AWS: ~2 RunInstances/sec per account-region, see Section 10). Well-written CapacityProvider implementations fan requests across multiple cloud accounts to multiply throughput.
Operator roll-up load. Each cluster's operator processes its own in-cluster CapacityRequests and sends one roll-up per cycle. At 250K pods per cluster, the operator's in-cluster informer handles watch events at the normal Kubernetes rate — nothing unusual. The protobuf roll-up compresses 250K CRs to ~15 entries. Operator memory is ~500MB for a 250K-CR cluster (informer cache).
Existing rightsizers — Uber's, Salesforce's, or any team's — typically do allocation-based autoscaling: watch utilization, maintain a target buffer, adjust replica counts, trigger the cluster autoscaler. They mix two concerns: deciding what capacity a workload should have (rightsizing), and getting that capacity provisioned (autoscaling).
BigFleet only handles the second concern. It provisions and reclaims nodes in response to CapacityRequests. It has no opinion on what capacity a workload should have.
This means existing rightsizers coexist with BigFleet rather than being replaced. The rightsizer continues to decide "service X should have 200 replicas with 4 cores each." That decision flows to the Kubernetes scheduler, which places pods. When pods can't be scheduled (insufficient cluster capacity), the cluster's per-pod CapacityRequest controller creates CRs. BigFleet provisions nodes. The cluster's scheduler places pods.
During migration from existing autoscaling tooling:
1. Phase 1 — BigFleet runs alongside. Existing rightsizers continue to trigger existing autoscalers. BigFleet runs in shadow mode, consuming roll-ups but not provisioning. Compare what BigFleet would have done against what the existing system did.
2. Phase 2 — Partial migration. Some clusters switch to BigFleet. Others keep existing autoscaling. The rightsizer is unchanged.
3. Phase 3 — Full migration. All clusters use BigFleet for provisioning. Existing rightsizers continue to determine workload sizes. Cluster Autoscaler / Karpenter are retired.
The rightsizer's "maintain a buffer" behaviour translates directly: it continues to configure workloads with a buffer (e.g., HPA target 70% utilization), and the buffer is reflected in the CapacityRequests. BigFleet provisions the capacity needed to satisfy the requests, buffer included.
BigFleet implements the capacity contract defined in Fleet-Scale Kubernetes:
Cluster Operator ──(ClusterCapacityNeeds protobuf)──► BigFleet Shard
│
BigFleet Shard ──(UpcomingNode data)──► Cluster Operator │
│
BigFleet Shard ──(Create/Delete)──► CapacityProvider │
│
BigFleet Shard ──(AvailableCapacity data)──► Cluster Operator
The per-cluster operator bridges the Kubernetes CRD world (CapacityRequest, UpcomingNode, AvailableCapacity) to BigFleet's protobuf interface. BigFleet never touches the Kubernetes API. It doesn't need to. It receives aggregated needs, makes decisions, and signals back through the operator.
This separation means BigFleet can be replaced, upgraded, or scaled independently of the clusters it manages. A cluster that loses contact with BigFleet keeps running. A BigFleet upgrade is invisible to running workloads.
Network topology between clusters and BigFleet. gRPC over the internet? A dedicated management network? VPN? Whatever works for your environment.
Authentication and authorization. How does BigFleet verify that a roll-up message is genuinely from cluster X? mTLS? OIDC? This is important but orthogonal to the autoscaling design.
Multi-tenancy. Can multiple teams share a BigFleet instance with isolated capacity pools? Yes, but the isolation model (quotas, budgets, priority tiers) is a policy layer on top of the core autoscaler.
Observability pipeline. BigFleet should export metrics (provisioning latency, decision throughput, utilization, cost). The specific pipeline (Prometheus, Datadog, custom) is an operational choice.
Cost accounting. BigFleet knows the cost of every machine. Aggregating this into per-team or per-workload cost reports is a reporting concern, not an autoscaling concern.
Shards own capacity. Machines are assigned to shards permanently. No borrowing, no distributed locking on the hot path. Rebalancing machine assignments across shards is a slow, deliberate operation managed by the global coordinator on a minutes-to-hours cycle.
Topology constraints do not cross shard boundaries. If a training job needs 64 GPUs on the same rack, the shard must have a rack with 64 GPUs in its inventory. If it doesn't, the request is unsatisfied — it does not trigger cross-shard coordination. This is why machine assignment uses topology-domain granularity (whole racks, not individual machines) and why the shard count is bounded by resource pool depth. If a specific shard consistently can't satisfy topology requests, it's a signal that machine assignments need rebalancing or the shard count is too high for the available resource pool.
Interruption and reclamation penalties are explicit fields on the CapacityRequest. Not derived from priority. Not inferred by BigFleet. The creator of the CapacityRequest specifies both. Interruption penalty plugs into effective_cost = price + (interruption_probability × interruption_penalty) to drive spot-vs-on-demand decisions. Reclamation penalty feeds into victim selection when BigFleet has to decide which machines to free, capturing operational value tied to specific hardware (GPU burn-in, warmed caches, characterised performance). Priority, interruption tolerance, and machine-specific value are three independent dimensions — a workload that's important (high priority), interruption-tolerant (low interruption penalty), and hardware-specific (high reclamation penalty) is a perfectly valid combination. BigFleet doesn't conflate them.
BigFleet does not manage cloud commitments. Reserved Instances, Savings Plans, and other financial commitments are procurement decisions made by humans or FinOps tools. BigFleet sees their effect through pricing: a provider reports committed capacity at a lower price, and BigFleet naturally prefers it because it's cheaper. BigFleet does not purchase, modify, or recommend commitments. Keeping this out of scope keeps BigFleet simple and keeps financial decisions where they belong.
When the fleet shrinks, homogeneity is preserved. If an operator decommissions 10,000 bare-metal machines in Amsterdam, BigFleet does not remove them all from one shard. The global coordinator removes machines proportionally across all shards — each shard loses its share of Amsterdam machines. The process is:
1. The global coordinator marks machines for decommission, balanced across shards.
2. Each shard drains the marked machines (graceful node shutdown, PDB-respected).
3. Once drained, the machines are returned to the provider for decommission.
4. Each shard's inventory shrinks proportionally. All shards remain homogeneous.
The same principle applies to fleet growth: new machines are distributed across shards at topology-domain granularity, maintaining each shard's ability to satisfy topology-constrained requests.
No quota. No entitlement. No admission control. BigFleet accepts any CapacityRequest from any cluster at any priority with any count. There is no mechanism by which a request is rejected for exceeding an allocation. This is a deliberate departure from the systems BigFleet borrows from: Borg enforces quota at admission and rejects jobs that exceed their allotment. Twine requires workloads to bind to entitlements granted in advance. BigFleet does neither.
Instead, priority is the sole throttling mechanism. You can ask for 10 million nodes at priority 0, and BigFleet will accept the request. You won't get the nodes — priority 1000000 work will occupy every available machine before your request is even considered — but no component tells you "that's too much." The request simply sits in the needs table, unsatisfied, until either the fleet has spare capacity or the cluster stops asking.
This has three consequences worth naming explicitly:
Operators must understand priority, not quota. There's no "you've used 80% of your allocation" metric. The meaningful signal is "my priority-X workloads are consistently unsatisfied" — which means either the fleet is saturated at that priority level or higher, or the request genuinely can't be satisfied (wrong instance type, unavailable zone, topology impossibility). Both require different responses, and BigFleet surfaces both through the shortfall mechanism.
Self-throttling is a consequence of priority, not a feature. A team running priority-0 batch workloads effectively has zero guaranteed capacity — they get whatever the fleet isn't using right now. A team running priority-1000000 production workloads effectively has reservation-grade capacity, bounded only by what the fleet physically contains. The gradient in between is continuous, not bucketed.
Budget, chargeback, and spending limits are policy on top. If an organisation wants to cap how much elastic capacity a team consumes, that lives in a layer above BigFleet — typically in the same place that generates CapacityRequests (Kueue, a custom controller, a platform team's admission webhook). BigFleet doesn't know or care about teams, tenants, or budgets. It provisions the cheapest capacity that satisfies the highest-priority request and leaves everything else for later.
This simplifies the architecture considerably. There's no quota service to run, no allocation database to maintain, no admission API to design. But it shifts responsibility: cost controls that most cluster management systems build in, BigFleet expects to live elsewhere.
Clusters are permanently assigned to shards. When a cluster's operator sends its first ClusterCapacityNeeds message, BigFleet assigns it to a shard at random (or least-loaded). That assignment is permanent. Clusters do not move between shards. BigFleet doesn't need to know about cluster lifecycle — it doesn't know when clusters are created or destroyed. A cluster that stops sending roll-ups simply has its needs table entries age out. A cluster that starts sending roll-ups gets assigned on first contact.
This means BigFleet doesn't interact with cluster lifecycle at all. No registration API. No deregistration. No coordination with fleet management tooling. A cluster exists in BigFleet's world from the moment it asks for capacity, and stops existing when it stops asking.
The shard count scales with the fleet. BigFleet starts as a single process for small fleets and scales out as the fleet grows:
| Fleet size | Clusters | Shards | Machines/shard | Notes |
|---|---|---|---|---|
| 10K nodes | ~2 | 1 | 10K | Single process. No coordination overhead. |
| 100K nodes | ~20 | 1 | 100K | Still single process. |
| 1M nodes | ~200 | 2–5 | 200K–500K | First shard split. Global coordinator is trivial. |
| 10M nodes | ~2,000 | 20 | 500K | Multiple cloud accounts per region for throughput. |
| 100M nodes | ~20,000 | 200 | 500K | Full two-tier hierarchy. |
Shard splits are triggered by the global coordinator when a shard's cluster count or machine count exceeds a threshold. The split is a slow operation: the global coordinator selects half the shard's clusters and half its machines, spins up a new shard controller, transfers the state, and redirects roll-ups. Since clusters are permanently assigned, a split is the one exception — it's the only time a cluster moves, and only because its shard is dividing in two.
Shard merges (when the fleet shrinks) are the reverse: two under-utilised shards combine their clusters and machines into one, and one controller shuts down.
[1] Verma et al. "Large-scale cluster management at Google with Borg." EuroSys '15. https://research.google/pubs/large-scale-cluster-management-at-google-with-borg/
[2] Schwarzkopf et al. "Omega: flexible, scalable schedulers for large compute clusters." EuroSys '13. https://research.google/pubs/omega-flexible-scalable-schedulers-for-large-compute-clusters/
[3] Tang et al. "Twine: A Unified Cluster Management System for Shared Infrastructure." OSDI '20. https://www.usenix.org/conference/osdi20/presentation/tang
[4] Lee et al. "Shard Manager: A Generic Shard Management Framework for Geo-distributed Applications." SOSP '21. https://dl.acm.org/doi/10.1145/3477132.3483546
[5] "Unified Scheduling System First Implemented on a Large Scale While Supporting Alibaba's Businesses During Double 11." https://www.alibabacloud.com/blog/598318
[6] Xiang et al. "Gödel: Unified Large-Scale Resource Management and Scheduling at ByteDance." SoCC '23. https://dl.acm.org/doi/10.1145/3620678.3624663
[7] AWS. "Reducing the Scope of Impact with Cell-Based Architecture." https://docs.aws.amazon.com/wellarchitected/latest/reducing-scope-of-impact-with-cell-based-architecture/
[8] MacCárthaigh. "Workload isolation using shuffle-sharding." Amazon Builders' Library. https://aws.amazon.com/builders-library/workload-isolation-using-shuffle-sharding/
[9] Kubernetes SIG Cluster Lifecycle. "Cluster API." https://cluster-api.sigs.k8s.io/
[10] Karpenter. https://karpenter.sh/