Ex-Googlers CockroachDB: A Scalable, Geo-Replicated, Transactional Datastore
github.com
github.com
I've been looking recently at long-term digital preservation systems -- tools designed to archive large amounts of data for decades. This is the Library of Alexandria problem -- how do we preserve all this data we're generating against once-in-a-century disasters?
So this 2005 paper lists thirteen different threats to long-term archives: Media Failure; Hardware Failure; Software Failure; Communication Errors; Failure of Network Services; Media & Hardware Obsolescence; Software Obsolescence; Operator Error; Natural Disaster; External Attack; Internal Attack; Economic Failure; Organizational Failure.[1]
Fault-tolerant distributed data stores are exciting, because they solve a bunch of those problems off the bat -- media failure, hardware failure, communication errors, failure of network services, hardware obsolescence, and natural disaster.
They also help to address software failure, software obsolescence, and economic failure, because archival projects are always strapped for resources and it's great to rely on tools that exist for totally distinct, commercially-valuable reasons.
But that still leaves operator error, external attack, and internal attack -- burning down the Library.
Hence my original question: are there distributed data stores that can be configured to resist intentional destruction of data?
[1] http://www.dlib.org/dlib/november05/rosenthal/11rosenthal.ht...
I guess you're asking whether there exists a distributed fault-tolerant with a form of version control (similar to git/cvs/perforce) as part of the native feature set.
Well, Git has checksums on everything.
The core of the issue is that humans view different information differently (child porn vs. Mona Lisa), whereas for computers, bits are bits and numbers are numbers. As long as child porn remains illegal and socially unacceptable, we'll want to enable attacks on data, i.e. for someone (usually internal operators) to be able to delete some kind of information, corrupt it or at least track it. Of course, this necessarily means that all information stored in the same data-store will be vulnerable.
So it would probably need to be write-only to prevent people from burning it down, which would necessarily mean that, once content is included, it cannot be modified or removed.
http://blog.dshr.org/2014/07/trac-certification-of-clockss-a...
But LOCKSS occupies a small niche. My hope is really that at some point a commercially-focused project with a ton of engineering effort and battle testing behind it will displace a lot of what LOCKSS has had to do manually. Seems like that might happen as web services get more and more distributed and fault-tolerant.
I'm curious, could you explain why? Google is an incredibly large company with many developers who never even touch their data systems so to me saying ex-Googler really doesn't mean anything beyond that they're probably a senior developer considering how rigorous (and honestly some old hat) their interview process is. But that doesn't change my viewpoint of the project at all.
With the tag being sign of pedigree. "A scalable DB from people who worked at a place with huge scale DBs".
Whether it is very convincing or not is another matter, but I am sure it gets more clicks/attention
.. so they've gone for CA and forgotten about P.
The section you're quoting is discussing a separate gossip protocol that is used to lazily propagate node state information. It does not affect the consistency of actual data replicas.
From the intro page:
"Cockroach is a distributed key/value datastore which supports ACID transactional semantics and versioned values as first-class features. The primary design goal is global consistency and survivability, hence the name. Cockroach aims to tolerate disk, machine, rack, and even datacenter failures with minimal latency disruption and no manual intervention."
If we read 'survivability' as 'availability', then that would suggest they've gone for CA. Although closer inspection reveals that their architecture seems to be made of shards each of which is maintained with Raft/Paxos. An evaluation of this by the Cambridge Computer laboratory is here: http://www.cl.cam.ac.uk/techreports/UCAM-CL-TR-857.pdf
That report makes two points relevant to this discussion. One is that a hard definition of C A and P can be difficult and that it's possible to achieve all three almost all of the time in real conditions. The other is from the conclusion:
"In particular, a [Raft] cluster can be rendered useless by a frequently disconnected node or a node with an asymmetric partition"
Let's reduce this case to a multi-master setup where a client can connect to any node and write to it. If a node fails outright and a client tries to connect, no big deal: the client chooses a different node, the failed node eventually comes back online later, catches up, then says "OK, write to me!" opening a listening socket.
However, if a partition happens, and client X writes to node A, client Y writes to node B, and then the two nodes cannot agree on the correct data, then you lose consistency. You can of course say that no node can be written to if other nodes are offline, which means the system is not highly available.
So their stated goal: "The primary design goal is global consistency and survivability..." either implies that high availability is not something they are going for, or that they are shooting for something that is not theoretically possible.
Of course all of the above is just my understanding, not necessarily fact, so please correct me if I'm wrong.
I strongly disagree. The "A" in CAP describes a system such that any single non-failing node can always make progress. That would be a nice property to have, but it's much stricter than is required for a real system.
If a distributed database is resilient to failure of a minority of nodes, it still makes sense to describe it as high-availability. And that is exactly what a consensus algorithm like Paxos or Raft gives you.
Unless this database can fly, I'd say they've chosen the right name. ;)
On the other hand, almost all my backend related features can be easily abstracted to those APIs.
Link? I see no relevant-looking "Facebook" or "TAO" in my recent RSS entries :(
MySQL: Weak consistency
Cassandra: No availability or weak consistency with datacenter failure
Am I the only person really sad that this didn't happen? after using apt over Tahoe-LAFS (over I2P - KillYourTV's PPA is on clearnet and I2P), I wanted this to be the default behavior for apt.
Hydra could have been another good choice :)
http://en.wikipedia.org/wiki/Hydra
See first entry at above Wikipedia page, about the many-header serpent.
On the other hand this question immediately popped into my mind: What kind of overhead does the GC incur, and how does it affect processes like a database where low latency is desired?
[0]: https://docs.google.com/document/d/16Y4IsnNRCN43Mx0NZc5YXZLo...
The eventual goal (golang v1.5 IIRC) is to have 40ms out of every 50ms available for actual processing. This is the kind of 'soft real time' that should provide good responsiveness for most clients most of the time.