Computing at scale, or, how Google has warped my brain
matt-welsh.blogspot.com
matt-welsh.blogspot.com
MPI is not good at fault tolerance, though MPI-3 will help some. MapReduce and Sawzall are well-suited to the problems they were made for: largely independent operations on very large data sets with occasional reductions and weak synchronization. Usually disk is a significant factor and strong scalability is not especially important.
But they really can't compete with MPI for other problem domains. Claiming that MPI doesn't scale is demonstrably false, as evidenced by PDE solvers running on 200k - 300k cores (Jaguar, Jugene, etc). This is a comfortable order of magnitude larger than Google is running, and the communication requirements of those algorithms are far more demanding, with essential all-to-all communication as well as low-latency global reductions (for which Blue Gene has dedicated hardware). The need for these low-latency all-to-all operations is actually very fundamental (formally proven). The algorithms also tend to have much stricter synchronization requirements, for example, nothing can be done until everyone gets the result of a dot product, because the next operation needs to happen in a consistent subspace (though MPI has fine asynchronous support). There has been plenty of work trying to make algorithms less synchronous, but weakly synchronous algorithms almost uniformly have worse algorithmic scalability (number of iterations to solution increases by more than a constant factor). And since problem sizes are usually chosen to fit in memory, the network doesn't get a free pass due to the disk being slow.
If you care about strong scalability, MapReduce and Sawzall are probably not good choices. Similarly, if your problem domain requires low-latency all-to-all, reductions, or "ghost exchange".
I would assume Google has well over over 20,000 cores. Or did you mean something else?
From what I have read, Google treats its global compute infrastructure in a reasonably abstract fashion. To the point where their production systems share resources across datacenters so that losing a datacenter has minimal impact. Thus, when Google people talk about "scale" they often mean getting it to work for long periods of time, efficiently, on flaky hardware, spread across the world, which PMI does not do.
The most common response is to perform global checkpoints often enough that hardware failure is not too expensive, but infrequently enough that your program performance doesn't suffer too much. Bear in mind that our jobs are not normally hitting disk heavily, so frequent checkpointing could easily cost an order of magnitude. When hardware failure occurs, the entire global job is killed and it is restarted from the last checkpoint. There are libraries (e.g. BLCR) that integrate with MPI implementations and resource managers for checkpoint, restart, and migration.
Note that it is not possible to locally reconstruct the state of the process that failed because there just isn't enough information unless you store all incoming messages, which can easily be a continuous gigabyte per second per node. Even if you had a fantastic storage system that gave you ring storage with bandwidth similar to the network, local reconstruction wouldn't do much good because the rest of the simulation could not continue until the replacement node had caught up, so the whole machine would be idle until during this time. If you have another job waiting to run, you could at least get something done, but the resource management implications of such a strategy are not a clear win.
So if hardware failure becomes so frequent that global checkpointing is too costly, the only practical recourse is to compute everything redundantly. Nobody in scientific computing is willing to pay double to avoid occasionally having to restart from the last checkpoint, so this strategy has not taken off (though there are prototype implementations).
First off individual computational errors are rarely important. EX: Simulating galactic evolution as long all the velocity's stay reasonable each individual calculation is fairly unimportant and bounds checking is fairly inexpensive.
Second, there is a minimal time constraint, losing 30 minutes of simulation time day is a reasonable sacrifice for gaining efficiency in other areas.
Third, computational resources tend to be expensive, local, and fixed. AKA, Cray Jaguar not all those spare CPU cycles running folding at home.
However, if you’re running VISA or World of Warcraft then you get a different set of optimizations.
This is completely wrong. Arithmetic errors (or memory errors) tend not to just mess up insignificant bits. If it occurs on integer data, then your program will probably seg-fault (because integers are usually indices into arrays) and if it occurs in floating point data, you are likely to either produce a tiny value (perhaps making a system singular) or a very large one. If you are solving an aerodynamic problem and compute a pressure of 10^80, then you have might as well have a supernova on the leading edge of your wing. And flipping a minor bit in a structural dynamics simulation could easily be the difference between the building standing and falling.
I would argue that data mining is actually more tolerant of such undetected errors because they are more likely to remain local and may stand out as obviously erroneous. People are unlikely to die as the result of an arithmetic error in data mining.
Second, there is a minimal time constraint,
There is not usually a real-time requirement, though there are exceptions, e.g. http://spie.org/x30406.xml, or search for "real-time PDE-constrained optimization". But by and large, we are willing to wait for restart from a checkpoint rather than sacrifice a huge amount of performance to get continual uptime. If you need guaranteed uptime, then there is no choice but to run everything redundantly, and that still won't allow you to handle a nontrivial network partition gracefully. (It's not a database, but there is something like the CAP Theorem here.)
Anyway, if you could not do this and accuracy is important, then you really would need to double check every calculation because there is no other way to tell if you had made a mistake.
"The cloud" has been around for a while, although implementation has changed.
Search team has a mixture of windows laptops, macs, linux desktops, linux laptops..and it didn't matter because everyone just ssh'd to a RHEL box and worked.
Yahoo used to be a BSD shop, then moved to linux. Search team used both Debian and RHEL in the beginning, then transitioned to RHEL fulltime.
I don't understand the obsession with RHEL. In my experience, Debian/Ubuntu will give the same benefits as RHEL sans the cost.
The rationale for moving from BSD to linux was practical. The gist was most of the vendors either do not do bsd or bsd was second class citizen. And BSD wasn't as quick as linux in catching up with esoteric hardware. But discarding Debian didn't actually make sense to me.
Actually I don't understand why Google would use virtualisation. Why would they split a physical box up? Making many hosts look like one, yes (is such a thing possible?), but not the other way around.
And AFS was cool: it was actually a global filesystem that many institutions used. you could cd /, ls and browse to nasa.gov amongst many thousands of others (much was protected!) (cern was mounted on /cern.ch)
Ease of administration. It turns out to be much easier to have ten virtual machines running one daemon each than one virtual machine running ten daemons. Fewer interactions, cleaner process boundaries (the boundary between one VM and another is pretty thick), easier debugging. Easier to swap out parts without affecting other parts. Much easier scaling. (If you need more servers for Service X, launch more VMs; if your machine needs more RAM to hold those VMs, launch the VMs on a different physical box -- the system doesn't need to be changed to support that because it's already architected to assume separate "machines" for everything.)
For more on the subject ask Ezra Zygmuntowicz or his blog and books. I stole all of this material from him.
It's funny, when you needed one OS instance per network server daemon on Windows in the 90s, everyone laughed at it. Now it counts as cutting edge thought.
Now it's a matter of choice for system owner. And of course all those single-daemon VMs run inside OS that lets numerous VMs run at the same time.
Not to mention Solaris 10 can be...interesting.
Yahoo still uses BSD. They're not completely BSD, of course; but they're not completely linux either.
Some new teams use BSD for their webapp deployment and database servers; can't confirm (no links) but I think the teams were of the opinion MySQL on BSD was more stable.
But as said, earlier Yahoo was big time BSD, investing in it and submitting patches, which no longer is the case. Search team did a mass migration to RHEL about 2 years back. BSD is still used here and there as deployment platform. My personal preference would be to go with BSD and save some buck but apparently management views it differently.
This situation sounds fine if you're a developer and your job requires only low-bandwidth text streams between you and your physical computer. But for anyone who does anything graphical, or for anyone at home who wants to enjoy media, the bandwidth and latency are big problems.
Maybe I should just assume that because this is on HN and written on a developer's blog, context is assumed. But it sure seems like a disconnect between hard programmers and the other 99% of the population.
1985: Interactive debuggers suck. PRINT() is your friend.
1990: Interactive debuggers have matured. No one uses PRINT() any more. Debugger questions are now even part of technical interviews.
2008: Concurrent processing and cloud computing have made interactive debugging difficult.
2010: printf() is your friend again.
Developing software is getting to be like fashion. Keep those old skinny jeans and workarounds in your closet. Sooner or later, they'll be in style again.
People in the Ruby community also use temporary puts/logging statements for exploratory debugging, but of course nobody in right mind commits such code to the repository.
[1] http://googletesting.blogspot.com/
[2] http://googletesting.blogspot.com/2007/01/introducing-testin...
* Strong, stern type checking has always been kind of a pain in the ass, but it's very effective at catching miscellaneous stupid errors. If debugging in the cloud is hard enough, then maybe overbearing type systems are the lesser of two hassles.
* Interactive programming is a big win. One of the best things about Lisp is the REPL, and how well-integrated it is with the editor. (At least, if you use something like SLIME.) If I'm programming on a distributed system, I definitely want to be able to quickly test out my code, preferably in a concurrent environment. You could solve this by having a bunch of servers for this, or with some fancy simulation environment, or a bunch of virtual machines, or something. The point is, I want errors to show up fast.
* Really good log analysis. Ideally, any framework you're using would automatically log everything it does, and then you would use some nice tools to figure out what happened, and where it went wrong. Maybe the Loggly guys will come up with some good stuff.
* Data-flow oriented abstractions like what Apache Pig offers. Pig lets you define a directed acyclic dataflow graph, in which each edge is some operation like "group by field x", or "filter by function f(fields)". It compiles all this into MapReduce jobs and runs them on Hadoop. It's probably a lot harder to mess up a query in Pig than with raw MapReduce, for sufficiently complicated queries. I'm not sure how much of a performance hit you take, though.
Or is developing software like business? Sell those skinny jeans and workarounds when they are in demand, then buy some back when they go out of style. You're now ready for their comeback and you have some mad money in your skinny jeans' pocket.
However, my use of debuggers has always been quite minimal, perhaps because despite the CSD experience, I am often behind the times. Or perhaps it was due to my early (not quite childhood) experience building real-time interrupt-rich programs in Sigma 5 assembly language.
I also read recently that most of the grownups don't use debuggers but have always stuck by printf or equivalent.
It kind of reminds me of an early attempt on my part, not wanting to do laundry in my early bachelor years, to start a wrinkled shirt fashion trend. I was not successful, at least for a large number of years. It finally did become fashion.
So stick with printf and build useful but lean logging habits would be my advice.
One is that remote machines can be rented instead of being bought.
Another is that the configuration of everything becomes easier. You buy new laptops by taking them out of the box, running a handful of app installers, and moving over a SSH key or two and an encrypted password file. Your career as amateur sysadmin might just be over. (If you miss the chance to build your own Linux kernel, of course, there's plenty of work left for you as a cloud sysadmin. But cloud system admin is also usually easier, because machines in data centers come with nifty administration tools, like SANs and dashboards, that benefit from economies of scale.)
Security: When you lose the "dumb terminal" you lose no data. You do lose the keys that are on the machine, but you can revoke those. And it is much, much easier to get consistent automated backups on a datacenter box than on a laptop.
There are definite problems with the "cloud" approach, such as "who owns the data and how do you trust it is being handled properly", but the advantages are strong these days.
I love moments like this when you realize as a developer that the bits that are the fruit of your labor are off executing in some far off physical location you'll never know about.