4 ms·
I think Citus is not really ACID. Spanner (and to an extent its less mature OSS descendants Cockroach and TiKV) has more comparable goals, but is fairly differ
by voidmain 8y ago
I think Citus is not really ACID.
Spanner (and to an extent its less mature OSS descendants Cockroach and TiKV) has more comparable goals, but is fairly different architecturally. For example, FoundationDB only requires N+1 replicas instead of 2N+1 to achieve N failure tolerance (even lots of databases with much weaker guarantees are in the latter category!), doesn't trust clocks at all, doesn't lose performance when transactions cross replica sets, and uses optimistic instead of pessimistic concurrency.
Also FoundationDB (and TiKV) make a distributed, transactional key/value store available as an API, while Spanner and Cockroach expose only a relational database layer. FoundationDB is designed philosophically with the idea that you want to have a single storage layer to manage operationally but should be able to mix and match data models and query engines above that layer.
On the other hand, FoundationDB doesn't currently have any full fledged high level database layer available. Someone will probably dig up our SQL layer (which was AGPL, I think) but I wouldn't really recommend using it in production because there is no active development team. Someone will probably try porting the SQL layers from TiDB and Cockroach.
Maybe Apple will open source more stuff in the future, but let's not get too greedy!
- infogulch 8y agoWait how can it be ACID if it can tolerate N failures from N+1 nodes? Doesn't that kill consistency almost by definition? What level of isolation does it support? Snapshot? Serializable?
- voidmain 8y agoIt is serializable and totally uncompromising. Philosophically pretty much everything defaults to the safest possible thing. It can't tolerate N failures from N+1 nodes. It can tolerate N failures with N+1 copies of your data. In a big cluster you have plenty of nodes but storing everything 5 times to tolerate 2 failures is really expensive.
- infogulch 8y ago> It can't tolerate N failures from N+1 copies of your data Sorry I got the terminology wrong, but that's a distinction without a difference. If it can tolerate N failures from N+1 copies, that means a network partition would allow any one copy to continue chugging along making changes by itself. You have two options: consistency is dropped and you downgrade to eventually consistent (at best), or availability is dropped meaning a single node can't make changes without a majority, which invalidates the N of N+1 failures claim. (Which is where the N failures of 2N+1 copies claim comes from in the first place: after N failures you still have a majority of copies.) Or there's some other magic quorum protocol I've never heard of that makes the majority problem disappear.
- voidmain 8y agoFoundationDB stores 2N+1 copies of some "coordination state" and does a consensus algorithm whenever it is updated. But this state doesn't contain a copy of your data; basically think of it as storing a replication configuration. It's very small and rarely changes. In the happy case, replication takes place using the replicas and quorum rules specified by this configuration. For example, you might require writes to succeed synchronously against all N+1 replicas of some transaction log. After N failures, there will still be 1 replica remaining with the latest transactions. But in order to proceed after any failures, you have to do a consensus transaction against a majority of replicas of the coordination state, to specify the new set of N+1 replicas you will be using. And you also make sure that the 1 replica you are recovering from knows you are doing it, so that it won't continue to accept writes under the old replication configuration. There can't be two partitions capable of committing transactions, because (in this case) you need either (a) All N+1 replicas of the log, so that you can commit synchronously, or (b) A majority (N+1 out of 2N+1) of the replicas of the coordination state, AND 1 replica of the log Sorry if this isn't a great explanation. Anyway it does work. I expect that you could rephrase this as an optimization of a consensus protocol, though I think it would be hard to build a performant and realistically featureful implementation that way.
- infogulch 8y agoNope your explanation made sense, thank you! When I wrote that I was wondering if it used a second 2N+1 dataset just for coordination & consensus. This has the benefit of separating data from consensus, allowing the N of N+1 data failure. But at the end of the day consistency still comes down to a N of 2N+1 failure tolerance of that second coordination state. It's smaller easier to replicate etc etc but it seems like it still has the same fault tolerance as just replicating the data 2N+1 times. It sounds like it's worked out great in practice for FDB. But you say it rarely changes... but wouldn't it have to change every time there's a change to the dataset? I feel like this means you have to do even more replication and consensus than just replicating the data without this second consensus state.
- wll 8y agoThe datacenter-aware mode documentation [0] says “Although data will always be triple replicated in this mode, it may not be replicated across all datacenters.” Why is that? [0] https://apple.github.io/foundationdb/configuration.html?datacenter-aware-mode https://apple.github.io/foundationdb/configuration.html?data...
- voidmain 8y agoI think it's just saying that it's willing to place two of the three replicas in a datacenter, for example if one of the three datacenters is down. This has downsides, since losing a datacenter will make it aggressively fill up disks, but mitigates against subsequent failures causing data loss. Most of the people who have run FoundationDB at scale have, for performance reasons, used configurations other than the "datacenter aware" mode for their inter region replication, so they may not be the strongest thing operationally. There is some work that from what I can see in the code is still in progress to build a new, almost magical inter-region replication mode that I am very excited about, which combines synchronous replication to a "satellite" datacenter within region with asynchronous replication between regions and recovery logic that will finish replication and fail over in case of a partial failure of a region. You get fast transaction commits (much less than the inter region ping time), can fail over to a secondary region automatically and safely (without losing any committed transactions) in the vast majority of circumstances, and in the worst case you can (manually, because you are accepting data loss!) give up very recently committed transactions to fail over.
- wll 8y agoHow would FoundationDB stay externally consistent with asynchronous cross-region replication? Thank you for your time and FoundationDB—along with @nlavezzo, and team(s)!
- voidmain 8y agoThe satellite mode that I described is an active/passive mode. One region is accepting reads and writes; the other is just replicating everything. When it looks like the active region is in trouble, the asynchronous replication is "finished up" before switching over to the other region. The multiple datacenters in each region ensure that usually a regional failure will be "slow enough" that this automatic process (which after all only takes hundreds of milliseconds to seconds) can usually complete before a region goes away. And this will be handled pretty transparently by the datastore. If a region is blown up instantly by an orbital laser cannon, then the database will go down and you will have to manually tell it to recover ACI in the other region, sacrificing the durability of whatever committed transactions in the lost region were destroyed by the laser cannon.
- hangonhn 8y ago"or example, FoundationDB only requires N+1 replicas instead of 2N+1 to achieve N failure tolerance" Wait, don't you need 3N+1 to tolerate N number of failures for it to be Byzantine fault tolerant? Is that not a goal of FoundationDB?
- voidmain 8y agoFoundationDB is not Byzantine fault tolerant.