3 ms·
I'll try to give you a quick introduction. The architecture talk I recorded for new engineers working on the product ran to four or five hours, I think :-). In
by voidmain 8y ago
I'll try to give you a quick introduction. The architecture talk I recorded for new engineers working on the product ran to four or five hours, I think :-). In short, it is serializable optimistic MVCC concurrency.
A FDB transaction roughly works like this, from the client's perspective:
1. Ask the distributed database for an appropriate (externally consistent) read version for the transaction
2. Do reads from a consistent MVCC snapshot at that read version. No matter what other activity is happening you see an unchanging snapshot of the database. Keep track of what (ranges of) data you have read
3. Keep track of the writes you would like to do locally.
4. If you read something that you have written in the same transaction, use the write to satisfy the read, providing the illusion of ordering within the transaction
5. When and if you decide to commit the transaction, send the read version, a list of ranges read and writes that you would like to do to the distributed database.
6. The distributed database assigns a write version to the transaction and determines if, between the read and write versions, any other transaction wrote anything that this transaction read. If so there is a conflict and this transaction is aborted (the writes are simply not performed). If not then all the writes happen atomically.
7. When the transaction is sufficiently durable the database tells the client and the client can consider the transaction committed (from an external consistency standpoint)
The implementations of 1 and 6 are not trivial, of course :-)
So a sufficiently "slow client" doing a read write transaction in a database with lots of contention might wind up retrying its own transaction indefinitely, but it can't stop other readers or writers from making progress.
It's still the case that if you want great performance overall you want to minimize conflicts between transactions!
- openasocket 8y agoThanks, that explanation is really helpful!
- alexashka 8y agoGreat write-up. Is this similar to how Software Transactional Memory (STM) is implemented? It sounds very very similar indeed.
- jbit84 8y agoThanks. Can you elaborate on how 6 is actually accomplished? Various earlier comments have hinted that the transactional authority (conflict checking) can actually scale 'horizontally' beyond the check-throughput that can be archived by a single node. Is that the case? and whats the magic sauce for doing that for multi-object transactions? :)
- voidmain 8y agoYes, conflict resolution is for most workloads a pretty small fraction of total resource use so you usually don't need a ton of resolvers (I think out of the box it still comes configured with just one?), but it can scale conflict resolution horizontally. The basic approach isn't super hard to understand, though the details are tricky. The resolvers partition the keyspace; a write ordering is imposed on transactions and then the conflict ranges of each transaction are divided among the resolvers; each resolver returns whether each transaction conflicts and transactions are aborted if there are any conflicts. (In general the resolution is sound, but not exact - it is possible for a transaction C to be aborted because it conflicts with another transaction B, but transaction B is also aborted because it conflicts with A (on another resolver), so C "could have" been committed. When Alec Grieser was an intern at FoundationDB he did some simulations showing that in horrible worst cases this inaccuracy could significantly hurt performance. But in practice I don't think there have been a lot of complaints about it.)
- ccleve 8y agoThis is a good explanation of how it happens on a single node. What do you do when the transaction is distributed? How do you achieve consensus? Is there a write up on it anywhere?
- wwilson 8y agoThe only thing that's different in a distributed cluster is the implementations of steps 1 and 6. As voidmain said, the details of that are not trivial, ESPECIALLY the details of how it never produces wrong answers during fault conditions. I don't know that there's been an exhaustive writeup of that part, but maybe one of us or somebody on the Apple team will put something together. It probably won't fit in an HN comment though! Or... maybe this is the part where I point out that the product is now open-source, and invite you to read the (mostly very well commented) code. :-)
- veesahni 8y agoThe documentation ( https://apple.github.io/foundationdb/technical-overview.html https://apple.github.io/foundationdb/technical-overview.html ) sells the product, but doesn't give a deep enough explanation. As a closed source product, that's understandable. Going forward as an opensource product, I hope to see some clarity on the "how it works"... Distributed, performant ACID sounds good, almost too good to be true. Not that I doubt it at the moment, I just want to understand it better :)
- tmoertel 8y agoDoes the implementation handle the case that you want to do a write that is conditioned on a prior read finding no corresponding record(s)?
- itp 8y agoOf course. FDB thinks about read and write conflict ranges, which are functions of the keys, not the values (or lack thereof). A read of a non-existent key conflicts with a write to that key. A read of a range of keys conflicts with a write to a key in that range, even if that key did not have an associated value at the time of the original read.
- techdragon 8y agoAny chance you’d be able to find that video, it sounds extremely interesting.