Menu
InfoQ Architecture·September 23, 2026

Modal's Sandbox Infrastructure: Scaling Beyond Kubernetes for 1 Million Concurrent Containers

Modal re-architected their sandbox infrastructure to scale to 1 million concurrent containers and tens of thousands of creations per second, moving beyond Kubernetes' limitations in centralized coordination and strongly consistent state. Their solution emphasizes horizontal scalability and a decentralized scheduling approach where workers act as their own source of truth, optimizing for extreme throughput and low latency cold starts.

Read original on InfoQ Architecture

The Challenge of Extreme Scale with Traditional Orchestration

Operating 1 million concurrent sandboxes and tens of thousands of creations per second presents a unique scaling challenge that traditional container orchestration systems like Kubernetes struggle with. The core issue lies in their reliance on centralized coordination and strongly consistent state management, particularly with components like the scheduler and `etcd` (the central durable store). At this scale, operations that are O(containers) or O(nodes) become bottlenecks, leading to severe issues under high pod creation rates or churn, as `etcd` is not natively shardable within a keyspace and both pods and nodes write to it multiple times.

Modal's engineers adopted a fundamental shift in their platform's philosophy: moving away from global coordination towards a load-balancing-like scheduling model. This design principle dictates that anything causing O(sandboxes) or O(nodes) load must be horizontally scalable by default, and the sandbox creation path must be as simple as possible. This approach directly addresses the limitations observed in Kubernetes by decentralizing critical functions.

Key Architectural Changes

  • Worker as Source of Truth: Each worker node maintains its own state, rather than relying on a central datastore. This eliminates the bottleneck of a single, shared state store for all operations.
  • Parallel Scheduling Fleet: Instead of a single, serialized scheduler, Modal deploys a fleet of scheduling servers that operate in parallel. This allows the scheduling layer to scale horizontally with demand.
  • Direct Worker Communication: Once a scheduler identifies a suitable worker, it communicates directly with that worker via RPC to request sandbox creation. Workers accept if they have free resources, otherwise they reject, simplifying the process and reducing coordination overhead.
  • Redis Stream for Worker State: The only remaining central bottleneck is workers publishing their state to a single Redis stream. However, load testing suggests this is viable for over 100,000 workers, demonstrating its robustness at extreme scale.
ℹ️

Achieving Hyper-Scalability

Modal achieved creating 1 million sandboxes in under a minute, with median startup-to-code time under 0.5 seconds. This performance is a direct result of their decentralized design, which avoids the scaling limitations inherent in traditional, globally coordinated systems.

Kubernetes alternativescontainer orchestrationdistributed schedulingserverlesshigh-scale infrastructuredecentralized systemsAI infrastructurecold start optimization

Comments

Loading comments...