Return of the Borg: How Twitter Rebuilt Google’s Secret Weapon
wired.com
wired.com
Let's say you have a Rails stack (app server + DB) that you want to deploy for testing in 3 different datacenters. If that works out, you then immediately need to deploy 10,000 instances each to 10 different datacenters. Oh yeah, and the storage needs to be distributed and universally available, in case any of the servers crash. Performance testing reveals that you can't have more than 10 servers per rack, otherwise you saturate the rack switch. You also need to account for power distribution redundancy, shifting traffic loads, etc.
Oh, and you want to do this with a single configuration that's manageable by a team of 3-4 people and have deployment be entirely automated and monitored.
The complexities behind this problem are simply enormous, almost too much to even comprehend. I am proud of the work we did at Google to attempt (yes, attempt) to solve this problem, and I know my friends at Twitter are doing great work as well. I'm mostly just happy that I can finally talk about it and that the badasses that do this work can get some credit.
Those people are building truly amazing things now based on that (for lack of a better term) incubator for the wider industry.
The stuff that has been built directly by Google is tremendous, to be sure. But the fundamental understandings that have come out of their R&D to do it will benefit everyone across the entire industry. It's magical to watch.
I have a huge amount of respect for the folks in SRE who would juggle clusters like you and I might adjust thermostats around the office. It was not a job for the feint of heart, and personally I don't think it was appreciated as much by the people who deployed on it as it should have been, but that is always somewhat true about the operations side of the house, you know you're doing ok when nobody is screaming at you.
Man, working at Google is just a different kind of experience in some ways. Not necessarily better, though I like it a lot, but definitely different.
IFAICT, Omega and Mesos (and YARN for what it worths) cannot really handle resource overcommit effectively. I wonder what they'll call it when they reinvent a better DRS :)
A naive approach is to measure idle CPU or RAM not allocated to a process. Then improving utilization metrics is simply a measure of fitting more computation onto one machine until there's not enough CPU time or RAM to fit any more.
This works fine for throughput-oriented workloads, but will cause immediate conflicts when imposed on teams who value low latency (e.g. the Search team at Google). So you end up in a situation where a team is being pressured to improve utilization at the expense of metrics that they care more about, and they refuse, and that's when you get grumbling around the water cooler about "wasteful propellerheads squandering expensive hardware!" / "ignorant beancounters micromanaging things they don't understand!".
I see this as a great example of how innovation happens. You have things that you wouldn't expect to go together working together almost out of sheer chance, and out of that comes something brilliant.
I also love that Hindman joined Twitter as an "Intern" after he was a consultant.
You could even get all fancy by using openvswitch to configure your own virtual network topology. It just doesn't seem so complicated what they're trying to accomplish.
Cheap hardware crashes. It crashes all the time. Furthermore, the application's needs themselves change all the time, depending on traffic curves and processing needs. The dynamic needs of Google's (and presumably Twitter's) completely heterogeneous application stacks don't lend themselves to simple virtualization and over-the-counter software. This is an incredibly tricky bin packing problem that was never quite solved in my 5 years at Google.
I don't know a lot about the "distributed frameworks" you mentioned, so perhaps they do this too. I kind of doubt it. If it were as easy as you think, I'm sure my friends at Twitter would be using what you mentioned.
(Why would it be insane? Because with dedicated machines your utilization levels are going to be crap, reconfiguring and rebalancing the machine allocations will be painful, and you'd need massive over-provisioning to achieve a sufficient level of fault-tolerance.)
Add in spinning jobs up and down quickly on demand rather than making them long-running. (This matters for a map-reduce, for instance.)
Add in making it easy to have dev/production running side by side with configurations that are as constant as possible, and not getting into each other's way.
Add in automatic failover from machine to machine, or from data center to data center as needed.
Add in making it easy to find/access applications set up and configured by someone else so that you can set rpcs their way.
After you keep adding things in, you understand why people think of this as a "data center operating system". And you realize that solving the basic problem - get jobs to run on a distributed set of machines - is only the visible first step in a long list of features that you want.
Another takeout from this, that virtualization is not Web-scale, it's more appropriate to IT and public cloud workloads, large companies like Google prefer bar-metal.
"The trace represents 29 day's worth of cell information from May 2011, on a cluster of about 11k machines"
A sample: over a 7 hour period, just that one cluster executed 3.5 million unique tasks.
Are these like the next generation of task schedulers, or are we talking about a higher abstraction level?