FoundationDB: A Distributed Key-Value Store
cacm.acm.org
cacm.acm.org
https://www.youtube.com/watch?v=4fFDFbi3toc
> We wanted FoundationDB to survive failures of machines, networks, disks, clocks, racks, data centers, file systems, etc., so we created a simulation framework closely tied to Flow. By replacing physical interfaces with shims, replacing the main epoll-based run loop with a time-based simulation, and running multiple logical processes as concurrent Flow Actors, Simulation is able to conduct a deterministic simulation of an entire FoundationDB cluster within a single-thread! Even better, we are able to execute this simulation in a deterministic way, enabling us to reproduce problems and add instrumentation ex post facto. This incredible capability enabled us to build FoundationDB exclusively in simulation for the first 18 months and ensure exceptional fault tolerance long before it sent its first real network packet. For a database with as strong a contract as the FoundationDB, testing is crucial, and over the years we have run the equivalent of a trillion CPU-hours of simulated stress testing.
It’s not a full-blown simulator (because generally the application code doesn’t even exist yet when you’re building the TLA+ model). But it can let you collect data and validate assumptions about your design before writing a single line of code.
In reality the design stage is a pretty critical phase so you need all the help you can get, so even if you don't like TLA+ you're way better off than not modeling at all.
As an example of the language specific thing, though, there's a library for Haskell I like that's very cool, called Spectacle, which also implements the temporal logic of TLA+ along with a model checker, but as a Haskell DSL. An interesting benefit of this is that you can model check actual real Haskell code that runs e.g. in your services, but I haven't taken this very far. There are also alternative solutions like Stateright for Rust. But again, not everyone has the benefit of these...
Unfortunately I got the product wrong; it was not Atlas, it was Realm Sync. All of the test-case generation stuff is in Section 5.
Using https://github.com/awslabs/shuttle which works on our real Rust code.
Documentation is still generally useful, and so is a model. You have to be committed to keeping both up to date as the code evolves.
This is such an interesting topic!
Some thoughts:
* I wonder if the approach could be used to implement debuggable replayability, with accurate tracing and profiling. A bit like what verdagon is doing with Vale.
* It could be used to integrate the event loop with tracing (rather than instrumentation with Jaegar)
* I really like the idea that "every object" is an event loop, which reminds me of Microsoft Orleans with its actor model for its grains.
* I am interested with actor and lightweight thread architectures.
* I am interested in the scalabiliy of nodejs event loop architecture and Win32 desktop application programming.
* I think this approach could be used to test and simulate microservices.
* Approach could be used to test GUIs with React Redux reducer style.
https://github.com/AaronFriel/hyhac/blob/master/test/Test/Hy...
Alas, I think Hyperdex development paused a few years later. It's a shame that it stopped then.
However, they (at least at the time most of the developers were at Apple, many have now moved to Snowflake and the Apple team has grown a little I think) haven't released or integrated their nightly cluster and performance testing systems into open, nor have they integrated them with GitHub Actions or Nightly runs or anything. My understanding is that this is "just" a lot of compute cluster/platform orchestration code on top of the tests that exist in the repository. So, while Apple or Snowflake integrates changes across hundreds of concurrent fuzzing simulations on whatever platforms they have, if you write patches yourself, you're stuck with long simulation runs. Maybe that's changed; I haven't kept up since the 7.0 series.
In practice if you write patches and they accept them, they will just do the testing in their runs for you, on a cluster far larger than what you could have. Failures reports will tell you how to reproduce them from the test files. As a contributor, testing the system on your own is mostly a matter of how much money or how many CPU cores you can personally stand to set on fire.
Someone could probably integrate this functionality into a Kubernetes operator or something so that outside engineers could run large scale simulations reliably. But it is really expensive and CPU/compute intense, no matter how you go about it.
[1] https://forums.foundationdb.org/t/how-to-use-foundationdb-un...
https://github.com/apple/foundationdb/tree/main/fdbserver/wo...
> Someone could probably integrate this functionality into a Kubernetes operator or something so that outside engineers could run large scale simulations reliably. But it is really expensive and CPU/compute intense, no matter how you go about it.
Maybe this.
And yes, I linked to the spec files because there actually isn't that much test code written in Flow I feel; the high-level specs in the .txt files can be mixed and matched so much to create a lot of variety from some small number of primitives, so that's really where all the good stuff is. Implementation vs interface, and all that.
However, since these topics are always filled with effusive praise in the comments, let me give an example of a distributed scenario where FDB has shortcomings: OLTP SQL.
First, FDB is clearly designed for “read often, update rarely” workloads, in a relative sense. It produces multiple consistent replicas which are consistently queryable at a past time stamp, without a transaction - excellent for that profile. However, its transaction consistency method is both optimistic and centralized, and can lead to difficulty writing during high contention and (brief) system-wide transaction downtime if there is a failover; while it will work, it’s not optimal for “write often, read once” workloads.
Secondly, while it is an ordered key value store - facilitating building SQL on top of it - the popular thought of layering SQL on top of the distributed layer comes with many shortcomings.
My key example of this is schema changes. Optimistic application, and keeping schema information entirely “above” the transaction layer, can make it extremely slow to apply changes to large tables, and possibly require taking them partially offline during the update. There are ways to manage this, but online schema changes will be a competitive advantage for other systems.
Even for read-only queries, you lose opportunities to push many types of predicates down to the storage node, where they can be executed with fewer round trips. Depending on how distributed your system is, this could add up to significant additional latency.
Afaik, all of the spanner-likes of the world push significant schema-specific information into their transaction layers - and utilize pessimistic locking - to facilitate these scenarios with competitive performance.
For reasons like these, I think FDB will find (and has found) the most success in warehousing scenarios, where individual datum are queried often once written, and updates come in at a slower pace than the reads.
Whether or not concurrency is optimistic (or done with locks, or whatever) doesn't really have a bearing on things. Any database is going to suffer if it has a bunch of updates to a specific hot keys that needs to be isolated (in the ACID sense). As long as your reads and writes are sufficiently spread out you'll avoid lock contention/optimistic transaction retries.
You speak to the real main limitation of FoundationDB when you talk about stuff like schema changes. There is a five-second transaction limit which in practice means that you cannot, for example, do a single giant transaction to change every row in a table. This was definitely a deliberate deliberate design choice, but not one without tradeoffs. The bad side is that if you want to be able to do something like this (lockout clients while you migrate a table) you need a different design that uses another strategy, like indirection. The good side is that screwed-up transactions that lock big chunks of your DB for a long time don't take down your system.
I find that the people who are relatively new to databases tend to wish that the five second limit was gone because it makes things simpler to code. People that are running them in production tend to like it more because it avoids a slew of production issues.
That said, I think for many situations a timeout like 30 or 60 seconds (with a warning at 10) would be a better operating point rather than the default 5 second cliff.
All databases do suffer under some red line of write contention; but optimistic databases will suffer more, and will start degrading at a lower level of contention. “Avoiding contention” is database optimization table stakes, and you should be structuring every schema you can to do so; but hot keys are almost inevitable when a certain class of real-time product scales, and they will show up in ways you do not expect. When it happens, you’d like your DBMS to give as much runway as possible before you have to make the tough changes to break through.
SQL-on-top becomes an issue for geographic distribution; without “pushing down” predicates, read-modify-write workloads, table joins, etc. on the client can incur significant round-trip time issuing queries. I think the lack of this is always going to present a persistent disadvantage vs selecting a competitor.
And again, given FDBs multiple-full-secondary model, it’s only a problem when working in real time, slower queries can work off a local secondary. But latest-data-latency is relevant for many applications.
A great example of how to best utilize FDB is Permazen [1], described well in its white paper [2].
Permazen is a Java library, so it can be utilized from any JVM language e.g. via Truffle you get Python, JavaScript, Ruby, WASM + any bytecode language. It supports any sorted K/V backend so you can build and test locally with a simple disk or in memory impl, or RocksDB, or even a regular SQL database. Then you can point it at FoundationDB later when you're ready for scaling.
Permazen is not a SQL implementation. Instead it's "language integrated" meaning you write queries using the Java collections library and some helpers, in particular, NavigableSet and NavigableMap. In effect you write and hard code your query plans. However, for this you get many of the same features an RDBMS would have and then some more, for example you get indexes, indexes with compound keys, strongly typed and enforced schemas with ONLINE updates, strong type safety during schema changes (which are allowed to be arbitrary), sophisticated transaction support, tight control over caching and transactional "copy out", watching fields or objects for changes, constraints and the equivalent of foreign key constraints with better validation semantics than what JPA or SQL gives you, you can define any custom data derivation function for new kinds of "index", a CLI for ad-hoc querying, and a GUI for exploration of the data.
Oh yes, it also has a Raft implementation, so if you want multi-cluster FDB with Raft-driven failover you could do that too (iirc, FDB doesn't have this out of the box).
And because the K/V format is stable, it has some helpers to write in memory stores to byte arrays and streams, so you can use it as a serialization format too.
FDB has something a bit like this in its Record layer, but it's nowhere near as powerful or well thought out. Permazen is obscure and not widely used, but it's been deployed to production as part of a large US 911 dispatching system and is maintained.
Incremental schema evolution is possible because Permazen stores schema data in the K/V store, along with a version for each persisted object (row), and upgrades objects on the fly when they're first accessed.
[2] https://cdn.jsdelivr.net/gh/permazen/permazen@master/permaze...
If instead of using some generic K/V backend, it made use of specific FDB features, it might be even better. Conflict ranges and snapshot reads have been useful for me for some background index building designs, and atomic ops have their uses.
> Oh yes, it also has a Raft implementation, so if you want multi-cluster FDB with Raft-driven failover you could do that too (iirc, FDB doesn't have this out of the box).
I don't know what you mean by this. Multiple FDB clusters?
Yes multiple FDB clusters. IIRC FDB replication doesn't support full geo-replication, or didn't. There's a post by me about it somewhere on their forums.
For reading it has a 5 second snapshot timeout that gets in the way. One can stitch multiple transactions together but that could mean losing snapshot isolation without further tricks.
In other words, even just for read-mostly workloads it has a few warts.
I always wanted my app to use a fully distributed database (for redundancy). I've been using RethinkDB in production for over 8 years now. I'm slowly rebuilding my app to use FoundationDB.
What I discovered when I started using FDB surprised me a bit. To make really good use of the database you can't really use a "database layer" and fully abstract it away from your app. Your code should be fully aware of transaction boundaries, for example. To make good use of versionstamps (an incredible feature) your code needs to be somewhat aware of them.
I think FDB is a great candidate for implementing a "user-friendly" database on top of it, and in fact several databases are doing exactly that (using FDB as a "lower layer"). But that abstracts away too much, at least for me.
The superficial take on FDB is "waah, where are my features? it doesn't do indexing? waaah, just use Postgres!".
But when you actually start mapping your app's data structures onto FDB optimally, you discover a whole new world. For example, I ended up writing my indexing code myself, in my language (Clojure). FDB gives you all the tools, and a strict serializable data model to work with — your language brings your data structures and your indexing functions. The combination is incredible. Once you define your index functions in your language, you will never want to look at SQL again. Plus, you get incredible features like versionstamps — I use them to replace RethinkDB changefeeds and implement super quick polling for recent changes.
Oh, and did I mention that it is a fully distributed database that correctly implements the strict serializable consistency model? There are very few dbs that can claim that. If you understand what that means, you probably know how incredible this is and I don't have to convince you. If you think you understand, I suggest you go and explore https://jepsen.io/consistency — carefully reading and learning about the differences in various consistency models.
I really worry that FoundationDB will not become popular because of its inherent complexity, while worse solutions (ahem, MongoDB) will be more fashionable.
I would encourage everyone to at least take a look at FDB. It really is something quite different.
Are you saying you were “in the room” when the decision was made? Can you elaborate?
Edit: indeed you were… https://www.snowflake.com/wp-content/uploads/2020/11/Rise-of...
AFAIK, no one has lost a byte to an FDB issue.
i love fdb, but, most people _should_ just use postgres. you should have a very precise explanation for why you want fdb instead of postgres. you're trading away a lot of things when you go to fdb.
> I would encourage everyone to at least take a look at FDB.
yes, make sure you're aware of it so you can spot the situations where it _is_ the right answer.
Ideally someone could implement the firestore or dynamodb api on top.
https://github.com/losfair/mvsqlite
Is basically distributed SQLite backed by FDB. I’ve been scared to use it since I don’t know rust and can’t attest to if mvcc had been implemented correctly.
In using this I actually realized how coupled the storage engine is to the storage system and how few open source projects make the storage engine easily swap-able.
mvsqlite seems to improve the transaction size [3], which is nice. Does it also improve the key/value limitations?
> Transaction size cannot exceed 10,000,000 bytes of affected data. [---] Keys cannot exceed 10,000 bytes in size. Values cannot exceed 100,000 bytes in size.
[1] https://apple.github.io/foundationdb/known-limitations.html
mvsqlite fixes the transaction size through its own transaction layer, from my understanding; I don't know how that would impact performance. The 10kb/100Kb key value limit is probably not fixable in any way, but it's not really a huge problem as a user in practice for FDB because you can just shard the value across two keys in a consistent transaction and it's fine. 10 kilobyte keys have pretty much never ever been an issue in my cases either; you can typically just do something like hash a really big key before insert and use that.
The record layer https://github.com/FoundationDB/fdb-record-layer which allows to store protobuf, and define the primary keys and index directly on those proto fields is truly amazing:
https://github.com/FoundationDB/fdb-record-layer/blob/main/d...
A less known but also great talk is the follow which talked about what the a few of the team worked on next, effectively trying to generalize the methodology to any computer program: https://www.youtube.com/watch?v=fFSPwJFXVlw
I liken the approach to being able to fuzz the execution space of the program, not just the inputs.
https://deno.com/kv on FoundationDB!
What I can tell you, for sure, is that if you find an issue with something as important and fundamental as data loss the team working on FoundationDB would take it super seriously.
1. https://www.datadoghq.com/blog/engineering/introducing-husky...
2. https://www.datadoghq.com/blog/engineering/husky-deep-dive/
3. https://www.youtube.com/watch?v=mNneCaZewTg
4. https://www.youtube.com/watch?v=1-zo9jqdRZU
I was involved with this project from the beginning and it would've taken significantly longer to deliver without FoundationDB.
https://news.ycombinator.com/item?id=16880404
Also, in the post itself, authors including Apple and Snowflake devs, it mentions it's run in production by Apple and Snowflake.
I haven't seen yet though what Apple uses it for.
https://machinelearning.apple.com/research/foundationdb-reco...
Say your system is "well tested" for example with a combination of unit tests, integration tests, stress tests, failure injection tests (a la Jepsen), and more.
There are probably still hundreds of nasty little bugs hidden. Not even thousands of experiment hours with a Jepsen-like approach would surface them because they are very, very unlikely. Antithesis will find these without breaking a sweat. It's designed to hunt for very unlikely (but possible) scenarios where your system misbehaves.
And here's the real kicker: once a bug is found, you can observe and step-through execution of your entire _distributed_ system. Similar to attaching a debugger to a single process but for an entire system composed of many clients and servers connected by a network.
It's language independent and doesn't require any modification to your system in order to use. It's pretty incredible. I would not believe this is possible if it hadn't seen it with my own eyes.
Presently, I am using it for work at Synadia (makers of NATS). NATS is like lego blocks to build all kinds of distributed systems in multiple languages, so it has a very large surface area. It's well tested, stable, and deployed successfully by many large and small projects and companies. Contemporary testing approaches can hardly find bugs. Antithesis can find very insidious edge cases where things break. And we proactively investigate and fix before any user/customer can be affected by these one-in-a-million nasty bugs, which would otherwise be very hard to find and resolve.
Also, they are hiring.
https://github.com/ccorcos/tuple-database/
I have found it conceptually similar to Relic or Datascript, but with strong preformance guarantees - something Relic considers a potential issue. It also solves the problem of using reactive queries to trigger things like popups and fullscreen requests, which must be run in the same event loop as user input.
https://github.com/wotbrew/relic https://github.com/tonsky/datascript
Having a full (fast!) database as my React state manager gives me LAMP nostalgia :)
Unfortunately they were acquired by Apple, only to resurface something like 10 years later. All momentum was gone, and I’m not really aware nor interested in where they stand. I’ll stick with my rusty old Postgres for a long time before I’d try anything else out.
Maybe this thing still exists in close source form at Apple? It wouldn't surprise me if it does and forms the basis of a Spanner alternative, they're big enough to need it. Or maybe they canned it pre/post acquisition.
Edit: ah, you've already mentioned the closed source layer that exists at Apple. There we go!
It’s similar to mongo (it’s nosql)
At the time of its release it was probably hard to justify having an order of magnitude more latency than competitors (of course they were not fault tolerant, but still).
Apple then bought it up and shut the open source down. They had to rebuild whole layers from scratch.
FoundationDB was a swear word.
And here's a news story - https://www.forbes.com/sites/benkepes/2015/03/25/a-cautionar...
The point wasn't to bash Apple - the point was that the team lost a lot of money and time and shipped and inferior product because Apple chose to purchase and close down part of the stack.
Is this such a strange use case, that there is not even a blog entry about it only some forum entries?
https://forums.foundationdb.org/t/designing-key-value-expira...
[1]: https://developer.apple.com/videos/play/wwdc2023/10164/?time...