Inside Google Spanner, the Largest Single Database on Earth
wired.com
wired.com
I mean, even if you had a picosecond accurate clock available for use inside a server farm, you would still need a way to query it with a known (not necessarily zero; just known) latency to synchronize several machines. Servers are not known latency machines (unless specialized hardware is involved).
How is that accomplished?
And what happens when two transactions happen below the system accuracy limit? (like two transactions pertaining to the same data, 20ns apart, in different servers; impossible to order).
Surely they have solved this, I just wonder how.
From the Spanner paper: "TrueTime explicitly represents time as a TTinterval, which is an interval with bounded time uncertainty (unlike standard time interfaces that give clients no notion of uncertainty)"
By representing time as intervals you can tell if one interval is definitely before or after another interval. If the intervals overlap however, then it means there's some uncertainty. I haven't quite got through the paper far enough to undersand how this kind of thing is handled. :-)
But I'm guessing that this kind of design means they're always consciously designing with error in mind. There's probably some acceptable amount of error in timing that they're able to carefully manage.
Edit: link to the paper for the lazy [pdf] http://static.googleusercontent.com/external_content/untrust...
Edit 2: It looks like the error typically ranges from 0 to 7ms averaging around 4ms. Outages can cause spikes in this error margin.
I think it's just a case of traditional NTP (which is itself based on atomic clocks and/or GPS) done over 'the Internet' is subject to too much latency that it becomes unreliable for their needs (e.g. Spanner).
By putting the equivalent of their own NTP master servers (based on GPS and atomic clocks) in each of their major data centres they solve that part of the problem. The possible LAN latency is much more controllable and reliable, to within tolerances that makes Spanner workable.
The clever bit isn't putting their own NTP master servers in each data centre, it's how Spanner works when using this info. The article puts far too much emphasis on the former whilst glossing over the latter.
I'm looking very forward to the paper on that protocol/system which will hopefully be forthcoming - They tease it with this: "This section describes the TrueTime API and sketches its implementation. We leave most of the details for another paper: our goal is to demonstrate the power of having such an API."
If the Spanner paper is as important as BigTable, ACID may become the new goal for those building distributed systems.
Full disclosure: I'm with FoundationDB, which is a distributed NoSQL database with high performance cross-node ACID transactions. http://www.foundationdb.com
They don't say they have an ACID NoSQL database, which is, to me, an oxymoron: ACIDity is useful if you have a powerful query language.
In the end don't you fear that you might simply reinvent SQL? Or, am I missing something?
For example, you could read the values associated with keys "a" and "b" and "c" and based on that information, make a change to the values associated with keys "c" and "d" and then commit the transaction. The transaction then either fails completely (if one of the keys you read had been written to since your read, in which case you re-try the transaction) or succeeds completely. That is an ACID transaction in a (NoSQL) key-value store.
It's great to have transactions when you have to cross reference data, on a key/value store, although you can have range queries I submit that kind of operations is less frequent.
But still, it sure is great to be able to run a batch of operations with the confidence it will be transactional, I'm sure there are many use cases that can benefit from it.
Oddly enough, it's strange that there aren't more NoSQL engines offering this feature as once you have MVCC you've done the hard part and AFAIK several NoSQL db have MVCC.
You can build all manner of higher level data structures with a transactional ordered k/v store, which is why our concept of "layers" is possible: www.foundationdb.com/#layers
* Consistency (all nodes see the same data at the same time)
* Availability (a guarantee that every request receives a response about whether it was successful or failed)
* Partition tolerance (the system continues to operate despite arbitrary message loss or failure of part of the system)
According to the theorem, a distributed system can satisfy any two of these guarantees at the same time, but not all three.
In practice, a network partition bad enough to bring down something like Spanner is pretty rare. (I worked at Google for four years, and I can't recall such a partition occurring.) But in theory it's certainly possible.
It doesn't say much more than the obvious. Obviously, if one node parts, you're either inconsistent or available.
I think it's a large misconception that these are competing ideas. They're not. They just represent different use cases: search results don't need to be consistent within seconds of events. Customer information, however, does.
Per the paper: "At least 300 applications within Google use Megastore (despite its relatively low performance) because its data model is simpler to manage than Bigtable’s, and because of its support for synchronous replication across datacenters."
So, the current trade off is ACID vs. performance, but I think the interesting point is that ACID is winning more and more now that a distributed ACID database is available to their developers.
Again, according to Google: "We believe it is better to have application programmers deal with performance problems due to overuse of transactions as bottlenecks arise, rather than always coding around the lack of transactions."
In any case, I don't think Google will stop using BigTable any time soon, especially with Spanner's massive latencies.
Which isn't at all comparable with Spanner, since Spanner works between clusters/datacenters, unlike FoundationDB.
"Distributed" is getting severely overloaded when it comes to databases. We need to stop calling databases with sub-millisecond latency to other nodes "distributed". They're not distributed, they're clustered, and the techniques you can use with them are vastly different.
Consistent transactions with a clustered database are a solved problem, full stop, and available to every developer with virtually any database backend, simply by building your transactions on top of Apache Zookeeper, and running whatever durable database you want underneath.
Doing consistent distributed transactions at a global scale, which is what Spanner does, is so far novel. Doing ACID transactions in a 24-node cluster (what is described at FoundationDB's website, for instance) isn't impressive these days, and if that's all Google had done, you'd see a collective yawn around the web.
"As Fikes points out, Google had to install GPS antennas on the roofs of its data centers and connect them to the hardware below."
This is usually one of the first things an enterprising sysadmin does at companies when they first start thinking about time - drop a GPS receiver on the roof (and they usually come up with a bunch of cool graphs showing where all the satellites are over time).
Soon thereafter, and a bit of reading about the NTP protocol, they realize that just adding:
server 0.pool.ntp.org
server 1.pool.ntp.org
server 2.pool.ntp.org
server 3.pool.ntp.org
to their ntp.conf is sufficient for 99.99% of all endeavors which require accurate time, outside of big physics, and, apparently Google's Spanner Database.This part was a bit incomplete:
"Typically, data-center operators keep their servers in sync using what’s called the Network Time Protocol, or NTP. This is essentially an online service that connects machines to the official atomic clocks that keep time for organizations across the world. But because it takes time to move information across a network, this method is never completely accurate,"
Much of the purpose (and math) behind the NTP protocol is to deal with network lag. And it does a pretty good job doing so.
Reading about the True Time Api at: http://static.googleusercontent.com/external_content/untrust...
"This implementation keeps uncertainty small (generally less than 10ms) by using multiple modern clock references (GPS and atomic clocks)"
So - apparently 10ms is their breakpoint - 10ms is about the limit of what you can expect out of NTP, so I guess it makes sense that if Google needs to do 10ms or better, something of their own invention would be required. Cool graph on the paper showing that 99.9% of variance across data centers thousands of kilometers apart are < 10ms deviation.
Finally - the entire purpose behind TrueTimeAPI is to sync up spanner's commits so they are consistent across the replicated/sharded nodes. Without Network, that database capability would come to a halt way before network time became a problem.
The bigger issue is the <10ms requirements. NTP over the Internet does not get you that consistently, and certainly not at the 99.9% success that Google was able to achieve with TrueTime API.
Also: > “We can commit data at two different locations — say the West Coast [of the United States] and Europe — and still have some agreed upon ordering between them,” Fikes says, “So, if the West Coast write happens first and then the one in Europe happens, the whole system knows that — and there’s no possibility of then being viewed in a different order.”
That's a large enough scale that you have to deal with relativity (light takes almost precisely 0.03 seconds to go from Palo Alto to Paris, eg). So in some sense there is no correct ordering. Anyone know how they deal with this? Have they just chosen some arbitrary point to make their reference frame, for purposes of ordering commits?
they've chosen gps time as a universal clock.
The part of the article that stood out to me is that Spanner is used in F1, the new backend datastore for AdWords. That's a significant vote of confidence.
That is an understatement.
"Advertising revenues made up 97% of our revenues in 2008 and 2009, and 96% of our revenues in 2010. We derive most of our additional revenues from offering display advertising management services to advertisers, ad agencies, and publishers, as well as licensing our enterprise products, search solutions, and web search technology."
2010 Revenue: 29.321 Billion
Therefore, advertising brought in >$28 Billion in 2010 (Therefore they average $1 Million in advertising revenue ever ~19 minutes, hence the laughable understatement of the article)
Source: http://investor.google.com/documents/20101231_google_10K.htm...
(http://www.home.agilent.com/en/pd-1000001383%3Aepsg%3Apro/pr...)
Using GPS is an operational convenience, not a necessity, I would think. If GPS didn't exist, they would have to have a master atomic clock, say, in Mountain View. Then they would have to bring each remote clock to Mountain View, synch it with the master, and then ship it (running continuously) to its final destination.
I remember perusing an HP catalog in the late eighties. You could buy a portable "traveling clock" for around $40,000. It was intended to be used to synchronize remote clocks. Here is a 1965 article from the Hewlett-Packard Journal, which describes the process:
(And if you're wondering, 15ns-accurate GPS clocks are $30 on eBay.)
I am sure that the U.S. armed forces have a very secure scheme to prevent this.
I'm a bit confused by this. How will this solve the situation when the first transaction renders a second transaction forbidden. To keep it simple, say an account with only $10 and two transactions trying to withdraw $10 each.
That is the case nearly everywhere, I think :-)