11 ms·
Parallel Commits: A New Atomic Commit Protocol for Distributed Transactions
- aboodman 7y agoWhy is it so important for transactions to be interactive? I feel like this is not a feature that I use frequently as an application developer and it makes distributed transactions so much harder. It seems like something that accidentally came along from SQL, and is imposing a large cost.
- zaphar 7y agoInteractive may be a little misleading here. Consider a common case in any java application at my current place of employment. 1. A controller starts a transaction. 2. Some code looks up a record in a database. 3. Based on the results of that lookup we run 1 of two possible writes to a table in that database. Knowing what kind of write you will do in the transaction requires knowing the result of the lookup which itself must be run in the transaction. Now take this small case and imagine it in the real world scenario where there might many such lookup specific write logic all wrapped in a single transaction and the space of possible combinations of statements you will want to execute in the transaction is large. None of this is interactive in the sense of someone working in a shell. But it is interactive in a way that makes the transaction useful in the real word.
- aboodman 7y agoExactly. What Calvin (and FaunaDB, I think) does is allow you to run an arbitrary function which takes input on the database. That function can do whatever it wants, including any number of reads, writes, and arbitrary logic based on the reads. But critically, it can't talk back to the calling client, except for to send the final result. This allows you to implement the pattern you describe, which I agree is common, but in a dramatically simpler way. Having the database, which is not really a single thing, but a swarm of computers spread across the globe separated by unreliable links who are trying to stay in consensus, pause their work to hear back from the client just seems ... well it seems miraculous that it can work at all. The Calvin way is so much easier, it seems like there must be some very good reason that it's not what CockroachDB does. But I've never heard what that reason is.
- foota 7y agoWouldn't you have to implement a DSL and all the parts related to it for that to work? Also things like serialization cost. What if I have a large local in memory structure I want to base my query results off of? Funny enough this is kind of similar to the problem solved by things like apache beam.
- aboodman 7y agoNo you'd need a complete programming language. But SQL is basically one already, most variants are already turing-complete. You'd also have to provide the programming environment with concepts of cursors and so on so they could page through data efficiently.
- Pxtl 7y agoEvery procedural layer I've ever used bolted onto SQL (pl/SQL, t-sql) has been absolute goddamned agony to use. SQL is a good (if dated) language for relational access and manipulation, but awful for procedural scripting.
- evanweaver 7y agoNobody in industry understood how to apply Calvin until we did at Fauna. That is the only reason; the rest is engineering path dependence.
- aboodman 7y agoif you don’t mind sharing, I’m curious to learn more what was the main challenge was applying Calvin. Thanks!
- evanweaver 7y agoSee @freels’ reply above. I think it is fair to say that the Calvin paper is visibly incomplete and expresses some constraints in a way that makes them seem insurmountable when they are not; specifically, they are only constraints within the log, but do not constrain the database experience overall. Applying Calvin to traditional RDBMS workloads was a very unlikely creative exercise because it required questioning these explicit constraints. The Spanner paper also leaves a lot unexplained, but it is less obvious until you are too far down the path to turn back. After all, it worked for Google. Calvin did not have that real-world proof. Combine that with the pessimism of the paper itself and nobody was willing to pick it up.
- BubRoss 7y agoThis article implies that it reduces latency.
- cryptica 7y agoI like the elegance and simplicity of two-phase commits. I didn't understand the criticism in the article; maybe it's something specific to CockroachDB. In my experience with two-phase commits, if the system crashes before a transaction is fully processed and committed, it should be fully reprocessed (from scratch) once the server restarts. In one of my open source projects with RethinkDB (which doesn't natively support atomic transactions), I implemented a distributed 'parallel' two-phase commit mechanism by assigning each pending record a shard key (an integer derived from a hash of the record id); then worker servers would decide which subset/range of DB records to process/commit based on a hash of their own server id (which would tell them which range of shard keys/records they were responsible for). Only when a record had been fully processed, its status would be updated as committed by writing a 'settled' flag on that record. Whenever a server failed and restarted, it would pick up processing from the last successful commit. If a server did not restart, the worker count would be updated and remaining workers would redistribute the partitions among themselves based on the shard keys of the records.
- irfansharif 7y agoThe "criticism" as it pertains to CockroachDB and off-the-shelf 2PC is less so about 2PC in isolation and more so about layering 2PC on top of consensus groups used to persist records. When any given txn does the "prepare" phase, it lays down markers for a possible upcoming commit. If the 2PC coordinator fails in an inopportune moment, there's a delay between the failure and the markers being cleaned up (whether or not the transaction is "reprocessed" or aborted). The reason why this delay is problematic is because any subsequent transactions that happen upon said markers, they just have to wait for the resolution ("commit"/"aborted") aka it blocks. So clearly recovery must be built into 2PC, i.e. the transaction state itself must be persisted. This is done so in the same way the markers/regular writes are, through consensus. But marking the transaction state as "committed" can only happen once we're guaranteed that all the individual write markers are persisted. Which adds a second round of consensus.
- ryanworl 7y agoThis is about reducing the number of message delays before the commit succeeds. Failure scenarios have to be handled to be correct, but this is a performance optimization primarily from what I can see.
- foota 7y agoWow, this is impressive.
- cdbattags 7y agoWhat I'll say might come off naive but I accept this at an attempt for a reductionist's viewpoint. A singular update that side affects *N (times N) where N is greater than 1 will always lead to either a race condition or latency.
- ComodoHacker 7y agoFrom what I understand, with this new commit protocol they managed to improve response time of a writing transaction by shifting some work of determining its final status to the readers. Am I understanding correctly that readers' performance will degrade by the same amount? While this is an achievement which some scenarios will definitely benefit from, like bulk loading or updating secondary indexes as mentioned in the article, what about other scenarios where read performance is more important, like the ones where there are much more readers than writers? Shouldn't there be a configuration option which commit protocol to use?
- irfansharif 7y ago> Am I understanding correctly that readers' performance will degrade by the same amount? Not quite. The "slow path" talked about is only applicable when the coordinator node is unavailable (presumably a rare event). If it's unavailable, there's nobody left to clean up the STAGING txn record, so the reader is tasked to do it itself. In normal conditions however, once the coordinator node receives acknowledgement for the successful persisting of all write intents and the "txn STAGING" record, it can simply record "txn COMMITTED" in memory (and return to the client, send off async intent resolution procedures, etc.) Any subsequent read requests that observes left over intents (yet to be resolved) are pointed to the coordinator node, which can simply consult the "txn COMMITTED" record in memory. This is all safe because the commit marker is not simply stored on the coordinator node, it's a distributed condition and can be reconstructed by any observer even if the coordinator failed.
- ComodoHacker 7y ago>Any subsequent read requests that observes left over intents (yet to be resolved) are pointed to the coordinator node, which can simply consult the "txn COMMITTED" record in memory. Another roundtrip performed by reader rather than writer? That's what I'm talking about. Though I understand it differently, as reader just waits until all writes and "txn COMMITTED" record arrive at it's node.
- jlokier 7y ago
- grogers 7y agoGiven how critical preventing future intent writes is to the protocol to ensure safety during recovery, it'd be nice to have more detail on how that works. Calling it an in memory data structure doesn't exactly inspire confidence.
- irfansharif 7y agoAgreed, we should be talking about the "timestamp cache" in more detail generally. While I'm here, looking at [1] helped me confirm how everything is kosher despite being in-memory. The timestamp cache pessimistically maintains a "low water mark", this always ratchets up monotonically and represents the earliest access timestamp of any key in the range. Writes happening at timestamps lower than this water mark are not let through, and bumping this watermark past the point of the observed write intent's timestamp is how slow inflight write intents are aborted by recovering read requests. On server restart, the timestamp cache is initialized with a low water mark of the current system time + maximum clock offset, so "future" write intents (or more accurately: write intents sent in the past but previously stuck in transit) are simply rejected. [1]: https://github.com/cockroachdb/cockroach/blob/master/pkg/storage/tscache/cache.go https://github.com/cockroachdb/cockroach/blob/master/pkg/sto...
- mcms 7y agoWhat prevents a transaction from being prematurely aborted while there are some intents in flight? From what I understand, transaction A can still be replicating intents and transaction B, not knowing that coordinator for transaction A is still at work, start a recovery process. This makes transaction A abort which could be prevented by waiting for intents of A to be successfully replicated.