Why the data center needs an operating system
radar.oreilly.com
radar.oreilly.com
Google SREs by last count were 1 engineer to 1000 machines. In 10 years the common devops engineer at a 200 person startup will leverage the same number of resources using layers of abstraction like Mesos and Kubernetes.
That number does not seem particularly impressive, if it is accurate. Even "traditional" well-run enterprise IT organizations are often in the 1 admin/SRE to 600-ish machines, so I have a hard time seeing that Google can only do ~2x as good at their scale and with their level of focus.
1 SRE to 5k machines, 10k machines, that makes more sense to me.
you can take a snapshot and say you are 1:n because today you have so many sre and so many machines, but it is very unlikely to be the same ratio down the road.
Admittedly, I probably have a bad impression of enterprise from consulting because who would pay for automation consulting at $$$ / hr when you do it pretty well with existing resources in the first place?
It seems much easier in my opinion to scale a single function (search or tweet) than the kinds of tasks that a healthcare company has to do like say...scanning faxes from doctors, applying OCR and properly placing them into a pharmacy order system.
When you remove the VM smokescreen and count physical boxes it's more like 1 person/100 machines, which is abysmal. I've seen order-of-magnitude people efficiency increases with automation like we're discussing here.
I don't know why many people think that they need to be at datacenter scale computing to benefit from abstractions like Mesos, it's completely wrong imo.
It's quite a big shift in mindset but it makes the life of everyone (dev and ops) so much easier when you stop having to think about single machines.
My next startup will be built on Mesosphere. Faster time to MVP and no "go dark for 18 months" when I have to scale.
Firstly, "distributed computing is the norm"? It's just not. Most businesses & app authors will never need to care about ultra-distributed computing, with all its problems and trade-offs. You can move faster with "local-only" computing and scale vertically very cheaply compared to a few years ago - 4 dedicated CPUs + hundreds of gigabytes of RAM save programmer hours, and get your problem solved faster.
For light scaling issues (compared to Google) Redis, MariaDB and other abstractions over local files have some great options for future scaling, and are well-trodden, obvious choices.
Secondly, who cares about "wasted" resources of a whole underutilised server when reliable dedicated servers are so cheap, and in such plentiful supply?
Thirdly, "organizations must employ armies of people to manually configure and maintain each individual application on each individual machine"? - in the 90s maybe! Surely anyone with more than a few applications to worry about is on board with some basic configuration management.
Twitter-size scaling is a "nice problem to have". For all but the best-funded & bullish companies, solve them only when you start to have them.
(my bias: I run a managed service provider in the UK - we tend to help customers scale vertically by shovelling server images around with minimum down time. We say "underused" dedicated server capacity at fixed monthly costs is usually cheaper than chasing the phantom of "optimum" AWS usage.)
The people bankrolling Google / Facebook / Twitter's electricity bills seem to care quite a bit.
There is another trope that gets repeated often, (and this is not even remotely directed at you, just a digression hopefully somewhat on topic) "performant languages runtimes are an anachronism, a bog slow language in which a programmer can code fast is way more useful than any of the performance bull crap". Typically the person repeating that would a be a webdev. However, in these large scale scenarios core infrastructural code can save orders of magnitude more in money in running costs than saving days in software development. So yeah at the interesting places algorithms and efficiency continue to matter. A reason that Google always managed to be ahead is partly due to how successful it was in minimizing running costs.
We do a lot of finding ENORMOUS performance problems with customer servers - simple stuff like a thundering herd, a vital missing index or a filesystem that's being overtaxed That's the kind of scale we work at. But those sorts of insights can make the difference between "help we might need a new server" and "oh thank god it's all working again".
/I used to work for a company that had a big Java app. We laughed at our client who needed 60 Rails servers to deliver worse performance than our single-instance app. But they probably saved more on dev costs than they spent on servers.
If we had a fabric (which is called Mesos ;) which allowed you to write elastic, distributed systems without the need for worrying about interconnecting hosts and segmenting hosts into static partitions, etc. wouldn't that be a big win?
Spark is also a great example for an app that was built directly on top of Mesos - the authors could focus on implementing the actual logic rather than worrying about interconnection. Same is true by the way for systems like Chronos and Marathon - all of which run on top of Mesos.
So i dont see distributed computing become the norm either. At least not in the next few years.
It's VAXocentrism for the 21st century. http://www.catb.org/jargon/html/V/vaxocentrism.html
- Pointer types are basically fungible in C (otherwise there would be no void-ptr type)
- There's no compile-time knowledge of the segment a pointer references to prevent you from dereferencing a pointer to an offset from segment A when segment B is loaded (compare this to Rust's parameterization of Box types by their allocator)
- Struct padding is painful and tacked on
- "unsigned char" isn't default even though it'd make much more sense for it to be (What "char"acter is negative? You can have a signed byte/octet, but a character is—in 1979, at least—basically an enum/sum type.
UNIX API fragmentation is still a problem, but my point is that POSIX was trying to unify existing implementations rather than a greenfields project to generate a set of portable APIs.
Mesos paper: http://people.csail.mit.edu/matei/papers/2011/nsdi_mesos.pdf
> Exposing machines as the abstraction to developers unnecessarily complicates the engineering, causing developers to build software constrained by machine-specific characteristics, like IP addresses and local storage.
It reminded me of the old RPC approach of making potentially any method call a remote method call. It didn't work out, simply because remote call latencies are orders of magnitudes higher than for in-process calls, so you need different (coarser) APIs for them.
By the same token, while abstractions are very welcome, they should still allow distinctions between local and non-local resources, for various definitions of "local".
It seems like what might be ideal would be a specification language that doesn't care where things are, an implementation that tries to deal with that automagically, and a way to specify portions (to all) precisely that is checked against the high-level specification.
Surely you meant 10's of milliseconds. But even so, fastest SSD random access latencies are in sub-millisecond ranges.
With the direction things are moving, "the data center as a computer" is absolutely the right approach.
Also, locality and latency can both be expressed in terms of placement rules and schedulers on Mesos can use those rules to guarantee or express preference for task placement that optimizes around reduced latencies.
The point that I take away is that the ability to express your needs in a declarative way (e.g., "place these two tasks such that they have such-and-such latency) is much more scalable, flexible and resilient than coding to machine-specific internals. The latter is easier to update and supports delegation of responsibilities.
John Wilkes of Google puts it this way:
"Our own experience has been that allowing our developers unfettered access to the internals of infrastructure systems has been a problem, and we're moving away from that model as fast as we can.
Constructing large-scale complex systems with many interdependencies leads to brittle, fragile systems if they rely on internal implementation mechanisms.
Allowing internal customers to rely on internal implementation mechanisms has made it hard to adopt new technologies, because we only know what knobs they set - not why.
The fix for both is similar: describe the desired end state, not how to get there."
That makes total sense to me, thanks.
There is no such thing as a datacenter OS.
Simplification: An OS is a kernel and it's associated base software.
The kernel drives the hardware. You need to talk to disks. Memory. what not. Kernel is needed. You need to talk to the kernel and tie these components together. You write software. Boom, you have an OS.
A datacenter isnt a disk and memory and devices. its a bunch of computers which themselves drive these devices.
If you invented a "datacenter" with a bunch of disks and an API to drive them, then a bunch of CPUs and an API to drive them, and so on, you'll end up with a supercomputer and a single, unreliable, unsafe OS (which is exactly why nobody makes super computers anymore. They make clutsters of computers. Cheaper, more reliable. Clusters. Ie... datacenters).
The only thing that could be needed is a universal API for resource access. Need a db? Here's an API. Need disk space? Here's an API. and so on.
This is exactly what AWS is and does. S3 doesnt expose the OS. Its just a filesystem API. ELBs arent an OS. They're a load balancer API.
It turns out that below that, there's an actual traditional OS because thats the way it works reliably.
You have other ways to interconnect these systems in smarter ways (see plan9) but its always running a "regular" OS in the end, too.
[1] http://www.cs.berkeley.edu/~rxin/db-papers/WarehouseScaleCom...
E.g. for an OS we have:
filesystem
scheduler
cron
Then you can go look at the history of these primitives so the same lessons don't have to be relearned. For example, the Linux kernel has gone through many iterations of its scheduler with cgroups, cfs etc. Why did they do that? Why were previous incarnations not good enough? etc.
Interested in learning more about or contributing to Mesos? Check out mesos.apache.org and follow @ApacheMesos on Twitter. We’re a growing community with users at companies like Twitter, Airbnb, Hubspot, OpenTable, eBay/Paypal, Netflix, Groupon, and more.
"The early catch phrase was to build a UNIX out of a lot of little systems, not a system out of a lot of little UNIXes."
What I am really keen to find out in the coming years is what MirageOS makes of this. If you are not familiar this article http://queue.acm.org/detail.cfm?id=2566628 explains it way better than I could. I wouldnt claim it is there yet but seems to be sitting right at an envious position full of realizable potential. Its written in OCaml to boot.
It's the only (semi-mainstream?) language that I know of that includes the infrastructure in the language itself to make the individual underlying machines (OS / hardware) appear irrelevant and allow programming to seamlessly span a group of computers.
Mesos is basically an application scheduler. It doesn't manage the base operating systems or machine provisioning. Mesos is concerned with ensuring that one or multiple applications are launched and running on a cluster of machines.
The Saltstack framework is the only thing I know of that would provide all the primitives and control necessary for an operator to truly control a full 5-300,000 node datacenter from one station. It's the only thing out there that will allow for super low latency response to commands across the cluster. It also can do configuration management etc if you want, or a lot of people use it to just trigger existing chef/puppet jobs.
I've been using Saltstack in conjunction with Mesos to build out this full datacenter-as-OS stack. Works... :-)
FWIW, Ansible also does low-latency communications through pipelining/accelerated mode. Ditto for the config management and triggering.
If I was to compare the two, I'd call YARN an application scheduler and Mesos a more meta scheduler. You could build yarn ontop of Mesos. You likely couldn't easily build Mesos ontop of Yarn. It should be relatively trivial to make a YARN Framework, which someone has already done with Mesos:
Mesos is closer to the metal but still is only a piece of what is needed. Picturing it as an OS is misleading while picturing Hadoop V2 as a proto-OS is not that much in accurate.
Twitter runs mesos in a multi 10k node cluster for user facing, production applications.
On another note; next up: "data center_s_ need an operating system". Very interested to see the move towards multi DC/availability zones.
Any singular system is going to fail.
the performance penalties you enjoy if you go down the shared disk path (which is still going to be a fun failure when it does fail)
I'm no Erlang expert, but it seems to me all the fundamental bricks are here (and were for multiple decades!) to build this "datacenter platform".
https://www.kernel.org/doc/ols/2005/ols2005v2-pages-243-258....
If you want containers use kubernetes by google!
Was in the process of putting cassandra on mesos.
If you want to mess around with it, you just need a digitalocean account and head over to mesosphere they'll set up 5, 7, 10 clusters for you.
I think the point of the article is to have one operating system, running multiple physical boxes. You could then in turn have something like containers running on top of that OS.
Here is a short tutorial for standing up Mesos on a single CoreOS instance: https://mesosphere.com/docs/tutorials/mesosphere-on-a-single...