A history of the distributed transactional database
infoq.com
infoq.com
In all those fraud cases banks were actually running ACID databases as they always do. Because ACID has nothing to do with "external consistency", unless you can run it in a human brain, which you can't. Interactions with humans can only be eventually consistent. But nobody even bothers to have ACID all the way to the web browser or APIs, so it's only limited to the application, but not beyond that. With all the frauds and double spending problems remaining unsolved.
Overall the usual FaunaDB attacks on other systems, especially AP. Strong eventual consistency of AP systems is actually much stronger, than ACID and is much much easier to verify and AP is where the cutting edge of distributed systems research happens. But more importantly, SEC can offer the best theoretically possible performance and latency for multi datancenter deployments. While distributed ACID transactions don't even show viability in multi datacenter deployments at the moment, forcing vendors to degrade into weak consistency models in such setups.
Not sure what you're getting at when you say ACID transactions don't work in a multi-DC context. FaunaDB aside, there are many production systems (commercial or internal) that provide them that I am aware of.
I expect they are talking about DC as _regions_, not DC as _availability zones_ (in the AWS parlance).
Planetary databases hit latency issues on cross-region ACID transactions. Whatever the sharding, a transaction modifying one piece whose leader is in Paris and another in Beijing will have at least the round-trip time between the two.
It is already unforgiving at that scale, and will become inadmissible at interplanetary scale, requiring some form of eventual consistency.
However, I still expect that we will retain ACID operations for local changes. Humans tolerate latencies that are necessary for planetary ACID transactions. The combination of the size of the Earth and the speed of light is a lucky draw.
We're indeed pretty lucky as far as transactions are concerned...
I don't think there's a great reference for the Percolator transaction architecture specifically, but there is a brief description here: https://blog.octo.com/en/my-reading-of-percolator-architectu...
One thing I learned recently is that FoundationDB is essentially the Percolator model (timestamp oracle), even though FoundationDB predated the Percolator paper.
The main difference from actual Percolator is that FoundationDB manages range locks through a logically separate pool of in-memory lock caches, instead of writing durable locks to the replicas, but the end result is the same. External consistency for read/write transactions within a single datacenter at a time.
EDIT: I was wondering why this conversation was giving me deja vu, and then I remembered that you and Dave discussed this in detail last month: https://news.ycombinator.com/item?id=18258805
https://pdfs.semanticscholar.org/9c08/4a56c1c625da116eca4742...
But if you read between the lines, and are savvy enough to look up the background research papers, you'll pick up most of the technology. The devil as they say however, is in the engineering details, and that's where you'll run into a lot of special sauce.
You'll want to look at their Transaction Monitoring Facility (TMF) in detail. They baked fault tolerance right into the bones of the architectural design, and TMF is an identifiable chunk of how that was done. You aren't going to be able to pull off the same kind of fault tolerance they're delivering at just the database level, too. Their thinking about fault tolerance goes way beyond, down into the hardware level though it stops short of the silicon.
I don't have the background to fully understand this stuff, but it claims "The first distributed SQL -- it offers distributed data, distributed execution, and distributed transactions." while every single NoSQL/NewSQL article says things like "Legacy databases (typically the centralized RDBMS) solve this problem by being undistributed."
Rather than build consensus-based high availability into the software, NonStop SQL relies on custom hardware and redundant networks that are designed to never fail. If they do fail, some portion of the data or some types of modifications become unavailable and the operator is expected to intervene and make choices about how to proceed.
They give an example of being unable to run an ALTER TABLE statement because two datacenters that each own a shard of the table are disconnected—the statement can't be run at either datacenter during the partition.
They give other examples where queries return partial results because a specific shard is down. This is potentially a correctness violation if it is not explicit, and it is definitely a loss of availability.
Modern distributed systems are instead designed to maintain almost total levels of availability on constantly failing hardware, and heal themselves immediately and autonomously without losing correctness.
'This model is equivalent to the multiprocessor RDBMS—which also uses a single physical clock, because it’s a single machine—but the system bus is replaced by the network. In practice, these systems give up multi-region scale out and are confined to a single datacenter.'
I don't think so, although Percolaotor has a single point of physical clock, but it is easy to eliminate the SPOF problem by using high-availability algorithms such as Raft or Paxos (Just like the design in TiDB). On the other hand, the strong consistency across multiple data centers, I think, without hardware devices such as TrueTime, it is difficult to have low latency and strong consistency while allowing writes happen in multiple data centers at the same time. The more realistic situation is that there is a primary data center (where both read and write occur), and multiple data centers guarantee high availability with strong consistency, we can achieve this goal by scheduling Paxos or Raft leader to the primary data center.
Another point worth mentioning is that, single transaction coordinator model is undervalued in community. Majority of systems have the following query pattern (sorted descending):
1. Single key read --> you don't need to contact coordinator
2. Single key write --> you don't need to contact coordinator. OCC with simple version number works here.
3. Multi-key asynchronous write (updating second, third keys using a pub/sub infrastructure) --> you don't need to contact coordinator in this case.
4. True multi-key transactions that should happen atomically --> need to contact coordinator
If a single coordinator can process 20-30K/s transactions of category 4, then total throughput of the system would be around 300K queries per second which is sufficient for many workloads.
Especially American banks. I can't think of any other retail banking market where banks abuse eventual consistency so bad that transfers take days to clear.