FoundationDB: A distributed unbundled transactional key value store
micahlerner.com
micahlerner.com
My post focuses on the correctness, proof of the correctness of FDB’s failure recovery, some of which is missing even from the paper itself.
Some of the paper authors reached out; and I corrected one issue on my blog post pointed out by one of the authors.
Certainly had a devil of a time trying to get it to work on windows.
As of version 6.0, a single FoundationDB cluster currently only supports one active, writable 'master' region + one hot failover region. Transactions are all essentially LAN so performance characteristics are great. It doesn't support the slower TiDB/CRDB style active/active configuration where all nodes are writable, data is globally distributed and transactions can span across all nodes in the cluster.
Do they have some kind of scheme whereby ACID is preserved for local transactions only, while far away regions are asynchronously replicated to with less guarantees on data consistency?
This has nice properties for blast containment.
It's also an "unbundled" low-level component that one could use as the foundation for a database engine or whatever. According to Microsoft, FASTER is not just "fast", but significantly faster than even some basic in-memory data structures that ship in the .NET standard library!
The downside is that it doesn't (yet) support some more advanced features like multi-server distributed mode.
However, that relative simplicity may be preferred in some scenarios...
Their distributed checkpointing/recovery thing (built on FASTER) is very interesting!
The graph in the upper right corner of page 9 is certainly impressive: https://tli2.github.io/assets/pdf/dpr-sigmod2021.pdf (it is indexed with number of virtual machines)
Where I think it really shines is when you couple it tightly with your application. Suddenly you get a fully transactional distributed database with clearly understandable semantics, and if you can show some flexibility in how you map your data onto the underlying primitives, you can get fantastic performance, as well as correctness.
> Transaction size cannot exceed 10,000,000 bytes of affected data
In my experience almost all transactions are smaller than that, but the there is a long tail of large ones even for the same operation (because some 1:n link usually has a small n, but occasionally has an n which is orders of magnitude larger).
> FoundationDB currently does not support transactions running for over five seconds.
Limiting the duration of write-transactions isn't too bad. But being able to use a longer lived readonly snapshot is very convenient.
Do higher level databases building on FoundationDB (e.g. the SQL layer) typically work around some of these limitations, or do they leave it to the application to avoid these?
AIUI, the transaction size does not include values that are read, and for snapshot reads does not include the keys either, so this 10 MB limit only really constrains transactions that are writing a large volume of data.
For those transactions, the documentation suggests writing the data first (using manny smaller transactions) and then only updating a pointer to that data in the final transaction.
I'm currently building a simple database I'm calling AgentDB on top of FDB. It implements a message-passing system within the database, with guaranteed exactly-once semantics. It's designed to allow business logic to be more easily expressed without having to worry about all the failure modes typically present in a distributed system.
For this, I process messages in batches, with one transaction per batch. If a transaction fails, I retry with a smaller batch size, so in my case the answer would be "yes, to an extent". The user of AgentDB still has to ensure that processing a single message doesn't exceed the transaction size, but they don't have to worry that a batch of messages would exceed that.
See section 4: https://www.foundationdb.org/files/record-layer-paper.pdf
https://nikita.melkozerov.dev/posts/2019/06/building-a-found...
https://forums.foundationdb.org/t/roles-classes-matrix/1340/...
Public source for that info, so I don’t get accused of NDA violations - https://youtu.be/HLE8chgw6LI
Timestamp 14:15 touches on CloudKit’s extensive FDB usage.
[1] https://medium.com/@siddontang/benchmark-foundationdb-with-g...
We don't view ourselves in competition with FoundationDB because most TiKV users want the ability to use TiDB with TiKV for SQL (and MySQL compatibility). If you know you don't ever need that then FDB can be a good choice (but TiKV can be as well).
https://github.com/foundationdb
EDIT: I forgot they were a commercial product before Apple open-sourced it after acquisition.
The problem is with implementation. The more you separate components, the more code you have to add to allow them to interoperate. The more code, the more complexity. The more complexity, the more bugs. In addition, when you continue to separate components out into new failure domains that have a higher probability of failure ("3 components running on 1 node" -> "3 components running on 3 separate nodes") you increase the chances of problems even more. So, while the design looks nice in theory, and I'm sure they have some wonderful reports about how thorough their test framework is, in practice it might be a tire fire (depending on how it's actually used).
Some people are going to go, "I've been using this in production for a year, it's rock solid!" Well, people said that about Riak, until they wanted to set it on fire when they found out how buggy the implementation was at actually handling various errors and edge cases, and that they each had to hire dedicated Riak programmers just to fix all its problems. You can make large-scale Riak clusters work "at scale" if you hire enough on-call to constantly fight fires and monkey-patch the bugs and juggle clusters and nodes and repair indexes, and hand-wave outages and flaws as due to some other issue ("the cluster was completely hosed and had to be recovered from backup because the disks got corrupted and a node went down and a network partitioned all at once, it's not Riak's fault").
(looking at this HN thread, it seems like my assumption may be correct: https://news.ycombinator.com/item?id=27424605)
that has got one guy complaining about fdb, and a couple of others saying positive things.
so in short, you've got no experience with fdb, wrote a couple of paragraphs of speculation and references to other systems, and then declared victory, hoping nobody would read the other thread, i guess.
i've written stuff against fdb, and i've seen it in non-trivial production. it's not a panacea, it's a useful point in the design space of databases, and does pretty well there.
Remember a decade ago when NoSQL came out, and everybody was hooping and hollering about how amazing the concept was, and people like me would go "Well, wait a minute, how well does it actually run in production?" And people on HN would shout us down because the cool new toy is always awesome. And lo and behold, most NoSQL databases are no better in production than Postgres, if not a tire fire. People who have fought these fires before can smell the smoke miles away.
Zookeeper bugs are very rare even in production and hard to find.
If they can uncover zookeeper bug, i think they did some serious testing.
The most glaring of these bugs are SEUs caused by cosmic rays. Unless your test framework is running on 10,000 machines, 24 hours a day, for 3 years, you will not receive the SEUs that will affect production and cause bugs which are literally impossible without randomly-generated cosmic rays.
The simpler of the bugs are buggy firmware. Or very specific sections of a protocol getting corrupted in very specific ways that only trigger specific kinds of logic in a program over a long period of time. Or simply rolling over a floating point that was never reached in test because "we thought 10,000 test nodes was enough" or "we didn't have 64 terabytes of RAM".
And more examples, like the implementation or the administrative tools just being programmed shittily. Even if you find every edge case, you can write a program badly which will simply not deal with it properly. Or not have the tools for the admins to be able to deal with every strange scenario (a common problem with distributed systems that haven't been run in production). Or 5 different bugs happening across 5 completely different systems at the same time in just the right order to cause catastrophic failure.
Complex systems just fuck up more. In fact, if your big complex distributed system isn't fucking up, it's very likely that it will fuck up and you just haven't seen it yet, which is way more dangerous than not knowing how or when it's going to fuck up.
And if you do work in distributed systems at that level, please elaborate in detail what you're trying to refute, because it's all "I feel this should be bad. We need to be careful in trusting them!", which is unconvincing.