Large-scale cluster management at Google with Borg
blog.acolyer.org
blog.acolyer.org
Exciting times!
[1] http://www.wired.com/2013/03/google-borg-twitter-mesos/
[2] https://www.youtube.com/watch?v=0ZFMlO98Jkc
- websearch
- google brain
- video transcoding for youtube
- appengine (think snapchat)
- gmail
Almost everything Google does runs on Borg.
So whatever else was going on, protein folding and drug docking and telescope simulation was soaking up whatever leftover cycles there were.
One of two possible designs immediately come to mind: (1) The MapReduce master interacts with Borg directly on behalf of the programmer and places a requirement on the task to run on or near any node that has a replica of the data required or (2) for each Borglet, there is an accompanying MapReduce agent running to accept RPCs from the MapReduce master.
(Furthermore, building out separate networks is extremely cost-prohibitive once you've expanded out to a total footprint in the low-1e6 ballpark.)
Think about it this way: In recent history, we've added L3 cache to commodity processors because keeping data on-chip is still faster, in 2015, than going off-chip to RAM. The same analog holds when trying to fetch from local disk versus remote disk.
Brendan Gregg actually provides a really good time-stretched view of various operations in his book on performance tuning. Go give that book a read; there's some _really_ good stuff there.
No, seriously.
Picture if all of your data is in RAM, somewhere. If you can move your program to that machine quickly, then that's a huge win.
For one, if I'm running a MapReduce, it's either likely on the same data someone else is investigating... Or I'm likely to have to re-run my MapReduce a few times, calculating exactly what I'm looking for.
So, I may be cheating by saying "Only the first time is the data cold. Every other time, the data is likely to be hot in RAM."
Because you'd be right to say, "Well, duh. I'm talking about the FIRST time. And since I'm talking about the first time, moving the data across the network wouldn't be that bad. After that, yes, of course it will still likely be hot, just as you describe."
My only real response is, I think it depends on how much data, versus how much code you're talking about. And can you maybe even broadcast the code to the right nodes, to even more effectively use the network.
"in which case you don't have a mapreduce, you have a coprocessor."
...a coprocessor which uses the MapReduce API. Yes, I think that is exactly what you have. And it works whether the data is local or remote, cold or hot. One unified API for all of those things, and it's probably optimized to work local and hot, because I'm betting more than half of all MapReduces actually work out to be local and hot.
Mapreduce workers can request shards (work) that have data local to the host it is running on. For large mapreduces this basically works because both the data and the tasks are in most of the cluster, and for small mapreduces it doesn't matter.
From the paper: "We use a Linux chroot jail as the primary security isolation mechanism between multiple tasks on the same machine… all Borg tasks run inside a Linux cgroup-based resource container."
I think was in RHEL 6.2 possibly even 6.0
You may run O(10) containers on a single machine. When you run O(10k)+ containers in a cluster you need new tools to manage things.
Things can in fact work isolated without chroot, too. Its just convenient.