Props to game servers- they atlre more complex than much of eterprise software.
However this is a big distinction - depending on a game, 10 - 100 game istances could run on one physical machine
You also need to manage the complexity of packing those game instances onto the same machine for your game which is a whole other fish.
My current project is doing 5x the traffic but with 20x the cores, it’s embarrassing. And that’s just counting “our” servers, which are only about 1/3 of the whole enterprise (heading toward 50% if I wasn’t on the scene). I look at all the waste in our project and then I think about why I would ever need 50k cores, let alone servers and I just can’t fathom it. Who is handling a tens of millions of requests per second? And what on earth are you fucking up so badly that you need 2 million servers? Is Google doing 4 billion requests per second?
At some point I have to ask myself if making it easy to manage more servers was really their best strategy. Often friction and constraints are where innovation comes from. When things are easy people don’t think about them until they are gone.
Speaking as someone less than one order of magnitude below that, it takes us about ~100-200 cores. (Depending on how you amortize shared infrastructure like our Kafka brokers, etc.) So even 10x'ing our infrastructure I can't imagine 50k servers.
Then that moment of dawning horror when you see that nobody knows what the fuck they're doing and everyone is faking it at best, and just a child playing dress-up at worst.
Google has certainly figured out a number of things, but anyone using 2.5 million servers in 2016 on a planet of only 7 billion people is playing dress-up.
Most servers do far, far more compute intensive things than handling connections - that’s a pretty meaningless number. I work in the mobility space, we’ve got servers that solve large/complex vehicle routing problems, and ideally they’re performing computation for just ONE user at a time.
Certainly 50,000 servers is a lot, but tonnes of large tech companies run 100s to 1000s of servers at a time.
All things being equal, CPU-bound is the exception, not the rule. Most every program we think of as an "application" is chat with some structure, persistence, access control added, and are indeed IO bound.
memory bandwidth bound
memory size bound
cache hit rate and bandwidth bound
TLB size bound
CPU decode and issue logic bound
CPU renaming / OOO buffer bound
kernel implementation bound (esp. locks, interrupts,...)
network physical layer bound
GPU/TPU bound (which are variations of the above, though mostly memory)
power and cooling bound
...
... so basically unless you're running out of functional units, you're hopefully I/O bound depending on where you consider i/o.
But in reality, most SW is just "crappy SW making poor use of resources bound." That's where we mostly are now as an industry. Bad language choices, terrible design, no cross-layer comprehension.
I agree, but the number one culprit is premature distribution, thanks to the widely pervasive cargo culting Amazon-style microservices.
Seems possible that horizontal scaling at every layer can be wasteful, but the alternatives are hard for me to even conceive.
The overhead that concerns me more is messaging. A network of N nodes has N! paths. In the general case it doesn't take long before the overhead of internal messaging absolutely dominates all other processing in the steady-state. Most people take a brute force approach of partitioning the network on purpose, or even intentionally bottle-necking (again) all the traffic. What's really hilarious is when you see people architecting microservices with kafka and apogee with all the uservice trimmings, only to deploy everything to a single physical rack, or even a single beefy machine.
I feel like the proper path to scalability goes through repeated breakage, because then you get to see how things actually break under load, and what can be done to fix it. For example, you can do a lot by moving logic to stored procs, or equivalently, moving your db into your application process. But people seem to think these are non-starters, for some reason, probably because the notion of stateless app servers as the key to horizontal scalability has become a Law, even though there are alternatives. But good luck questioning foundational assumptions - the risk is just too damn high.
"What is are the natural, performant operations for the layer below me? How can I construct my solution from these operations instead of pretending the layers are orthogonal?"
A good example would be the linear read/write performance of harddisks. If you had the option to avoid random access and take advantage of read-ahead and other methods you'd have seen much better performance than an approach that ignored this behavior.
There are many, many examples of this.
I see some clouds deploying servers with 100 Gbps NICs and I wonder what percentage of deployed applications could get anywhere near that…
Not all applications are web pages where the lifecycle of a given connection is "establish TCP, establish TLS, receive series of requests and produce series of responses mostly by hitting external caches or DBs" [or the QUIC variant of this]. That problem space was one of the very first scale challenges to arrive in 1998 and was one of the very first that actually got addressed. It's not the challenge space now and has not been, basically, for 20+ years, except perhaps the whole scaling-of-database problem, which has been dealt with by sharding, and the distribution of flows problem, which has mostly been solved with clever applications of intensely performant _scale-up_ hardware in the form of modern switch NPUs running variations of ECMP and clever uses of anycast, DNS load balancing, routing, etc.
All of this stuff was pretty common by the mid-2000s.
But bottlenecks in systems always exist. They move around. They can be anywhere in the stack.
Connections is the simplest one that got a lot of attention for twenty years - real pre-emptive threads (in Linux, and Solaris's weird diversion into m:n), select() scalability (both in terms of bookkeeping and in terms of basic limits - that is, the lack thereof) giving rise to kqueue on freebsd, WFMO() on NT and years of attempts on Linux to get something that actually worked, which took awhile, _then_ the c10k problem, and so on.
After connectivity you have issues in layer 3 - TCP offload, cost of TLS, etc. Userland to kernel copies (basic stuff like sendfile() to different userland driver schemes). And so on. Physical layer - servers move from NICs with lots of copies, to ring based with scatter-gather, to TSO, crypto offload, to ... while going from 10mb to 100, 1000, 2.5g, 10g, 40g, 100G and sooner or later 400G on server will not be as uncommon as it is now.
But as your networking capacity and throughput scale, you start bumping into other things. Once you're talking high speed links and various schemes to get the kernel out of the way, you are mostly - not always, but mostly - talking about data movement problems. Elephant flows have their own system level problems in networks, and for data moving and staging you actually don't want the host _cpu_ involved if you can avoid it, let alone the kernel. Now you are in the area of doing (remote)->NIC----PCIe---->NVME (or --->GPU) directly, if you can. Now your NVME storage device becomes the bottleneck.
90s era supercomputing clusters had all of these problems with slightly different technologies, AI clusters have them today. These are not connection limited, they do not scale with people. Their primary scale challenge is utilization/CAPEX, but that's a longer discussion.
Can you elaborate on that part?
Realistically you also usually need to perform some non-trivial work from time to time for some non-trivial portion of those connections, which will further load your server, but still.
[0] https://phoenixframework.org/blog/the-road-to-2-million-webs...
The more likely load balancer outcome would be DNS split on inbound client IPs, and scaling out until each load balancer handles the appropriate amount of traffic (by some measure and scale out if exceeded).