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 ArchitectureOperating 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.
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.