Live data from Hacker News

Consistency Without Clocks: FaunaDB's Distributed Transaction Protocol

fauna.com

81–90 of 105 posts

Re: Consistency Without Clocks: FaunaDB's Distributed Transaction Protocol

#81

We're here, ready for your questions to be consistently replicated across the WAN at the lowest latency information science allows ;-)

have a question about phantom read issue, i.e. two doctors are oncall for the night, but we have the constraint that at least one is needed to be oncall. If we have two transaction checking whether there is at least 2 doctors oncall, then remove one, seems we could end up with no one oncall with FaunaDB?

If so, seems this violate serializability?

Re: Consistency Without Clocks: FaunaDB's Distributed Transaction Protocol

#82
post #68

Earlier quoted context omitted.

But opting into the log for read-only txns would destroy performance for many apps, both in terms of latency and throughput. This is rarely practical since so many apps are read-heavy. I assume you don't recommend that to your customers, especially in global deployments. Causal tokens are practical from a performance point of view, but do require application awareness and changes, which presents a different set of is…

Really putting me on the spot when our CTO is out of town. My initial reply was ambiguous and appears to suggest that the transaction highwater mark/causal token preserves strict serializability. It preserves serializability. Strict serializability requires the log. Let me try to explain more clearly and you tell me what you would call it. It is subtle because there are three coordination points, the log, the replica…

Sorry to put you on the spot. Perhaps your CTO can chime in when he gets back.

Serializable tends to sound really good to database people, because it's the "highest" isolation level available in most databases. However, those databases are built for a single machine, and they really offer "strict serializable" without calling it that, because they maintain a single log which serializes all writes, and all reads operate as of the latest time. With the advent of distributed databases, the "strict" part is becoming more relevant, because it's no longer something you just get without realizing it.

I gave an example of the "blog post" anomaly that occurs in a system that only guarantees serializability. As another example, imagine I am using FaunaDB to store bank account balances. As the developer, I design a system where any time a customer receives money to their account, I send them a text giving the amount and their new balance, along with a link to the website where they can get more details. Except, when they click the link and go to the website, they don't see any change to their balance. Furthermore, when they refresh the page, sometimes they do see their new balance, and then they refresh again, and it disappears!

More investigation reveals that FaunaDB is in a partitioned state, and was serving stale bank account balances from some replicas. The replicas were behind a load balancer, and sometimes the web client hit one replica, and sometimes another, which explains the alternating balances. Is it possible to solve this problem using clever application logic, causal tokens, or saying "don't do that"? Sure. But a consistent database is supposed to prevent this scenario, not the application. Also, note that the "this is no different than increased client-server latency" argument does not hold here. People expect that if they finish action A (like writing to the DB), then do action B (like clicking a web link that reads the DB), that B will always see the results of A. The stale bank account balance isn't a delayed value; it's a wrong value.

You're correct that most systems allow some variant of "secondary" or "follower" reads that can return stale data, in exchange for better performance. I'm all for this kind of feature, as long as the developer understands the tradeoff they're making. Moreover, I think it's appropriate for FaunaDB to offer the same. But I'd advise against making unqualified claims of "complete guarantees of serializability and consistency" and the "latency profile of an eventually-consistent system like Apache Cassandra". FaunaDB is able to offer one or the other, but not both at once.

Re: Consistency Without Clocks: FaunaDB's Distributed Transaction Protocol

#83
post #78
post #59

The paper contains this claim: > FaunaDB is an elegant, software-only solution for achieving global ACID transactions, with complete guarantees of serializability and consistency. The paper makes it sound like FaunaDB claims strict serializability for all transactions (like Spanner). This means that if Txn(A) ends before Txn(B) begins, then Txn(B) must be guaranteed to see all data written by Txn(A). However, what if…

>It reads the value of X on replica R2, which does not yet reflect W(X). Why does someone want this functionality? Do you expect, say, Oracle to work this way? I honestly don't see this as a useful feature in a database that's operating at scale, but I'm open to hearing reasons for needing this.

Not sure I fully understand your question, but I'll try to answer:

FaunaDB (like other highly available databases) is set up to have multiple replicas of the same data. You can direct reads to any of the replicas. For example if replica #1 is unavailable due to a network outage, you can read replica #2 instead.

This thread is about the behavior of FaunaDB, which allows reads to replicas that have older stale values (i.e. they're not up-to-date with latest data). This causes the various consistency anomalies that I've been describing. The reason FaunaDB works like this is to make reads faster and to better distribute read load. Allowing secondary/follower/slave reads like this is a common tactic for increasing read performance.

Does that answer your question?

Re: Consistency Without Clocks: FaunaDB's Distributed Transaction Protocol

#84
“Consisteny without clocks” seems like a misnomer considering Raft indexes are tantamount to a logical clock. Perhaps it should be “consistency without wall clocks.” But that’s not really novel by itself, except in the context of geo-distributed databases.

Re: Consistency Without Clocks: FaunaDB's Distributed Transaction Protocol

#85
post #82

Earlier quoted context omitted.

Really putting me on the spot when our CTO is out of town. My initial reply was ambiguous and appears to suggest that the transaction highwater mark/causal token preserves strict serializability. It preserves serializability. Strict serializability requires the log. Let me try to explain more clearly and you tell me what you would call it. It is subtle because there are three coordination points, the log, the replica…

Sorry to put you on the spot. Perhaps your CTO can chime in when he gets back. Serializable tends to sound really good to database people, because it's the "highest" isolation level available in most databases. However, those databases are built for a single machine, and they really offer "strict serializable" without calling it that, because they maintain a single log which serializes all writes, and all reads opera…

Having the read information sent from one location (presumably where the write occurs) and then read from another is a good example of one of the few times this anomaly can occur.

I think it's a stretch to say that including a single timestamp in the link is too hard for developers. It happens by default across queries within a process with the Fauna drivers, so the database is preventing it; the only place you have to think about it is multi-process, multi-location reads. Also, having a loadbalancer jump back and forth between two replicas would not normally occur because of a geo routing, a soft form of datacenter affinity, but it is possible. Note that in your example, even if the link didn't include the timestamp, once they visit the webpage, the webserver should begin managing the timestamp in process or in a session cookie, preventing the the balance from disappearing once seen for the first time.

In my experience as a DBA the legacy database experience is much worse because in practice there are always physically distinct read replicas in the mix for scalability and hot failover that have no consistency guarantees at all, due to various faults like statement-based replication desynchronization, the lack of any ability to share log cursor information, and the like. I am not making a worse is better argument. I'm just pointing out that you don't get any C at all for the feeble amount of A gained under CAP in practice with most “single node” systems.

You are correct that the latency argument does not hold if you do not pass the timestamp information across queries. But if the observer can observe multiple replicas in the database topology, it is always possible to use the observer itself as the communication channel for the timestamp to avoid the anomaly. If you want to maintain strict serializability in all cases, but minimum-latency reads, the writer could block on transaction application acknowledgement from all replicas. Obviously that gives up availability for writes under partition, but it may be reasonable in your asynchronous example in order to avoid sending the txt message too early. Perhaps, for example, the message is "go to your local ATM" instead of a link with a timestamp. Although I would expect the ATM had the latency budget to do a linearized read even in a global context.

We are always happy to make the descriptions of the behavior and consistency levels more precise.

Re: Consistency Without Clocks: FaunaDB's Distributed Transaction Protocol

#86
I've published this article 5 months ago. Mine doesn't require global consensus for all commits, just local one and it works from there.

What's even more annoying is that I tweeted this article to them and they didn't say anything.

https://medium.com/p/optimistic-acid-transactions-4f193844bb...

Re: Consistency Without Clocks: FaunaDB's Distributed Transaction Protocol

#87

I've published this article 5 months ago. Mine doesn't require global consensus for all commits, just local one and it works from there. What's even more annoying is that I tweeted this article to them and they didn't say anything. https://medium.com/p/optimistic-acid-transactions-4f193844bb...

Sorry, let me clarify. Tweeted to them at the time. So, back in May.

Re: Consistency Without Clocks: FaunaDB's Distributed Transaction Protocol

#88

I've published this article 5 months ago. Mine doesn't require global consensus for all commits, just local one and it works from there. What's even more annoying is that I tweeted this article to them and they didn't say anything. https://medium.com/p/optimistic-acid-transactions-4f193844bb...

Did you implement it?

Re: Consistency Without Clocks: FaunaDB's Distributed Transaction Protocol

#89
I'm really confused how this scales to high transaction rates. If the replica has to redo all the reads (which means talking to multiple nodes) before it can make a commit/abort decision for the transaction this could take tens of microseconds if all nodes are in the same datacenter (if serving from RAM). Since it also has to process transactions from the log in order, that seems like it would limit the transaction rate to tens of thousands of TPS? Forget about distributing a replica across data centers or having enough data that it may not be RAM resident.

Is this actually how it works or am I missing something important?

Re: Consistency Without Clocks: FaunaDB's Distributed Transaction Protocol

#90

I've published this article 5 months ago. Mine doesn't require global consensus for all commits, just local one and it works from there. What's even more annoying is that I tweeted this article to them and they didn't say anything. https://medium.com/p/optimistic-acid-transactions-4f193844bb...

Did you implement it?

I'm planning on releasing a prototype as open source in a few months. I'm also planning on releasing my patent into the public domain.
Post reply on HN